diff --git a/internal/storage/document.go b/internal/storage/document.go index 224e649..09e0569 100644 --- a/internal/storage/document.go +++ b/internal/storage/document.go @@ -14,10 +14,13 @@ // Документ является основной единицей хранения в СУБД futriis. // Lock-free: Document.fields переведён на atomic.Value для wait-free доступа. + + package storage import ( "fmt" + "math/rand" "strings" "sync" "sync/atomic" @@ -28,10 +31,33 @@ import ( "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 представляет документ в коллекции (аналог строки в реляционной СУБД) type Document struct { ID string `msgpack:"_id"` // Уникальный идентификатор документа - fieldsPtr atomic.Value // map[string]interface{} - lock-free хранилище полей + fieldsPtr atomic.Value // *fieldVersionWrapper - lock-free хранилище полей с версией CreatedAt int64 `msgpack:"created_at"` // Время создания (Unix миллисекунды) UpdatedAt int64 `msgpack:"updated_at"` // Время последнего обновления DeletedAt int64 `msgpack:"deleted_at"` // Время удаления (Unix миллисекунды, 0 = не удалён) @@ -80,7 +106,10 @@ func NewDocument() *Document { Compressed: false, OriginalSize: 0, } - d.fieldsPtr.Store(make(map[string]interface{})) + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: make(map[string]interface{}), + version: 0, + }) return d } @@ -96,49 +125,132 @@ func NewDocumentWithID(id string) *Document { Compressed: false, OriginalSize: 0, } - d.fieldsPtr.Store(make(map[string]interface{})) + 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) func (d *Document) loadFields() map[string]interface{} { - val := d.fieldsPtr.Load() - if val == nil { + wrapper := d.loadFieldsWrapper() + if wrapper == nil || wrapper.fields == nil { return make(map[string]interface{}) } - return val.(map[string]interface{}) + return wrapper.fields } // storeFields сохраняет карту полей (lock-free) func (d *Document) storeFields(newMap map[string]interface{}) { - d.fieldsPtr.Store(newMap) + oldWrapper := d.loadFieldsWrapper() + newVersion := uint64(0) + if oldWrapper != nil { + newVersion = oldWrapper.version + 1 + } + d.fieldsPtr.Store(&fieldVersionWrapper{ + fields: newMap, + version: newVersion, + }) } -// SetField устанавливает значение поля документа (lock-free) -func (d *Document) SetField(name string, value interface{}) { +// 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 { - oldFields := d.loadFields() + oldWrapper := d.loadFieldsWrapper() + + // Создаём копию полей для модификации newFields := make(map[string]interface{}) - for k, v := range oldFields { + for k, v := range oldWrapper.fields { newFields[k] = v } - newFields[name] = value - if d.compareAndSwapFields(oldFields, newFields) { - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - - // Аудит изменения поля - AuditFieldOperation("UPDATE", "", "", d.ID, name, value) - return + // Применяем обновление + 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 } } } -// compareAndSwapFields выполняет CAS операцию для полей (lock-free) -func (d *Document) compareAndSwapFields(old, new map[string]interface{}) bool { - return d.fieldsPtr.CompareAndSwap(old, new) +// 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.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + + // Аудит изменения поля + AuditFieldOperation("UPDATE", "", "", d.ID, name, value) } // GetField возвращает значение поля документа (lock-free) @@ -151,30 +263,19 @@ func (d *Document) GetField(name string) (interface{}, error) { } // DeleteField удаляет поле из документа (lock-free) +// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock func (d *Document) DeleteField(name string) { - for { - oldFields := d.loadFields() - if _, exists := oldFields[name]; !exists { - return - } - - newFields := make(map[string]interface{}) - for k, v := range oldFields { - if k != name { - newFields[k] = v - } - } - - if d.compareAndSwapFields(oldFields, newFields) { - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - - // Аудит удаления поля - AuditFieldOperation("DELETE", "", "", d.ID, name, nil) - return - } - } + d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { + delete(fields, name) + return fields + }) + + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + + // Аудит удаления поля + AuditFieldOperation("DELETE", "", "", d.ID, name, nil) } // HasField проверяет наличие поля в документе (lock-free) @@ -246,7 +347,7 @@ func (d *Document) Serialize() ([]byte, error) { Compressed: d.Compressed, OriginalSize: d.OriginalSize, } - docCopy.fieldsPtr.Store(fields) + docCopy.storeFields(fields) data, err := serializer.Marshal(docCopy) if err != nil { @@ -285,7 +386,7 @@ func (d *Document) Deserialize(data []byte) error { return err } d.ID = doc.ID - d.fieldsPtr.Store(doc.loadFields()) + d.storeFields(doc.loadFields()) d.CreatedAt = doc.CreatedAt d.UpdatedAt = doc.UpdatedAt d.DeletedAt = doc.DeletedAt @@ -298,7 +399,7 @@ func (d *Document) Deserialize(data []byte) error { return err } d.ID = doc.ID - d.fieldsPtr.Store(doc.loadFields()) + d.storeFields(doc.loadFields()) d.CreatedAt = doc.CreatedAt d.UpdatedAt = doc.UpdatedAt d.DeletedAt = doc.DeletedAt @@ -329,30 +430,29 @@ func (d *Document) Clone() *Document { for k, v := range fields { copiedFields[k] = deepCopyValue(v) } - clone.fieldsPtr.Store(copiedFields) + clone.storeFields(copiedFields) return clone } // Update применяет обновление к документу (атомарно, lock-free) +// ИСПРАВЛЕНО: Использует casFieldsWithBackoff для защиты от ABA и livelock func (d *Document) Update(updates map[string]interface{}) error { - for { - oldFields := d.loadFields() - newFields := make(map[string]interface{}) - for k, v := range oldFields { - newFields[k] = v - } + success := d.casFieldsWithBackoff(func(fields map[string]interface{}) map[string]interface{} { for k, v := range updates { - newFields[k] = v - } - - if d.compareAndSwapFields(oldFields, newFields) { - d.UpdatedAt = time.Now().UnixMilli() - d.Version++ - d.Compressed = false - return nil + fields[k] = v } + return fields + }) + + if success { + d.UpdatedAt = time.Now().UnixMilli() + d.Version++ + d.Compressed = false + return nil } + + return fmt.Errorf("failed to update document after multiple attempts (possible high contention)") } // SoftDelete мягко удаляет документ (устанавливает метку времени удаления) @@ -631,4 +731,4 @@ func (d *Document) SetNestedField(path string, value interface{}) error { d.UpdatedAt = time.Now().UnixMilli() d.Compressed = false return nil -} +} \ No newline at end of file