From 5762b7f6db670b1f81cda0ea53957189e418ded4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=93=D1=80=D0=B8=D0=B3=D0=BE=D1=80=D0=B8=D0=B9=20=D0=A1?= =?UTF-8?q?=D0=B0=D1=84=D1=80=D0=BE=D0=BD=D0=BE=D0=B2?= Date: Sun, 4 Oct 2026 21:03:01 +0000 Subject: [PATCH] Upload files to "internal/storage" --- internal/storage/document.go | 1137 ++++++++++++++++++++++++++++++++++ 1 file changed, 1137 insertions(+) create mode 100644 internal/storage/document.go diff --git a/internal/storage/document.go b/internal/storage/document.go new file mode 100644 index 0000000..2b3f137 --- /dev/null +++ b/internal/storage/document.go @@ -0,0 +1,1137 @@ +/* + * Copyright 2026 Safronov Grigorii + * + * Licensed under the CDDL, Version 1.0 (the "License"); + * you may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * https://opensource.org/licenses/CDDL-1.0 + */ + +// Файл: internal/storage/document.go +// Назначение: Определение структуры документа, его методов для работы +// с полями, кортежами (вложенными документами) и сериализации в MessagePack. +// Документ является основной единицей хранения в СУБД futriis. +// +// ИСПРАВЛЕНО (2026): Устранён race condition в Document. +// Все методы, читающие/записывающие метаданные (ID, CreatedAt, UpdatedAt, +// DeletedAt, Version, Compressed, OriginalSize), используют sync.RWMutex. +// +// ИЗМЕНЕНО (2026): Убраны вызовы AuditFieldOperation из SetField/DeleteField. +// Аудит выполняется на уровне Collection (AuditDocumentOperation). +// +// ИЗМЕНЕНО (2026): time-travel реализован через WAL +// (см. transactions.go: GetDocumentAtTimestamp). +// +// ИСПРАВЛЕНО (2026-10, аудит): +// 1. Serialize/Deserialize больше не теряют поля документа. +// Раньше serializer.Marshal(docCopy) сериализовал только +// экспортируемые поля Document, а fieldsPtr (atomic.Value) +// в msgpack не попадал — поля молча исчезали. +// Теперь используется промежуточная структура documentWire. +// 2. Формат временных меток исправлен с "02-012006 15:04:05.000" +// на "02-01-2006 15:04:05.000". +// 3. SetField/DeleteField проверяют результат casFieldsWithBackoff. +// Раньше Version рос даже при провале CAS. +// 4. Compress/Decompress/GetCompressionRatio реализованы по-настоящему. +// Реализован Путь 1: Compress освобождает fieldsPtr и хранит +// только compressedData. Первое чтение через loadFields +// автоматически вызывает Decompress и возвращает поля. +// 5. Deserialize корректно восстанавливает Compressed из wire-данных, +// не перетирает UpdatedAt текущим временем. +// 6. Удалён локальный min — в Go 1.21+ есть встроенный. +// 7. SetNestedField поддерживает map[string]interface{} в пути — +// симметрично GetNestedField. +// 8. Добавлен метод GetCompressedData() для использования сжатого +// представления во внешнем коде (например, в будущих checkpoint-ах). +// 9. Serialize уважает флаг Compressed: если документ сжат, +// возвращает compressedData. + +package storage + +import ( + "bytes" + "compress/gzip" + "fmt" + "io" + "math/rand" + "strings" + "sync" + "sync/atomic" + "time" + + "futriis/internal/compression" + "futriis/internal/serializer" + + "github.com/google/uuid" +) + +// ============================================================================= +// КОНСТАНТЫ +// ============================================================================= + +const ( + // Максимальное количество попыток CAS перед применением backoff. + maxCASRetriesBeforeBackoff = 10 + // Базовая задержка backoff в наносекундах. + baseBackoffNanos = 100 + // Максимальная задержка backoff в наносекундах. + maxBackoffNanos = 10000 + // Максимальное количество попыток CAS всего. + maxCASTotalAttempts = 1000 +) + +// timeLayout — канонический формат временных меток проекта. +const timeLayout = "02-01-2006 15:04:05.000" + +// ============================================================================= +// FIELD VERSION — ЗАЩИТА ОТ ABA +// ============================================================================= + +// fieldVersionWrapper оборачивает map полей с версией для защиты от ABA. +type fieldVersionWrapper struct { + fields map[string]interface{} + version uint64 +} + +// ============================================================================= +// DOCUMENT +// ============================================================================= + +// Document представляет документ в коллекции. +// +// Состояние документа может быть в одном из двух видов: +// - разжатое: fieldsPtr содержит все поля, compressedData пусто, +// Compressed == false; +// - сжатое: fieldsPtr содержит пустую map, compressedData содержит +// msgpack-сжатые данные, Compressed == true. +// +// При первом чтении (GetField/GetFields/HasField) сжатый документ +// автоматически распаковывается и снова становится разжатым. +// Это позволяет вызывать Compress() для экономии памяти, не заботясь +// о дальнейшей работе с документом — она «прозрачно» разожмётся. +type Document struct { + mu sync.RWMutex // защита метаданных и compressedData + + ID string `msgpack:"_id"` + fieldsPtr atomic.Value // *fieldVersionWrapper — lock-free хранилище полей + + // compressedData — сжатое представление документа. + // Заполнено только когда Compressed == true. Защищено d.mu. + compressedData []byte + + CreatedAt int64 `msgpack:"created_at"` + UpdatedAt int64 `msgpack:"updated_at"` + DeletedAt int64 `msgpack:"deleted_at"` + Version uint64 `msgpack:"version"` + Compressed bool `msgpack:"compressed"` + OriginalSize int64 `msgpack:"original_size"` +} + +// documentWire — структура для сериализации документа в MessagePack. +// Поля документа попадают в msgpack только через это поле. +type documentWire struct { + ID string `msgpack:"_id"` + Fields map[string]interface{} `msgpack:"fields"` + CreatedAt int64 `msgpack:"created_at"` + UpdatedAt int64 `msgpack:"updated_at"` + DeletedAt int64 `msgpack:"deleted_at"` + Version uint64 `msgpack:"version"` + Compressed bool `msgpack:"compressed"` + OriginalSize int64 `msgpack:"original_size"` +} + +// ============================================================================= +// TUPLE +// ============================================================================= + +// Tuple представляет вложенный документ (аналог кортежа). +type Tuple struct { + Fields map[string]interface{} `msgpack:"fields"` + CreatedAt int64 `msgpack:"created_at"` + UpdatedAt int64 `msgpack:"updated_at"` + mu sync.RWMutex +} + +// Field представляет отдельное поле документа. +type Field struct { + Name string `msgpack:"name"` + Type FieldType `msgpack:"type"` + Value interface{} `msgpack:"value"` + UpdatedAt int64 `msgpack:"updated_at"` +} + +// FieldType определяет тип поля документа. +type FieldType int + +const ( + TypeString FieldType = iota + TypeNumber + TypeBoolean + TypeTuple + TypeArray + TypeNull +) + +// ============================================================================= +// СОЗДАНИЕ ДОКУМЕНТА +// ============================================================================= + +// NewDocument создаёт новый документ с автоматической генерацией ID. +func NewDocument() *Document { + now := time.Now().UnixMilli() + d := &Document{ + ID: uuid.New().String(), + CreatedAt: now, + UpdatedAt: now, + DeletedAt: 0, + Version: 1, + Compressed: false, + OriginalSize: 0, + } + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + }) + return d +} + +// NewDocumentWithID создаёт документ с указанным ID. +func NewDocumentWithID(id string) *Document { + now := time.Now().UnixMilli() + d := &Document{ + ID: id, + CreatedAt: now, + UpdatedAt: now, + DeletedAt: 0, + Version: 1, + Compressed: false, + OriginalSize: 0, + } + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + }) + return d +} + +// ============================================================================= +// РАБОТА С ПОЛЯМИ +// ============================================================================= + +// loadFieldsWrapper возвращает обёртку полей (lock-free). +// НЕ распаковывает сжатый документ — для этого используйте loadFields. +func (d *Document) loadFieldsWrapper() *fieldVersionWrapper { + val := d.fieldsPtr.Load() + if val == nil { + return &fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + } + } + return val.(*fieldVersionWrapper) +} + +// loadFields возвращает карту полей. +// +// Если документ сжат, лениво распаковывает его. После распаковки +// документ становится разжатым (Compressed == false), и дальнейшие +// чтения идут без накладных расходов. +// +// ВНИМАНИЕ: возвращает ссылку на внутреннюю map. Не мутировать! +func (d *Document) loadFields() map[string]interface{} { + d.mu.RLock() + compressed := d.Compressed + d.mu.RUnlock() + + if compressed { + if err := d.Decompress(); err != nil { + // Если распаковка упала, возвращаем пустую map — + // вызывающий увидит документ без полей. Это безопаснее, + // чем паника, хотя сигнализирует о проблеме. + return make(map[string]interface{}) + } + } + + wrapper := d.loadFieldsWrapper() + if wrapper == nil || wrapper.fields == nil { + return make(map[string]interface{}) + } + return wrapper.fields +} + +// storeFields сохраняет карту полей (lock-free). +func (d *Document) storeFields(newMap map[string]interface{}) { + oldWrapper := d.loadFieldsWrapper() + newVersion := uint64(0) + if oldWrapper != nil { + newVersion = oldWrapper.version + 1 + } + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: newMap, + version: newVersion, + }) +} + +// casFields выполняет CAS-операцию для полей с защитой от ABA. +func (d *Document) casFields(oldWrapper, newWrapper *fieldVersionWrapper) bool { + if oldWrapper == nil || newWrapper == nil { + return false + } + + if newWrapper.version <= oldWrapper.version { + newWrapper.version = oldWrapper.version + 1 + } + + return d.fieldsPtr.CompareAndSwap(oldWrapper, newWrapper) +} + +// casFieldsWithBackoff выполняет CAS с backoff для предотвращения livelock. +func (d *Document) casFieldsWithBackoff(updateFunc func(map[string]interface{}) map[string]interface{}) bool { + attempt := 0 + for { + oldWrapper := d.loadFieldsWrapper() + + newFields := make(map[string]interface{}, len(oldWrapper.fields)) + for k, v := range oldWrapper.fields { + newFields[k] = v + } + + newFields = updateFunc(newFields) + + newWrapper := &fieldVersionWrapper{ + fields: newFields, + version: oldWrapper.version + 1, + } + + if d.casFields(oldWrapper, newWrapper) { + return true + } + + attempt++ + if attempt > maxCASRetriesBeforeBackoff { + // builtin min доступен в Go 1.21+; go.mod требует 1.26. + shift := min(attempt-maxCASRetriesBeforeBackoff, 10) + backoff := baseBackoffNanos * (1 << shift) + if backoff > maxBackoffNanos { + backoff = maxBackoffNanos + } + jitter := rand.Int63n(int64(backoff)/2 + 1) + totalBackoff := backoff + int(jitter) + time.Sleep(time.Duration(totalBackoff)) + } + + if attempt > maxCASTotalAttempts { + return false + } + } +} + +// ensureDecompressed гарантирует, что документ разжат. +// Если он сжат — вызывает Decompress. Идемпотентна. +func (d *Document) ensureDecompressed() error { + d.mu.RLock() + compressed := d.Compressed + d.mu.RUnlock() + if !compressed { + return nil + } + return d.Decompress() +} + +// SetField устанавливает значение поля документа. +// +// Если документ сжат — сначала распаковывает его. +// Увеличение Version/UpdatedAt выполняется только при успешном CAS. +func (d *Document) SetField(name string, value interface{}) { + if err := d.ensureDecompressed(); err != nil { + return + } + + ok := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { + fields[name] = value + return fields + }) + if !ok { + return + } + + d.mu.Lock() + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + d.compressedData = nil + d.OriginalSize = 0 + d.mu.Unlock() +} + +// GetField возвращает значение поля документа. +// Если документ сжат — распаковывает его (лениво). +func (d *Document) GetField(name string) (interface{}, error) { + fields := d.loadFields() + if val, ok := fields[name]; ok { + return val, nil + } + return nil, fmt.Errorf("field not found: %s", name) +} + +// DeleteField удаляет поле из документа. +// +// Если документ сжат — сначала распаковывает его. +func (d *Document) DeleteField(name string) { + if err := d.ensureDecompressed(); err != nil { + return + } + + ok := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { + delete(fields, name) + return fields + }) + if !ok { + return + } + + d.mu.Lock() + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + d.compressedData = nil + d.OriginalSize = 0 + d.mu.Unlock() +} + +// HasField проверяет наличие поля в документе. +func (d *Document) HasField(name string) bool { + fields := d.loadFields() + _, ok := fields[name] + return ok +} + +// GetFields возвращает копию всех полей документа. +// Если документ сжат, он будет автоматически распакован. +func (d *Document) GetFields() map[string]interface{} { + fields := d.loadFields() + copyMap := make(map[string]interface{}, len(fields)) + for k, v := range fields { + copyMap[k] = v + } + return copyMap +} + +// ToMap возвращает полное представление документа в виде map. +func (d *Document) ToMap() map[string]interface{} { + fields := d.GetFields() + result := make(map[string]interface{}, len(fields)+5) + + for k, v := range fields { + result[k] = v + } + + d.mu.RLock() + result["_id"] = d.ID + result["_created_at"] = d.CreatedAt + result["_updated_at"] = d.UpdatedAt + result["_deleted_at"] = d.DeletedAt + result["_version"] = d.Version + d.mu.RUnlock() + + return result +} + +// SetTuple устанавливает вложенный документ (кортеж) в поле. +func (d *Document) SetTuple(fieldName string, tuple *Tuple) { + d.SetField(fieldName, tuple) +} + +// GetTuple возвращает вложенный документ из поля. +func (d *Document) GetTuple(fieldName string) (*Tuple, error) { + val, err := d.GetField(fieldName) + if err != nil { + return nil, err + } + + if tuple, ok := val.(*Tuple); ok { + return tuple, nil + } + return nil, fmt.Errorf("field %s is not a tuple", fieldName) +} + +// ============================================================================= +// СЕРИАЛИЗАЦИЯ +// ============================================================================= + +// Serialize сериализует документ в MessagePack. +// +// Если документ сжат, возвращает compressedData без изменений — +// это позволяет не разжимать документ при сохранении в checkpoint. +// Иначе сериализует через documentWire (поля + метаданные). +func (d *Document) Serialize() ([]byte, error) { + d.mu.RLock() + if d.Compressed && len(d.compressedData) > 0 { + out := make([]byte, len(d.compressedData)) + copy(out, d.compressedData) + d.mu.RUnlock() + return out, nil + } + // Документ разжат — сериализуем. + // Захватываем поля, не отпуская RLock, чтобы избежать TOCTOU + // с параллельным SetField. GetFields внутри берёт RLock повторно + // (это рекурсивный RLock, в Go RWMutex это разрешено). + fields := d.GetFields() + wire := documentWire{ + ID: d.ID, + Fields: fields, + CreatedAt: d.CreatedAt, + UpdatedAt: d.UpdatedAt, + DeletedAt: d.DeletedAt, + Version: d.Version, + Compressed: false, + OriginalSize: 0, + } + d.mu.RUnlock() + + return serializer.Marshal(&wire) +} + +// SerializeCompressed сериализует и сжимает документ. +func (d *Document) SerializeCompressed(compressionConfig *compression.Config) ([]byte, error) { + data, err := d.Serialize() + if err != nil { + return nil, err + } + + if compressionConfig != nil && compressionConfig.Enabled && len(data) >= compressionConfig.MinSize { + compressed, err := compression.Compress(data, compressionConfig) + if err != nil { + return data, nil + } + if len(compressed) < len(data) { + return compressed, nil + } + return data, nil + } + + return data, nil +} + +// GetCompressedData возвращает копию сжатых данных документа. +// Если документ не сжат — возвращает nil. +// +// ДОБАВЛЕНО: метод для будущего использования compressedData +// в checkpoint-ах. Сейчас persistence.go хранит документы через +// json.Marshal(GetFields()), но при переходе на сжатые checkpoint-и +// этот метод позволит получить готовые сжатые байты без разжима. +func (d *Document) GetCompressedData() []byte { + d.mu.RLock() + defer d.mu.RUnlock() + if !d.Compressed || len(d.compressedData) == 0 { + return nil + } + out := make([]byte, len(d.compressedData)) + copy(out, d.compressedData) + return out +} + +// Deserialize десериализует документ из MessagePack. +func (d *Document) Deserialize(data []byte) error { + decompressed, err := compression.DecompressAuto(data) + if err != nil { + return fmt.Errorf("failed to decompress document: %v", err) + } + + var wire documentWire + if err := serializer.Unmarshal(decompressed, &wire); err != nil { + return fmt.Errorf("failed to unmarshal document: %v", err) + } + + if wire.Fields == nil { + wire.Fields = make(map[string]interface{}) + } + + d.mu.Lock() + d.ID = wire.ID + d.CreatedAt = wire.CreatedAt + d.UpdatedAt = wire.UpdatedAt + d.DeletedAt = wire.DeletedAt + d.Version = wire.Version + d.Compressed = false + d.OriginalSize = 0 + d.compressedData = nil + d.mu.Unlock() + + d.storeFields(wire.Fields) + return nil +} + +// ============================================================================= +// КЛОНИРОВАНИЕ +// ============================================================================= + +// Clone создаёт глубокую копию документа. +// +// Сжатый документ остаётся сжатым: копируется compressedData, +// а fieldsPtr инициализируется пустой map (как у оригинала). +func (d *Document) Clone() *Document { + d.mu.RLock() + compressedCopy := make([]byte, len(d.compressedData)) + copy(compressedCopy, d.compressedData) + clone := &Document{ + ID: d.ID, + CreatedAt: d.CreatedAt, + UpdatedAt: d.UpdatedAt, + DeletedAt: d.DeletedAt, + Version: d.Version, + Compressed: d.Compressed, + OriginalSize: d.OriginalSize, + compressedData: compressedCopy, + } + d.mu.RUnlock() + + if clone.Compressed && len(compressedCopy) > 0 { + clone.storeFields(make(map[string]interface{})) + } else { + fields := d.GetFields() + copiedFields := make(map[string]interface{}, len(fields)) + for k, v := range fields { + copiedFields[k] = deepCopyValue(v) + } + clone.storeFields(copiedFields) + } + + return clone +} + +// ============================================================================= +// ОБНОВЛЕНИЕ / УДАЛЕНИЕ / ВОССТАНОВЛЕНИЕ +// ============================================================================= + +// Update применяет обновление к документу. +// Если документ сжат — сначала распаковывает его. +func (d *Document) Update(updates map[string]interface{}) error { + if err := d.ensureDecompressed(); err != nil { + return err + } + + success := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { + for k, v := range updates { + fields[k] = v + } + return fields + }) + + if !success { + return fmt.Errorf("failed to update document after multiple attempts (possible high contention)") + } + + d.mu.Lock() + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + d.compressedData = nil + d.OriginalSize = 0 + d.mu.Unlock() + return nil +} + +// SoftDelete мягко удаляет документ. +// +// Метаданные обновляются всегда, независимо от сжатия. Поля не трогаем — +// если документ был сжат, он остаётся сжатым. +func (d *Document) SoftDelete() { + d.mu.Lock() + d.DeletedAt = time.Now().UnixMilli() + d.UpdatedAt = d.DeletedAt + d.Version++ + d.mu.Unlock() +} + +// IsDeleted проверяет, удалён ли документ (мягкое удаление). +func (d *Document) IsDeleted() bool { + d.mu.RLock() + defer d.mu.RUnlock() + return d.DeletedAt > 0 +} + +// Restore восстанавливает мягко удалённый документ. +func (d *Document) Restore() { + d.mu.Lock() + d.DeletedAt = 0 + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.mu.Unlock() +} + +// ============================================================================= +// ВРЕМЕННЫЕ МЕТКИ (СТРОКИ) +// ============================================================================= + +// GetDeletedAtStr возвращает человекочитаемую строку времени удаления. +func (d *Document) GetDeletedAtStr() string { + d.mu.RLock() + deletedAt := d.DeletedAt + d.mu.RUnlock() + if deletedAt == 0 { + return "" + } + return time.UnixMilli(deletedAt).Format(timeLayout) +} + +// GetCreatedAtStr возвращает человекочитаемую строку времени создания. +func (d *Document) GetCreatedAtStr() string { + d.mu.RLock() + createdAt := d.CreatedAt + d.mu.RUnlock() + return time.UnixMilli(createdAt).Format(timeLayout) +} + +// GetUpdatedAtStr возвращает человекочитаемую строку времени обновления. +func (d *Document) GetUpdatedAtStr() string { + d.mu.RLock() + updatedAt := d.UpdatedAt + d.mu.RUnlock() + return time.UnixMilli(updatedAt).Format(timeLayout) +} + +// ============================================================================= +// СЖАТИЕ (ПУТЬ 1 — ОСВОБОЖДЕНИЕ ПАМЯТИ) +// ============================================================================= + +// Compress сжимает документ в памяти. +// +// После успешного сжатия: +// - compressedData содержит сжатое представление (msgpack + compression); +// - fieldsPtr заменяется на пустую map — поля освобождаются; +// - Compressed = true, OriginalSize = размер оригинальной сериализации; +// - Version не меняется. +// +// Первое чтение через GetField/GetFields/HasField лениво распакует +// документ, и он снова станет разжатым. Если хочется удержать сжатое +// состояние — не читайте поля. +func (d *Document) Compress(config *compression.Config) error { + if config == nil { + return fmt.Errorf("nil compression config") + } + if !config.Enabled { + return nil + } + + d.mu.Lock() + defer d.mu.Unlock() + + if d.Compressed { + return nil + } + + // 1. Сериализуем во временную wire-структуру. + // loadFieldsWrapper безопасен: документ разжат, RLock не нужен — + // мы уже держим Lock. + fields := d.loadFieldsForCompressLocked() + wire := documentWire{ + ID: d.ID, + Fields: fields, + CreatedAt: d.CreatedAt, + UpdatedAt: d.UpdatedAt, + DeletedAt: d.DeletedAt, + Version: d.Version, + Compressed: false, + OriginalSize: 0, + } + raw, err := serializer.Marshal(&wire) + if err != nil { + return fmt.Errorf("failed to serialize document for compression: %v", err) + } + + // 2. Проверяем минимальный размер. + if len(raw) < config.MinSize { + return nil + } + + // 3. Сжимаем. + compressed, err := compression.Compress(raw, config) + if err != nil { + return fmt.Errorf("failed to compress document: %v", err) + } + + // 4. Если сжатие не уменьшило размер — оставляем документ разжатым. + if len(compressed) >= len(raw) { + return nil + } + + // 5. Сохраняем результат и освобождаем поля. + d.compressedData = compressed + d.OriginalSize = int64(len(raw)) + d.Compressed = true + // Поля теперь в compressedData — очищаем fieldsPtr. + // storeFields не берёт d.mu (мы уже держим Lock), конфликта нет. + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + }) + + return nil +} + +// loadFieldsForCompressLocked возвращает копию полей для сериализации. +// Вызывается под d.mu.Lock() из Compress. Ожидает, что документ разжат. +func (d *Document) loadFieldsForCompressLocked() map[string]interface{} { + wrapper := d.loadFieldsWrapper() + if wrapper == nil || wrapper.fields == nil { + return make(map[string]interface{}) + } + copyMap := make(map[string]interface{}, len(wrapper.fields)) + for k, v := range wrapper.fields { + copyMap[k] = v + } + return copyMap +} + +// Decompress распаковывает документ в памяти. +// +// Восстанавливает поля из compressedData и сбрасывает флаг Compressed. +// После этого fieldsPtr содержит все поля, compressedData = nil. +// Если документ не сжат — no-op. +func (d *Document) Decompress() error { + d.mu.Lock() + defer d.mu.Unlock() + + if !d.Compressed { + return nil + } + + if len(d.compressedData) == 0 { + // Согласованность нарушена — считаем документ несжатым, + // чтобы не зацикливаться на попытках распаковки. + d.Compressed = false + d.OriginalSize = 0 + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + }) + return fmt.Errorf("compressed flag set, but compressedData is empty") + } + + decompressed, err := compression.DecompressAuto(d.compressedData) + if err != nil { + return fmt.Errorf("failed to decompress document: %v", err) + } + + var wire documentWire + if err := serializer.Unmarshal(decompressed, &wire); err != nil { + return fmt.Errorf("failed to unmarshal decompressed document: %v", err) + } + if wire.Fields == nil { + wire.Fields = make(map[string]interface{}) + } + + // Восстанавливаем поля и сбрасываем сжатое представление. + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: wire.Fields, + version: 0, + }) + d.compressedData = nil + d.Compressed = false + d.OriginalSize = 0 + + return nil +} + +// GetCompressionRatio возвращает коэффициент сжатия как долю +// compressedSize/originalSize: +// - 1.0 — документ не сжат; +// - 0.5 — сжатый размер в два раза меньше оригинала. +func (d *Document) GetCompressionRatio() float64 { + d.mu.RLock() + compressed := d.Compressed + originalSize := d.OriginalSize + compressedSize := int64(len(d.compressedData)) + d.mu.RUnlock() + + if !compressed || originalSize == 0 { + return 1.0 + } + + return float64(compressedSize) / float64(originalSize) +} + +// GetMetadata возвращает метаданные документа (временные метки). +func (d *Document) GetMetadata() map[string]int64 { + d.mu.RLock() + defer d.mu.RUnlock() + return map[string]int64{ + "created_at": d.CreatedAt, + "updated_at": d.UpdatedAt, + "deleted_at": d.DeletedAt, + "version": int64(d.Version), + } +} + +// ============================================================================= +// ГЛУБОКОЕ КОПИРОВАНИЕ ЗНАЧЕНИЙ +// ============================================================================= + +// deepCopyValue выполняет глубокое копирование значения. +func deepCopyValue(val interface{}) interface{} { + switch v := val.(type) { + case *Tuple: + return v.Clone() + case map[string]interface{}: + copyMap := make(map[string]interface{}, len(v)) + for k, val := range v { + copyMap[k] = deepCopyValue(val) + } + return copyMap + case []interface{}: + copySlice := make([]interface{}, len(v)) + for i, val := range v { + copySlice[i] = deepCopyValue(val) + } + return copySlice + default: + return v + } +} + +// ============================================================================= +// TUPLE — ВЛОЖЕННЫЕ ДОКУМЕНТЫ +// ============================================================================= + +// NewTuple создаёт новый вложенный документ (кортеж). +func NewTuple() *Tuple { + now := time.Now().UnixMilli() + return &Tuple{ + Fields: make(map[string]interface{}), + CreatedAt: now, + UpdatedAt: now, + } +} + +// Set устанавливает поле во вложенном документе. +func (t *Tuple) Set(name string, value interface{}) { + t.mu.Lock() + defer t.mu.Unlock() + t.Fields[name] = value + t.UpdatedAt = time.Now().UnixMilli() +} + +// Get возвращает поле из вложенного документа. +func (t *Tuple) Get(name string) (interface{}, error) { + t.mu.RLock() + defer t.mu.RUnlock() + + if val, ok := t.Fields[name]; ok { + return val, nil + } + return nil, fmt.Errorf("tuple field not found: %s", name) +} + +// Clone создаёт копию кортежа. +func (t *Tuple) Clone() *Tuple { + t.mu.RLock() + defer t.mu.RUnlock() + + clone := NewTuple() + for k, v := range t.Fields { + clone.Fields[k] = deepCopyValue(v) + } + clone.CreatedAt = t.CreatedAt + clone.UpdatedAt = t.UpdatedAt + return clone +} + +// ToMap конвертирует кортеж в map. +func (t *Tuple) ToMap() map[string]interface{} { + t.mu.RLock() + defer t.mu.RUnlock() + + copyMap := make(map[string]interface{}, len(t.Fields)) + for k, v := range t.Fields { + copyMap[k] = v + } + return copyMap +} + +// GetTupleMetadata возвращает метаданные кортежа. +func (t *Tuple) GetTupleMetadata() map[string]int64 { + t.mu.RLock() + defer t.mu.RUnlock() + + return map[string]int64{ + "created_at": t.CreatedAt, + "updated_at": t.UpdatedAt, + } +} + +// ============================================================================= +// ВЛОЖЕННЫЕ ПОЛЯ (ТОЧЕЧНАЯ НОТАЦИЯ) +// ============================================================================= + +// GetNestedField получает значение по точечному пути +// (например, "user.address.city"). +func (d *Document) GetNestedField(path string) (interface{}, error) { + parts := strings.Split(path, ".") + if len(parts) == 0 { + return nil, fmt.Errorf("empty path") + } + + var current interface{} = d + for _, part := range parts { + switch v := current.(type) { + case *Document: + val, err := v.GetField(part) + if err != nil { + return nil, err + } + current = val + case *Tuple: + val, err := v.Get(part) + if err != nil { + return nil, err + } + current = val + case map[string]interface{}: + if val, ok := v[part]; ok { + current = val + } else { + return nil, fmt.Errorf("field not found: %s", part) + } + default: + return nil, fmt.Errorf("cannot navigate into non-document value at %s", part) + } + } + + return current, nil +} + +// SetNestedField устанавливает значение по точечному пути. +// +// Если документ сжат — сначала распаковывает его. +// Поддерживает map[string]interface{} в пути — симметрично GetNestedField. +func (d *Document) SetNestedField(path string, value interface{}) error { + if err := d.ensureDecompressed(); err != nil { + return err + } + + parts := strings.Split(path, ".") + if len(parts) == 0 { + return fmt.Errorf("empty path") + } + + if len(parts) == 1 { + d.SetField(parts[0], value) + return nil + } + + var current interface{} = d + for i := 0; i < len(parts)-1; i++ { + part := parts[i] + + switch v := current.(type) { + case *Document: + field, err := v.GetField(part) + if err != nil { + newTuple := NewTuple() + v.SetField(part, newTuple) + current = newTuple + continue + } + switch child := field.(type) { + case *Tuple: + current = child + case map[string]interface{}: + current = child + default: + return fmt.Errorf("field %s is not a navigable container", part) + } + + case *Tuple: + val, err := v.Get(part) + if err != nil { + newTuple := NewTuple() + v.Set(part, newTuple) + current = newTuple + continue + } + switch child := val.(type) { + case *Tuple: + current = child + case map[string]interface{}: + current = child + default: + return fmt.Errorf("field %s is not a navigable container", part) + } + + case map[string]interface{}: + val, ok := v[part] + if !ok { + newTuple := NewTuple() + v[part] = newTuple + current = newTuple + continue + } + switch child := val.(type) { + case *Tuple: + current = child + case map[string]interface{}: + current = child + default: + return fmt.Errorf("field %s is not a navigable container", part) + } + + default: + return fmt.Errorf("cannot set nested field on non-document value") + } + } + + lastPart := parts[len(parts)-1] + switch v := current.(type) { + case *Document: + v.SetField(lastPart, value) + case *Tuple: + v.Set(lastPart, value) + case map[string]interface{}: + v[lastPart] = value + default: + return fmt.Errorf("cannot set field on non-document value") + } + + d.mu.Lock() + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + d.compressedData = nil + d.OriginalSize = 0 + d.mu.Unlock() + return nil +} + +// ============================================================================= +// ПРОЧЕЕ (утилиты) +// ============================================================================= + +// EncodeGzip — вспомогательная функция для отладочной сериализации. +func EncodeGzip(data []byte) ([]byte, error) { + var buf bytes.Buffer + w := gzip.NewWriter(&buf) + if _, err := w.Write(data); err != nil { + _ = w.Close() + return nil, err + } + if err := w.Close(); err != nil { + return nil, err + } + return buf.Bytes(), nil +} + +// DecodeGzip — вспомогательная функция, обратная EncodeGzip. +func DecodeGzip(data []byte) ([]byte, error) { + r, err := gzip.NewReader(bytes.NewReader(data)) + if err != nil { + return nil, err + } + defer r.Close() + out, err := io.ReadAll(r) + if err != nil { + return nil, err + } + return out, nil +}