diff --git a/internal/storage/document.go b/internal/storage/document.go deleted file mode 100644 index 6d17b98..0000000 --- a/internal/storage/document.go +++ /dev/null @@ -1,786 +0,0 @@ -/* - * 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. -// Lock-free: Document.fields переведён на atomic.Value для wait-free доступа. -// -// ИСПРАВЛЕНО (2026): Устранён race condition в Document. -// Все методы, читающие/записывающие метаданные (ID, CreatedAt, UpdatedAt, -// DeletedAt, Version, Compressed, OriginalSize), теперь используют -// sync.RWMutex для защиты от гонок. -// -// ИЗМЕНЕНО (2026): Убраны вызовы AuditFieldOperation из SetField/DeleteField -// (они были неинформативными — db/coll пустые). Аудит выполняется на уровне -// Collection (AuditDocumentOperation). -// -// ИЗМЕНЕНО (2026): Удалены поля Version, Compressed, OriginalSize из -// MVCC-контекста — они остались, но больше не используются для версионирования. -// Версия документа увеличивается при каждом изменении, но time-travel -// реализован через WAL (см. transactions.go: GetDocumentAtTimestamp). -// -// ИЗМЕНЕНО (2026-10-03): формат временных меток изменён с "2006-01-02 15:04:05.000" -// на "02-012006 15:04:05.000" для совместимости с OpenIndiana. - -package storage - -import ( - "fmt" - "math/rand" - "strings" - "sync" - "sync/atomic" - "time" - - "futriis/internal/compression" - "futriis/internal/serializer" - "github.com/google/uuid" -) - -// ============================================================================= -// КОНСТАНТЫ ДЛЯ ЗАЩИТЫ ОТ LIVELOCK -// ============================================================================= - -const ( - // Максимальное количество попыток CAS перед применением backoff - maxCASRetriesBeforeBackoff = 10 - // Базовая задержка backoff в наносекундах - baseBackoffNanos = 100 - // Максимальная задержка backoff в наносекундах - maxBackoffNanos = 10000 -) - -// ============================================================================= -// FIELD VERSION - ЗАЩИТА ОТ ABA -// ============================================================================= - -// fieldVersionWrapper оборачивает map полей с версией для защиты от ABA. -type fieldVersionWrapper struct { - fields map[string]interface{} - version uint64 -} - -// Document представляет документ в коллекции (аналог строки в реляционной СУБД). -// -// ИСПРАВЛЕНО: добавлен sync.RWMutex для защиты метаданных документа -// (ID, CreatedAt, UpdatedAt, DeletedAt, Version, Compressed, OriginalSize) -// от гонок при конкурентном чтении/записи. -type Document struct { - mu sync.RWMutex // защита метаданных документа - - ID string `msgpack:"_id"` // Уникальный идентификатор документа - fieldsPtr atomic.Value // *fieldVersionWrapper - lock-free хранилище полей с версией - CreatedAt int64 `msgpack:"created_at"` // Время создания (Unix миллисекунды) - UpdatedAt int64 `msgpack:"updated_at"` // Время последнего обновления - DeletedAt int64 `msgpack:"deleted_at"` // Время удаления (Unix миллисекунды, 0 = не удалён) - Version uint64 `msgpack:"version"` // Версия документа - Compressed bool `msgpack:"compressed"` // Флаг, сжат ли документ - OriginalSize int64 `msgpack:"original_size"` // Оригинальный размер до сжатия -} - -// 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). -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 загружает карту полей (lock-free). -// ВНИМАНИЕ: возвращает ссылку на внутреннюю map. Не мутировать! -func (d *Document) loadFields() 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{}) - 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 { - backoff := baseBackoffNanos * (1 << min(attempt-maxCASRetriesBeforeBackoff, 10)) - if backoff > maxBackoffNanos { - backoff = maxBackoffNanos - } - jitter := rand.Int63n(int64(backoff)/2 + 1) - totalBackoff := backoff + int(jitter) - time.Sleep(time.Duration(totalBackoff)) - } - - if attempt > 1000 { - return false - } - } -} - -// min возвращает минимальное из двух чисел. -func min(a, b int) int { - if a < b { - return a - } - return b -} - -// SetField устанавливает значение поля документа (lock-free). -// -// ИЗМЕНЕНО: убран вызов AuditFieldOperation (был неинформативным). -func (d *Document) SetField(name string, value interface{}) { - d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { - fields[name] = value - return fields - }) - - d.mu.Lock() - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - d.mu.Unlock() -} - -// GetField возвращает значение поля документа (lock-free). -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 удаляет поле из документа (lock-free). -// -// ИЗМЕНЕНО: убран вызов AuditFieldOperation. -func (d *Document) DeleteField(name string) { - d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { - delete(fields, name) - return fields - }) - - d.mu.Lock() - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - d.mu.Unlock() -} - -// HasField проверяет наличие поля в документе (lock-free). -func (d *Document) HasField(name string) bool { - fields := d.loadFields() - _, ok := fields[name] - return ok -} - -// GetFields возвращает копию всех полей документа (lock-free). -func (d *Document) GetFields() map[string]interface{} { - fields := d.loadFields() - copy := make(map[string]interface{}) - for k, v := range fields { - copy[k] = v - } - return copy -} - -// ToMap возвращает полное представление документа в виде map. -func (d *Document) ToMap() map[string]interface{} { - fields := d.GetFields() - result := make(map[string]interface{}) - - 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 с поддержкой сжатия. -func (d *Document) Serialize() ([]byte, error) { - fields := d.GetFields() - - d.mu.RLock() - docCopy := &Document{ - ID: d.ID, - CreatedAt: d.CreatedAt, - UpdatedAt: d.UpdatedAt, - DeletedAt: d.DeletedAt, - Version: d.Version, - Compressed: d.Compressed, - OriginalSize: d.OriginalSize, - } - d.mu.RUnlock() - - docCopy.storeFields(fields) - - data, err := serializer.Marshal(docCopy) - if err != nil { - return nil, err - } - - return data, nil -} - -// 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 - } - return compressed, nil - } - - return data, nil -} - -// Deserialize десериализует документ из MessagePack. -func (d *Document) Deserialize(data []byte) error { - decompressed, err := compression.DecompressAuto(data) - if err == nil && len(decompressed) < len(data) { - var doc Document - if err := serializer.Unmarshal(decompressed, &doc); err != nil { - return err - } - d.mu.Lock() - d.ID = doc.ID - d.CreatedAt = doc.CreatedAt - d.UpdatedAt = doc.UpdatedAt - d.DeletedAt = doc.DeletedAt - d.Version = doc.Version - d.Compressed = true - d.OriginalSize = int64(len(decompressed)) - d.mu.Unlock() - d.storeFields(doc.loadFields()) - } else { - var doc Document - if err := serializer.Unmarshal(data, &doc); err != nil { - return err - } - d.mu.Lock() - d.ID = doc.ID - d.CreatedAt = doc.CreatedAt - d.UpdatedAt = doc.UpdatedAt - d.DeletedAt = doc.DeletedAt - d.Version = doc.Version - d.Compressed = false - d.OriginalSize = 0 - d.mu.Unlock() - d.storeFields(doc.loadFields()) - } - - d.mu.Lock() - d.UpdatedAt = time.Now().UnixMilli() - d.mu.Unlock() - return nil -} - -// Clone создаёт глубокую копию документа (lock-free). -func (d *Document) Clone() *Document { - fields := d.GetFields() - - d.mu.RLock() - clone := &Document{ - ID: d.ID, - CreatedAt: d.CreatedAt, - UpdatedAt: d.UpdatedAt, - DeletedAt: d.DeletedAt, - Version: d.Version, - Compressed: d.Compressed, - OriginalSize: d.OriginalSize, - } - d.mu.RUnlock() - - copiedFields := make(map[string]interface{}) - for k, v := range fields { - copiedFields[k] = deepCopyValue(v) - } - clone.storeFields(copiedFields) - - return clone -} - -// Update применяет обновление к документу (атомарно, lock-free). -func (d *Document) Update(updates map[string]interface{}) error { - success := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { - for k, v := range updates { - fields[k] = v - } - return fields - }) - - if success { - d.mu.Lock() - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - d.mu.Unlock() - return nil - } - - return fmt.Errorf("failed to update document after multiple attempts (possible high contention)") -} - -// 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 возвращает человекочитаемую строку времени удаления. -// -// ИЗМЕНЕНО (2026-10-03): формат изменён на "02-012006 15:04:05.000". -func (d *Document) GetDeletedAtStr() string { - d.mu.RLock() - deletedAt := d.DeletedAt - d.mu.RUnlock() - if deletedAt == 0 { - return "" - } - return time.UnixMilli(deletedAt).Format("02-012006 15:04:05.000") -} - -// GetCreatedAtStr возвращает человекочитаемую строку времени создания. -// -// ИЗМЕНЕНО (2026-10-03): формат изменён на "02-012006 15:04:05.000". -func (d *Document) GetCreatedAtStr() string { - d.mu.RLock() - createdAt := d.CreatedAt - d.mu.RUnlock() - return time.UnixMilli(createdAt).Format("02-012006 15:04:05.000") -} - -// GetUpdatedAtStr возвращает человекочитаемую строку времени обновления. -// -// ИЗМЕНЕНО (2026-10-03): формат изменён на "02-012006 15:04:05.000". -func (d *Document) GetUpdatedAtStr() string { - d.mu.RLock() - updatedAt := d.UpdatedAt - d.mu.RUnlock() - return time.UnixMilli(updatedAt).Format("02-012006 15:04:05.000") -} - -// Compress сжимает документ в памяти. -func (d *Document) Compress(config *compression.Config) error { - d.mu.Lock() - defer d.mu.Unlock() - - if d.Compressed { - return nil - } - - fields := d.loadFields() - originalSize := len(fields) - if originalSize < config.MinSize { - return nil - } - - d.Compressed = true - d.OriginalSize = int64(originalSize) - - return nil -} - -// Decompress распаковывает документ в памяти. -func (d *Document) Decompress() error { - d.mu.Lock() - defer d.mu.Unlock() - - if !d.Compressed { - return nil - } - - d.Compressed = false - d.OriginalSize = 0 - - return nil -} - -// GetCompressionRatio возвращает коэффициент сжатия. -func (d *Document) GetCompressionRatio() float64 { - d.mu.RLock() - compressed := d.Compressed - originalSize := d.OriginalSize - d.mu.RUnlock() - - if !compressed || originalSize == 0 { - return 1.0 - } - - fields := d.loadFields() - currentSize := len(fields) - return float64(currentSize) / 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{}: - copy := make(map[string]interface{}) - for k, val := range v { - copy[k] = deepCopyValue(val) - } - return copy - case []interface{}: - copy := make([]interface{}, len(v)) - for i, val := range v { - copy[i] = deepCopyValue(val) - } - return copy - default: - return v - } -} - -// 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() - - copy := make(map[string]interface{}) - for k, v := range t.Fields { - copy[k] = v - } - return copy -} - -// 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") - } - - 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 устанавливает значение по точечному пути. -func (d *Document) SetNestedField(path string, value interface{}) error { - 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: - if !v.HasField(part) { - newTuple := NewTuple() - v.SetField(part, newTuple) - current = newTuple - } else { - field, _ := v.GetField(part) - if tuple, ok := field.(*Tuple); ok { - current = tuple - } else { - return fmt.Errorf("field %s is not a tuple", part) - } - } - case *Tuple: - if val, err := v.Get(part); err == nil { - if tuple, ok := val.(*Tuple); ok { - current = tuple - } else { - return fmt.Errorf("field %s is not a tuple", part) - } - } else { - newTuple := NewTuple() - v.Set(part, newTuple) - current = newTuple - } - 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) - default: - return fmt.Errorf("cannot set field on non-document value") - } - - d.mu.Lock() - d.UpdatedAt = time.Now().UnixMilli() - d.Compressed = false - d.mu.Unlock() - return nil -}