From 96dc6311eaba5765ee999be719b17a441c2f9c78 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: Sat, 3 Oct 2026 22:48:32 +0000 Subject: [PATCH] Upload files to "internal/storage" --- internal/storage/document.go | 786 +++++++++++++++++++++++++++++++++++ 1 file changed, 786 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..6d17b98 --- /dev/null +++ b/internal/storage/document.go @@ -0,0 +1,786 @@ +/* + * 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 +}