Files
futriix/internal/storage/document.go
T

1137 lines
36 KiB
Go
Raw Normal View History

2026-10-04 21:03:01 +00:00
/*
* 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
}