diff --git a/internal/storage/document.go b/internal/storage/document.go deleted file mode 100644 index 36058b8..0000000 --- a/internal/storage/document.go +++ /dev/null @@ -1,812 +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 для защиты от гонок. - -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 (lock-free) -// ИСПРАВЛЕНО: Использует версионирование для защиты от ABA-проблемы -func (d *Document) casFields(oldWrapper, newWrapper *fieldVersionWrapper) bool { - if oldWrapper == nil || newWrapper == nil { - return false - } - - // Проверяем, что версия увеличилась (защита от ABA) - if newWrapper.version <= oldWrapper.version { - newWrapper.version = oldWrapper.version + 1 - } - - return d.fieldsPtr.CompareAndSwap(oldWrapper, newWrapper) -} - -// casFieldsWithBackoff выполняет CAS с backoff для предотвращения livelock -// ИСПРАВЛЕНО: Добавлен экспоненциальный backoff для предотвращения livelock -// ИСПРАВЛЕНО: Исправлены типы int/int64 для устранения ошибок компилятора -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 - } - - // Backoff при неудаче для предотвращения livelock - attempt++ - if attempt > maxCASRetriesBeforeBackoff { - // Экспоненциальный backoff с jitter - backoff := baseBackoffNanos * (1 << min(attempt-maxCASRetriesBeforeBackoff, 10)) - if backoff > maxBackoffNanos { - backoff = maxBackoffNanos - } - // ИСПРАВЛЕНО: Приводим типы для совместимости с rand.Int63n (int64) - // Добавляем jitter для избежания синхронизации потоков - 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) -// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock -// ИСПРАВЛЕНО: Метаданные защищены мьютексом -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() - - // Аудит изменения поля - AuditFieldOperation("UPDATE", "", "", d.ID, name, value) -} - -// 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) -// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock -// ИСПРАВЛЕНО: Метаданные защищены мьютексом -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() - - // Аудит удаления поля - AuditFieldOperation("DELETE", "", "", d.ID, name, nil) -} - -// 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 -// Включает метаданные (_id, _created_at, _updated_at, _deleted_at, _version) и все поля -// ИСПРАВЛЕНО: метаданные читаются под мьютексом -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) -// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock -// ИСПРАВЛЕНО: Метаданные защищены мьютексом -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 возвращает человекочитаемую строку времени удаления -// ИСПРАВЛЕНО: метаданные читаются под мьютексом -func (d *Document) GetDeletedAtStr() string { - d.mu.RLock() - deletedAt := d.DeletedAt - d.mu.RUnlock() - if deletedAt == 0 { - return "" - } - return time.UnixMilli(deletedAt).Format("2006-01-02 15:04:05.000") -} - -// GetCreatedAtStr возвращает человекочитаемую строку времени создания -// ИСПРАВЛЕНО: метаданные читаются под мьютексом -func (d *Document) GetCreatedAtStr() string { - d.mu.RLock() - createdAt := d.CreatedAt - d.mu.RUnlock() - return time.UnixMilli(createdAt).Format("2006-01-02 15:04:05.000") -} - -// GetUpdatedAtStr возвращает человекочитаемую строку времени обновления -// ИСПРАВЛЕНО: метаданные читаются под мьютексом -func (d *Document) GetUpdatedAtStr() string { - d.mu.RLock() - updatedAt := d.UpdatedAt - d.mu.RUnlock() - return time.UnixMilli(updatedAt).Format("2006-01-02 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 - } - - // Для простоты реализации используем подход с чтением-модификацией-записью - // В production коде потребуется более сложная lock-free структура - - // Сначала проверяем путь - 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 -}