From ab48d50988dc1188bde67d061e450fa45422640d 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 20:31:37 +0000 Subject: [PATCH] Upload files to "internal/storage" --- internal/storage/runtime_limits.go | 451 ++++++ internal/storage/skiplist.go | 705 +++++++++ internal/storage/transactions.go | 2118 ++++++++++++++++++++++++++++ 3 files changed, 3274 insertions(+) create mode 100644 internal/storage/runtime_limits.go create mode 100644 internal/storage/skiplist.go create mode 100644 internal/storage/transactions.go diff --git a/internal/storage/runtime_limits.go b/internal/storage/runtime_limits.go new file mode 100644 index 0000000..9c901d4 --- /dev/null +++ b/internal/storage/runtime_limits.go @@ -0,0 +1,451 @@ +/* + * 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/runtime_limits.go +// Назначение: Ограничения на размер коллекции/документа в рантайме +// Добавлен механизм eviction при нехватке памяти (OOM protection) +// Eviction проверяет активные batch-операции перед удалением +// +// ИСПРАВЛЕНО (2026-10): после замены ACID-транзакций на batch-операции +// функция isDocInActiveTransaction ссылалась на удалённые globalTxManager +// и Transaction. Заменена на isDocInActiveBatch, который проверяет +// реестр активных batch'ей в globalBatchManager. +// +// Логика та же: не вытеснять документы, которые в данный момент +// участвуют в незакоммиченной batch-операции. Batch живёт коротко +// (создаётся, коммитится и удаляется в рамках одной функции), но +// между AddUpdate и Commit документ ещё не изменён, и eviction +// мог бы его удалить — защита сохраняется. + +package storage + +import ( + "fmt" + "runtime" + "sync" + "sync/atomic" + "time" +) + +// LoggerInterface определяет интерфейс для логирования +type LoggerInterface interface { + Debug(msg string) + Info(msg string) + Error(msg string) + Warn(msg string) +} + +// EvictionPolicy определяет политику вытеснения +type EvictionPolicy int + +const ( + EvictionNone EvictionPolicy = iota + EvictionLRU + EvictionTTL + EvictionOldest +) + +// RuntimeLimitsManager управляет runtime-ограничениями +type RuntimeLimitsManager struct { + mu sync.RWMutex + globalMaxDocSize int64 + globalMaxCollSize int64 + globalMaxDocsPerColl int64 + globalMaxMemory int64 + collectionOverrides map[string]*CollectionLimits + metrics *LimitMetrics + logger LoggerInterface + enabled bool + evictionPolicy EvictionPolicy + memoryThreshold float64 + evictionChan chan string + stopChan chan struct{} + wg sync.WaitGroup +} + +// CollectionLimits содержит лимиты для конкретной коллекции +type CollectionLimits struct { + MaxDocSize int64 + MaxCollectionSize int64 + MaxDocuments int64 + EvictionPolicy EvictionPolicy + LastUpdated int64 +} + +// LimitMetrics хранит метрики ограничений +type LimitMetrics struct { + RejectedBySize atomic.Uint64 + RejectedByDocCount atomic.Uint64 + RejectedByCollSize atomic.Uint64 + RejectedByMemory atomic.Uint64 + EvictedDocuments atomic.Uint64 + EvictedBytes atomic.Uint64 + LastCheckTime atomic.Int64 + LastEvictionTime atomic.Int64 + SkippedEviction atomic.Uint64 // Пропущено из-за активных batch-операций +} + +// RuntimeLimitsConfig содержит конфигурацию ограничений +type RuntimeLimitsConfig struct { + Enabled bool `json:"enabled"` + GlobalMaxDocSizeMB int `json:"global_max_doc_size_mb"` + GlobalMaxCollSizeMB int64 `json:"global_max_coll_size_mb"` + GlobalMaxDocsPerColl int64 `json:"global_max_docs_per_coll"` + GlobalMaxMemoryMB int64 `json:"global_max_memory_mb"` + EvictionPolicy EvictionPolicy `json:"eviction_policy"` + MemoryThreshold float64 `json:"memory_threshold"` +} + +// DefaultRuntimeLimitsConfig возвращает конфигурацию по умолчанию +func DefaultRuntimeLimitsConfig() *RuntimeLimitsConfig { + return &RuntimeLimitsConfig{ + Enabled: true, + GlobalMaxDocSizeMB: 16, + GlobalMaxCollSizeMB: 10240, + GlobalMaxDocsPerColl: 10000000, + GlobalMaxMemoryMB: 0, + EvictionPolicy: EvictionLRU, + MemoryThreshold: 0.85, + } +} + +// NewRuntimeLimitsManager создаёт новый менеджер ограничений +func NewRuntimeLimitsManager(cfg *RuntimeLimitsConfig, logger LoggerInterface) *RuntimeLimitsManager { + if cfg == nil { + cfg = DefaultRuntimeLimitsConfig() + } + + maxMemory := cfg.GlobalMaxMemoryMB * 1024 * 1024 + if maxMemory <= 0 { + var memStats runtime.MemStats + runtime.ReadMemStats(&memStats) + maxMemory = int64(float64(memStats.Sys) * 0.8) + } + + rlm := &RuntimeLimitsManager{ + globalMaxDocSize: int64(cfg.GlobalMaxDocSizeMB) * 1024 * 1024, + globalMaxCollSize: cfg.GlobalMaxCollSizeMB * 1024 * 1024, + globalMaxDocsPerColl: cfg.GlobalMaxDocsPerColl, + globalMaxMemory: maxMemory, + collectionOverrides: make(map[string]*CollectionLimits), + metrics: &LimitMetrics{}, + logger: logger, + enabled: cfg.Enabled, + evictionPolicy: cfg.EvictionPolicy, + memoryThreshold: cfg.MemoryThreshold, + evictionChan: make(chan string, 100), + stopChan: make(chan struct{}), + } + + rlm.wg.Add(1) + go rlm.memoryMonitorLoop() + + if logger != nil { + logger.Debug(fmt.Sprintf("Runtime limits manager initialized: maxDoc=%dMB, maxColl=%dMB, maxDocs=%d, maxMemory=%dMB, eviction=%d", + cfg.GlobalMaxDocSizeMB, cfg.GlobalMaxCollSizeMB, cfg.GlobalMaxDocsPerColl, maxMemory/(1024*1024), cfg.EvictionPolicy)) + } + + return rlm +} + +func (rlm *RuntimeLimitsManager) memoryMonitorLoop() { + defer rlm.wg.Done() + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + if !rlm.enabled { + continue + } + var memStats runtime.MemStats + runtime.ReadMemStats(&memStats) + currentUsage := int64(memStats.Alloc) + usageRatio := float64(currentUsage) / float64(rlm.globalMaxMemory) + + if usageRatio > rlm.memoryThreshold { + if rlm.logger != nil { + rlm.logger.Warn(fmt.Sprintf("Memory usage %.2f%% exceeds threshold %.2f%%, triggering eviction", + usageRatio*100, rlm.memoryThreshold*100)) + } + rlm.triggerEviction() + } + case <-rlm.stopChan: + return + } + } +} + +func (rlm *RuntimeLimitsManager) triggerEviction() { + select { + case rlm.evictionChan <- "global": + default: + } +} + +// isDocInActiveBatch проверяет, участвует ли документ в активной +// batch-операции. +// +// ИСПРАВЛЕНО (2026-10): ранее метод назывался isDocInActiveTransaction +// и обращался к globalTxManager и Transaction, которые были удалены при +// переходе на batch-операции. Теперь проверяем реестр активных batch'ей +// в globalBatchManager (см. transactions.go). +// +// Защита сохраняется: batch живёт коротко (создаётся и коммитится в +// рамках одной функции), но между AddUpdate и Commit документ ещё не +// изменён, и eviction мог бы его удалить. +func (rlm *RuntimeLimitsManager) isDocInActiveBatch(docID, dbName, collName string) bool { + bm := GetBatchManager() + if bm == nil { + return false + } + return bm.IsDocInActiveBatch(docID, dbName, collName) +} + +// EvictFromCollection выполняет вытеснение документов из коллекции. +// +// Проверка активных batch-операций перед удалением. +func (rlm *RuntimeLimitsManager) EvictFromCollection(coll *Collection, targetBytes int64) (int64, error) { + if !rlm.enabled { + return 0, nil + } + + evictedBytes := int64(0) + evictedCount := int64(0) + skippedCount := int64(0) + + docs := coll.GetAllDocumentsIncludingDeleted() + + // Сначала вытесняем удалённые документы + for _, doc := range docs { + if evictedBytes >= targetBytes { + break + } + + if doc.IsDeleted() { + // Проверяем активные batch-операции + if rlm.isDocInActiveBatch(doc.ID, coll.DBName(), coll.Name()) { + skippedCount++ + continue + } + + size := doc.OriginalSize + if size == 0 { + size = 1024 + } + + if err := coll.PermanentDelete(doc.ID); err == nil { + evictedBytes += size + evictedCount++ + } + } + } + + // Затем вытесняем самые старые + if evictedBytes < targetBytes && rlm.evictionPolicy == EvictionLRU { + for _, doc := range docs { + if evictedBytes >= targetBytes { + break + } + + if !doc.IsDeleted() { + // Проверяем активные batch-операции + if rlm.isDocInActiveBatch(doc.ID, coll.DBName(), coll.Name()) { + skippedCount++ + continue + } + + size := doc.OriginalSize + if size == 0 { + size = 1024 + } + + if coll.metadata.Settings.SoftDelete { + if err := coll.Delete(doc.ID); err == nil { + evictedBytes += size + evictedCount++ + } + } else { + if err := coll.PermanentDelete(doc.ID); err == nil { + evictedBytes += size + evictedCount++ + } + } + } + } + } + + rlm.metrics.EvictedDocuments.Add(uint64(evictedCount)) + rlm.metrics.EvictedBytes.Add(uint64(evictedBytes)) + if skippedCount > 0 { + rlm.metrics.SkippedEviction.Add(uint64(skippedCount)) + } + rlm.metrics.LastEvictionTime.Store(time.Now().UnixMilli()) + + if rlm.logger != nil { + rlm.logger.Info(fmt.Sprintf("Evicted %d documents (%d bytes) from collection %s.%s (skipped %d in active batches)", + evictedCount, evictedBytes, coll.DBName(), coll.Name(), skippedCount)) + } + + return evictedBytes, nil +} + +func (rlm *RuntimeLimitsManager) ValidateDocumentSize(dbName, collName string, docSize int64) error { + if !rlm.enabled { + return nil + } + + limit := rlm.getCollectionLimit(dbName, collName) + maxSize := rlm.globalMaxDocSize + if limit != nil && limit.MaxDocSize > 0 { + maxSize = limit.MaxDocSize + } + + if docSize > maxSize { + rlm.metrics.RejectedBySize.Add(1) + return fmt.Errorf("document size %d bytes exceeds limit %d bytes", docSize, maxSize) + } + return nil +} + +func (rlm *RuntimeLimitsManager) ValidateCollectionSize(coll *Collection, newDocSize int64) error { + if !rlm.enabled { + return nil + } + + limit := rlm.getCollectionLimit(coll.DBName(), coll.Name()) + maxSize := rlm.globalMaxCollSize + if limit != nil && limit.MaxCollectionSize > 0 { + maxSize = limit.MaxCollectionSize + } + + currentSize := coll.Size() + if currentSize+newDocSize > maxSize { + neededBytes := currentSize + newDocSize - maxSize + if evicted, err := rlm.EvictFromCollection(coll, neededBytes); err == nil && evicted >= neededBytes { + return nil + } + + rlm.metrics.RejectedByCollSize.Add(1) + return fmt.Errorf("collection size would exceed limit %d bytes (current: %d, new: %d)", + maxSize, currentSize, newDocSize) + } + return nil +} + +func (rlm *RuntimeLimitsManager) ValidateDocumentCount(coll *Collection) error { + if !rlm.enabled { + return nil + } + + limit := rlm.getCollectionLimit(coll.DBName(), coll.Name()) + maxDocs := rlm.globalMaxDocsPerColl + if limit != nil && limit.MaxDocuments > 0 { + maxDocs = limit.MaxDocuments + } + + currentCount := coll.Count() + if currentCount >= maxDocs { + targetBytes := int64((currentCount - maxDocs + 1) * 1024) + if evicted, err := rlm.EvictFromCollection(coll, targetBytes); err == nil && evicted > 0 { + if coll.Count() < maxDocs { + return nil + } + } + + rlm.metrics.RejectedByDocCount.Add(1) + return fmt.Errorf("collection has reached maximum document count %d", maxDocs) + } + return nil +} + +func (rlm *RuntimeLimitsManager) CheckMemoryUsage() error { + if !rlm.enabled { + return nil + } + + var memStats runtime.MemStats + runtime.ReadMemStats(&memStats) + + currentUsage := int64(memStats.Alloc) + if currentUsage > rlm.globalMaxMemory { + rlm.metrics.RejectedByMemory.Add(1) + return fmt.Errorf("memory usage %d bytes exceeds limit %d bytes", currentUsage, rlm.globalMaxMemory) + } + return nil +} + +func (rlm *RuntimeLimitsManager) getCollectionLimit(dbName, collName string) *CollectionLimits { + rlm.mu.RLock() + defer rlm.mu.RUnlock() + + key := fmt.Sprintf("%s.%s", dbName, collName) + if limits, ok := rlm.collectionOverrides[key]; ok { + return limits + } + return nil +} + +func (rlm *RuntimeLimitsManager) SetCollectionLimits(dbName, collName string, maxDocSizeMB int, maxCollSizeMB int64, maxDocuments int64) { + rlm.mu.Lock() + defer rlm.mu.Unlock() + + key := fmt.Sprintf("%s.%s", dbName, collName) + rlm.collectionOverrides[key] = &CollectionLimits{ + MaxDocSize: int64(maxDocSizeMB) * 1024 * 1024, + MaxCollectionSize: maxCollSizeMB * 1024 * 1024, + MaxDocuments: maxDocuments, + EvictionPolicy: rlm.evictionPolicy, + LastUpdated: time.Now().UnixMilli(), + } + + if rlm.logger != nil { + rlm.logger.Info(fmt.Sprintf("Set limits for %s: maxDoc=%dMB, maxColl=%dMB, maxDocs=%d", + key, maxDocSizeMB, maxCollSizeMB, maxDocuments)) + } +} + +func (rlm *RuntimeLimitsManager) RemoveCollectionLimits(dbName, collName string) { + rlm.mu.Lock() + defer rlm.mu.Unlock() + + key := fmt.Sprintf("%s.%s", dbName, collName) + delete(rlm.collectionOverrides, key) + + if rlm.logger != nil { + rlm.logger.Info(fmt.Sprintf("Removed limits override for %s", key)) + } +} + +func (rlm *RuntimeLimitsManager) GetMetrics() map[string]interface{} { + return map[string]interface{}{ + "rejected_by_size": rlm.metrics.RejectedBySize.Load(), + "rejected_by_doc_count": rlm.metrics.RejectedByDocCount.Load(), + "rejected_by_coll_size": rlm.metrics.RejectedByCollSize.Load(), + "rejected_by_memory": rlm.metrics.RejectedByMemory.Load(), + "evicted_documents": rlm.metrics.EvictedDocuments.Load(), + "evicted_bytes": rlm.metrics.EvictedBytes.Load(), + "skipped_eviction": rlm.metrics.SkippedEviction.Load(), + "last_eviction_time": rlm.metrics.LastEvictionTime.Load(), + "global_max_doc_size_mb": rlm.globalMaxDocSize / (1024 * 1024), + "global_max_coll_size_mb": rlm.globalMaxCollSize / (1024 * 1024), + "global_max_docs_per_coll": rlm.globalMaxDocsPerColl, + "global_max_memory_mb": rlm.globalMaxMemory / (1024 * 1024), + "eviction_policy": rlm.evictionPolicy, + "memory_threshold": rlm.memoryThreshold, + "enabled": rlm.enabled, + } +} + +func (rlm *RuntimeLimitsManager) Stop() { + close(rlm.stopChan) + rlm.wg.Wait() +} diff --git a/internal/storage/skiplist.go b/internal/storage/skiplist.go new file mode 100644 index 0000000..a0da781 --- /dev/null +++ b/internal/storage/skiplist.go @@ -0,0 +1,705 @@ +/* + * 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/skiplist.go +// Назначение: Lock-free skip list (список с пропусками) для инвертированного +// индекса коллекции. Заменяет sync.Map в Index.data. +// +// ИСПРАВЛЕНО (2026-10): +// - Баг #6: рассинхронизация size при воскрешении узла (Insert после Delete). +// Теперь проверяется wasDeleted и size увеличивается корректно. +// - Баг #7: дубликат на верхних уровнях при конкурентной вставке одного +// ключа с разными newLevel. Решение: двухфазная вставка — сначала +// уровень 0 (гарантированно единственный узел с таким ключом), затем +// верхние уровни (best-effort, дубликаты безопасны). +// +// КЛЮЧЕВОЙ ИНВАРИАНТ: +// - На уровне 0 не может быть двух узлов с одинаковым ключом. +// - На уровнях > 0 дубликаты допустимы (они лишь замедляют поиск, но не +// ломают корректность, потому что Search всегда спускается до уровня 0 +// и там находит единственный узел с нужным ключом). +// +// ДРУГИЕ ИСПРАВЛЕНИЯ (2026): +// - Устранён data race при обновлении value существующего узла: +// value хранится в atomic.Value. +// - Исправлен обход Range: спуск на базовый уровень теперь корректный. +// - Исправлена вставка на уровнях выше текущего sl.level. +// - Идемпотентный Delete. +// - RangeWithKeyPrefix корректно останавливается при выходе за префикс. + +package storage + +import ( + "fmt" + "math/rand" + "strings" + "sync" + "sync/atomic" +) + +// ============================================================================= +// КОНСТАНТЫ +// ============================================================================= + +const ( + // Максимальное количество уровней в skip list. + maxSkipListLevel = 32 + + // Вероятность "подъёма" на следующий уровень. + skipListProbability = 0.5 + + // Порог, после которого Range запускает физическую очистку + // помеченных (tombstone) узлов. + skipListSweepThreshold = 1024 +) + +// ============================================================================= +// ТИПЫ +// ============================================================================= + +// skipListNode представляет узел skip list. +type skipListNode struct { + // key — значение индексированного поля. + key interface{} + + // value — ID документа. Хранится в atomic.Value для защиты от + // data race при обновлении существующего ключа. + // atomic.Value не может хранить nil строку напрямую, поэтому + // используется обёртка *string. + value atomic.Value // *string + + // deleted — tombstone-флаг. + deleted atomic.Bool + + // forward — массив атомарных указателей на следующий узел + // на каждом уровне. + forward []atomic.Value // *skipListNode +} + +// loadValue атомарно загружает value. +func (n *skipListNode) loadValue() string { + val := n.value.Load() + if val == nil { + return "" + } + ptr, ok := val.(*string) + if !ok || ptr == nil { + return "" + } + return *ptr +} + +// storeValue атомарно сохраняет value. +func (n *skipListNode) storeValue(v string) { + n.value.Store(&v) +} + +// loadForward атомарно загружает указатель на следующий узел на уровне. +func (n *skipListNode) loadForward(level int) *skipListNode { + if level < 0 || level >= len(n.forward) { + return nil + } + val := n.forward[level].Load() + if val == nil { + return nil + } + ptr, ok := val.(*skipListNode) + if !ok { + return nil + } + return ptr +} + +// storeForward атомарно сохраняет указатель на следующий узел на уровне. +func (n *skipListNode) storeForward(level int, next *skipListNode) { + if level < 0 || level >= len(n.forward) { + return + } + n.forward[level].Store(next) +} + +// casForward атомарно заменяет указатель, если он равен old. +func (n *skipListNode) casForward(level int, old, new *skipListNode) bool { + if level < 0 || level >= len(n.forward) { + return false + } + return n.forward[level].CompareAndSwap(old, new) +} + +// SkipList — lock-free список с пропусками. +type SkipList struct { + head *skipListNode + level atomic.Int32 + size atomic.Int64 + + rngMu sync.Mutex + rng *rand.Rand +} + +// NewSkipList создаёт новый пустой skip list. +func NewSkipList() *SkipList { + head := &skipListNode{ + key: nil, + forward: make([]atomic.Value, maxSkipListLevel), + } + // Инициализируем forward-указатели типизированным nil. + for i := 0; i < maxSkipListLevel; i++ { + head.forward[i].Store((*skipListNode)(nil)) + } + // Инициализируем value head-а, чтобы atomic.Value имел + // установленный concrete type. + emptyStr := "" + head.value.Store(&emptyStr) + + sl := &SkipList{ + head: head, + rng: rand.New(rand.NewSource(int64(0x5eed5eed))), + } + sl.level.Store(0) + sl.size.Store(0) + return sl +} + +// randomLevel генерирует случайный уровень для нового узла. +func (sl *SkipList) randomLevel() int { + sl.rngMu.Lock() + defer sl.rngMu.Unlock() + + level := 0 + for level < maxSkipListLevel-1 && sl.rng.Float64() < skipListProbability { + level++ + } + return level +} + +// ============================================================================= +// ПОИСК ПРЕДШЕСТВЕННИКОВ (вспомогательный метод) +// ============================================================================= + +// findPredecessors заполняет update[] указателями на узлы, после которых +// нужно вставлять/искать узел с заданным ключом на каждом уровне. +// +// Возвращает: +// - update[] — массив предшественников (длина maxSkipListLevel) +// - existing — узел на уровне 0 с ключом == key, либо nil +// +// Этот метод используется и в Insert, и в Delete. +func (sl *SkipList) findPredecessors(key interface{}) ([]*skipListNode, *skipListNode) { + update := make([]*skipListNode, maxSkipListLevel) + + currentLevel := int(sl.level.Load()) + + // Спускаемся с верхнего уровня до базового. + x := sl.head + for i := currentLevel; i >= 0; i-- { + for { + next := x.loadForward(i) + if next == nil { + break + } + if next.deleted.Load() { + // Пытаемся физически вытолкнуть tombstone. + if x.casForward(i, next, next.loadForward(i)) { + continue + } + continue + } + if compareSkipListKeys(next.key, key) < 0 { + x = next + continue + } + break + } + update[i] = x + } + + // Для уровней выше currentLevel предшественник — head. + // (Это может понадобиться, если новый уровень больше текущего.) + for i := currentLevel + 1; i < maxSkipListLevel; i++ { + update[i] = sl.head + } + + // Ищем существующий узел на уровне 0. + existing := update[0].loadForward(0) + for existing != nil && compareSkipListKeys(existing.key, key) < 0 { + existing = existing.loadForward(0) + } + if existing != nil && compareSkipListKeys(existing.key, key) != 0 { + existing = nil + } + + return update, existing +} + +// ============================================================================= +// ВСТАВКА (ИСПРАВЛЕНО: баги #6 и #7) +// ============================================================================= + +// Insert добавляет пару (key, value) в skip list. +// Если ключ уже существует, значение перезаписывается. +// +// РЕАЛИЗАЦИЯ (двухфазная вставка): +// +// Фаза 1: вставка на уровень 0. +// - Если ключ уже есть на уровне 0 — обновляем value (с корректной +// обработкой wasDeleted для size). +// - Иначе — CAS-вставка нового узла на уровень 0. Если CAS провалился +// из-за конкурента — retry. +// - После успешной вставки на уровень 0 ключ ГАРАНТИРОВАННО есть в списке +// (в единственном экземпляре на уровне 0). +// +// Фаза 2: вставка на верхние уровни (best-effort). +// - Для каждого уровня i от 1 до newLevel пытаемся CAS-вставить узел. +// - Если CAS провалился — не страшно: узел уже есть на уровне 0, +// а верхние уровни — только ускорители поиска. Дубликаты на верхних +// уровнях безопасны, потому что Search всегда спускается до уровня 0. +func (sl *SkipList) Insert(key interface{}, value string) { + if key == nil { + return + } + + // Внешний цикл — retry при провале CAS на уровне 0. + for attempt := 0; attempt < 1000; attempt++ { + newLevel := sl.randomLevel() + + // Поднимаем максимальный уровень, если нужно. + for { + currentLevel := int(sl.level.Load()) + if newLevel <= currentLevel { + break + } + if sl.level.CompareAndSwap(int32(currentLevel), int32(newLevel)) { + break + } + } + + // Находим предшественников и существующий узел. + update, existing := sl.findPredecessors(key) + + // Случай 1: ключ уже есть на уровне 0. + // Обновляем value и, если узел был tombstone, воскрешаем его. + if existing != nil { + wasDeleted := existing.deleted.Load() + existing.storeValue(value) + existing.deleted.Store(false) + + // ИСПРАВЛЕНО (баг #6): если узел был tombstone, + // увеличиваем size, потому что Delete его уже уменьшил. + if wasDeleted { + sl.size.Add(1) + } + return + } + + // Случай 2: ключа нет — создаём новый узел и вставляем + // сначала на уровень 0. + newNode := &skipListNode{ + key: key, + forward: make([]atomic.Value, newLevel+1), + } + for i := 0; i <= newLevel; i++ { + newNode.forward[i].Store((*skipListNode)(nil)) + } + newNode.storeValue(value) + + // Пытаемся вставить на уровень 0. + next0 := update[0].loadForward(0) + // Защита от ситуации, когда между findPredecessors и сейчас + // кто-то вставил узел с таким же ключом. + for next0 != nil && compareSkipListKeys(next0.key, key) < 0 { + next0 = next0.loadForward(0) + } + + // Если между findPredecessors и сейчас появился узел с нашим + // ключом — обновляем его value и выходим. + if next0 != nil && compareSkipListKeys(next0.key, key) == 0 { + wasDeleted := next0.deleted.Load() + next0.storeValue(value) + next0.deleted.Store(false) + if wasDeleted { + sl.size.Add(1) + } + return + } + + newNode.storeForward(0, next0) + if !update[0].casForward(0, next0, newNode) { + // CAS провалился — кто-то другой вставил узел между + // findPredecessors и сейчас. Retry. + continue + } + + // Узел успешно вставлен на уровень 0. Теперь увеличиваем size + // (единственное место, где мы это делаем для новых узлов). + sl.size.Add(1) + + // ============================================================ + // ФАЗА 2: вставка на верхние уровни (best-effort). + // ============================================================ + // + // Дубликаты на верхних уровнях безопасны: Search всегда спускается + // до уровня 0 и там находит единственный узел с нужным ключом. + // Если CAS на верхнем уровне провалился — просто пропускаем уровень. + for i := 1; i <= newLevel; i++ { + if update[i] == nil { + // Предшественник не определён — пропускаем уровень. + continue + } + // Перечитываем next на этом уровне (мог измениться). + next := update[i].loadForward(i) + newNode.storeForward(i, next) + // CAS best-effort: если провалился — не retry, просто пропускаем. + update[i].casForward(i, next, newNode) + } + + return + } + + // Если за 1000 попыток не удалось — что-то не так. + // Молча выходим (не паникуем). +} + +// ============================================================================= +// ПОИСК +// ============================================================================= + +// Search ищет значение по ключу. +func (sl *SkipList) Search(key interface{}) (string, bool) { + if key == nil { + return "", false + } + + x := sl.head + for i := int(sl.level.Load()); i >= 0; i-- { + for { + next := x.loadForward(i) + if next == nil { + break + } + if next.deleted.Load() { + // Пропускаем tombstone. + if x.casForward(i, next, next.loadForward(i)) { + continue + } + continue + } + if compareSkipListKeys(next.key, key) < 0 { + x = next + continue + } + break + } + } + + candidate := x.loadForward(0) + if candidate == nil { + return "", false + } + if compareSkipListKeys(candidate.key, key) == 0 { + if candidate.deleted.Load() { + return "", false + } + return candidate.loadValue(), true + } + return "", false +} + +// SearchAll возвращает все значения, ключ которых равен заданному. +// +// ВАЖНО: текущая реализация индекса хранит только один docID на ключ, +// поэтому метод возвращает 0 или 1 значение. Для полноценной поддержки +// неуникальных индексов потребуется изменить структуру узла. +func (sl *SkipList) SearchAll(key interface{}) []string { + val, ok := sl.Search(key) + if !ok { + return nil + } + return []string{val} +} + +// ============================================================================= +// УДАЛЕНИЕ +// ============================================================================= + +// Delete помечает узел с заданным ключом как удалённый (tombstone). +// Идемпотентен: повторный вызов не уменьшает size. +func (sl *SkipList) Delete(key interface{}) { + if key == nil { + return + } + + _, existing := sl.findPredecessors(key) + if existing == nil { + return + } + + // CAS как атомарная проверка+установка. + // Уменьшаем size только если мы первыми поставили tombstone. + if existing.deleted.CompareAndSwap(false, true) { + sl.size.Add(-1) + } +} + +// ============================================================================= +// ОБХОД (RANGE) +// ============================================================================= + +// Range вызывает fn для каждой пары (key, value) в порядке возрастания ключей. +// Если fn возвращает false — обход прекращается. +func (sl *SkipList) Range(fn func(key interface{}, value string) bool) { + prev := sl.head + tombstoneCount := 0 + + for { + next := prev.loadForward(0) + if next == nil { + break + } + + if next.deleted.Load() { + tombstoneCount++ + // Пытаемся физически удалить узел из базового уровня. + if prev.casForward(0, next, next.loadForward(0)) { + if tombstoneCount >= skipListSweepThreshold { + sl.sweepUpperLevels() + tombstoneCount = 0 + } + continue + } + // CAS не удался — двигаемся дальше. + prev = next + continue + } + + if !fn(next.key, next.loadValue()) { + return + } + + prev = next + } +} + +// RangeWithKeyPrefix вызывает fn для каждой пары, ключ которой является +// строкой с заданным префиксом. +// +// Благодаря сортированности skip list, обход можно прервать, как только +// ключ перестал соответствовать префиксу. +func (sl *SkipList) RangeWithKeyPrefix(prefix string, fn func(key interface{}, value string) bool) { + sl.Range(func(key interface{}, value string) bool { + keyStr, ok := key.(string) + if !ok { + return true + } + if strings.HasPrefix(keyStr, prefix) { + return fn(key, value) + } + // Если ключ — строка и он больше префикса, дальнейшие ключи + // (отсортированные) тоже не подойдут. + if keyStr > prefix { + return false + } + return true + }) +} + +// sweepUpperLevels проходит по верхним уровням и физически выталкивает +// tombstone-узлы. +func (sl *SkipList) sweepUpperLevels() { + for i := 1; i < maxSkipListLevel; i++ { + x := sl.head + for { + next := x.loadForward(i) + if next == nil { + break + } + if next.deleted.Load() { + if x.casForward(i, next, next.loadForward(i)) { + continue + } + } + x = next + } + } +} + +// ============================================================================= +// СЛУЖЕБНЫЕ МЕТОДЫ +// ============================================================================= + +// Size возвращает приблизительное количество активных элементов. +func (sl *SkipList) Size() int64 { return sl.size.Load() } + +// Len — синоним Size. +func (sl *SkipList) Len() int64 { return sl.size.Load() } + +// IsEmpty возвращает true, если список пуст. +func (sl *SkipList) IsEmpty() bool { return sl.size.Load() == 0 } + +// Clear помечает все узлы как удалённые. +func (sl *SkipList) Clear() { + sl.Range(func(key interface{}, value string) bool { + sl.Delete(key) + return true + }) + sl.size.Store(0) +} + +// ============================================================================= +// СРАВНЕНИЕ КЛЮЧЕЙ +// ============================================================================= + +// compareSkipListKeys сравнивает два ключа. +func compareSkipListKeys(a, b interface{}) int { + if a == nil && b == nil { + return 0 + } + if a == nil { + return -1 + } + if b == nil { + return 1 + } + + switch av := a.(type) { + case string: + if bv, ok := b.(string); ok { + return strings.Compare(av, bv) + } + case int: + if bv, ok := toInt64(b); ok { + return compareInt64(int64(av), bv) + } + case int64: + if bv, ok := toInt64(b); ok { + return compareInt64(av, bv) + } + case float64: + if bv, ok := toFloat64Comparable(b); ok { + if av < bv { + return -1 + } + if av > bv { + return 1 + } + return 0 + } + case bool: + if bv, ok := b.(bool); ok { + if av == bv { + return 0 + } + if !av { + return -1 + } + return 1 + } + case []byte: + if bv, ok := b.([]byte); ok { + return strings.Compare(string(av), string(bv)) + } + } + + return strings.Compare(fmt.Sprintf("%v", a), fmt.Sprintf("%v", b)) +} + +func toInt64(v interface{}) (int64, bool) { + switch x := v.(type) { + case int: + return int64(x), true + case int32: + return int64(x), true + case int64: + return x, true + case uint: + return int64(x), true + case uint32: + return int64(x), true + case uint64: + return int64(x), true + case float64: + return int64(x), true + case float32: + return int64(x), true + default: + return 0, false + } +} + +func toFloat64Comparable(v interface{}) (float64, bool) { + switch x := v.(type) { + case int: + return float64(x), true + case int32: + return float64(x), true + case int64: + return float64(x), true + case float32: + return float64(x), true + case float64: + return x, true + default: + return 0, false + } +} + +func compareInt64(a, b int64) int { + if a < b { + return -1 + } + if a > b { + return 1 + } + return 0 +} + +// ============================================================================= +// ITERATOR +// ============================================================================= + +// SkipListIterator — итератор по skip list. +type SkipListIterator struct { + sl *SkipList + current *skipListNode + started bool +} + +// NewIterator создаёт новый итератор. +func (sl *SkipList) NewIterator() *SkipListIterator { + return &SkipListIterator{sl: sl} +} + +// Next переходит к следующему элементу. +func (it *SkipListIterator) Next() (interface{}, string, bool) { + if it.sl == nil { + return nil, "", false + } + + if !it.started { + it.started = true + it.current = it.sl.head + } + + for { + next := it.current.loadForward(0) + if next == nil { + return nil, "", false + } + it.current = next + if !next.deleted.Load() { + return next.key, next.loadValue(), true + } + } +} + +// Close освобождает ресурсы итератора. +func (it *SkipListIterator) Close() { + it.sl = nil + it.current = nil +} diff --git a/internal/storage/transactions.go b/internal/storage/transactions.go new file mode 100644 index 0000000..15a23d4 --- /dev/null +++ b/internal/storage/transactions.go @@ -0,0 +1,2118 @@ +/* + * 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/transactions.go +// Назначение: Batch-операции (замена полноценных ACID-транзакций) + +// Time-travel queries на основе WAL. +// +// АРХИТЕКТУРНЫЕ ИЗМЕНЕНИЯ (2026): +// - Удалены все мьютексы и MVCC (MVCCManager, VisibilityMap, +// ReadTimestampCache, DocumentVersion). Они давали огромную сложность +// без реальной пользы для in-memory СУБД. +// - Полноценные ACID-транзакции заменены на batch-операции. +// Batch — это группа операций, которые применяются атомарно: +// либо все, либо ни одна. Гарантия атомарности обеспечивается +// через WAL: сначала пишем весь batch в WAL, потом применяем. +// - SAGA-шаги используют batch для гарантии атомарности шага. +// - Time-travel реализован через WAL: для получения версии документа +// на момент времени проигрываем WAL до нужного LSN. +// - Удалён AsyncRecoveryManager (был источником race conditions). +// Восстановление теперь синхронное при старте. +// +// ИСПРАВЛЕНО (2026-10): добавлен реестр activeBatches в BatchManager, +// который используется eviction-механизмом (runtime_limits.go), чтобы +// не вытеснять документы, участвующие в незакоммиченной batch-операции. +// Ранее runtime_limits.go ссылался на удалённые globalTxManager и +// Transaction — теперь эти ссылки заменены на GetBatchManager() и *Batch. + +package storage + +import ( + "bufio" + "encoding/binary" + "encoding/json" + "fmt" + "io" + "os" + "path/filepath" + "sort" + "sync" + "sync/atomic" + "time" + + "futriis/internal/config" +) + +// ============================================================================= +// БАЗОВЫЕ ТИПЫ +// ============================================================================= + +type TransactionID uint64 + +// BatchOperation — одна операция в batch. +type BatchOperation struct { + Type string `json:"type"` // insert, update, delete, restore + Database string `json:"database"` + Collection string `json:"collection"` + DocumentID string `json:"document_id"` + Data map[string]interface{} `json:"data,omitempty"` + Version uint64 `json:"version,omitempty"` +} + +// Batch представляет группу операций, применяемых атомарно. +type Batch struct { + ID TransactionID `json:"id"` + Operations []BatchOperation `json:"operations"` + Timestamp int64 `json:"timestamp"` + Database string `json:"database,omitempty"` + Collection string `json:"collection,omitempty"` + LSN uint64 `json:"lsn,omitempty"` +} + +// WALRecord — запись в WAL. +type WALRecord struct { + CRC uint32 `json:"crc"` + Length uint32 `json:"length"` + Type byte `json:"type"` // 1 = batch + Data []byte `json:"data"` + Timestamp int64 `json:"timestamp"` + LSN uint64 `json:"lsn"` +} + +// ============================================================================= +// КОНСТАНТЫ +// ============================================================================= + +const ( + WALSegmentSize = 64 * 1024 * 1024 + WALSegmentPrefix = "wal_segment_" + WALIndexPrefix = "wal_index_" + + VersionPruneInterval = 5 * time.Minute + + SagaStateDir = "saga_states" + SagaStateFilePrefix = "saga_state_" + SagaStateFileSuffix = ".json" + SagaCheckpointInterval = 30 * time.Second + SagaMaxRetries = 5 + SagaRetryBackoff = 100 * time.Millisecond + + SagaOrphanTimeout = 5 * time.Minute + SagaCleanupInterval = 1 * time.Minute +) + +// ============================================================================= +// CRC32 +// ============================================================================= + +var crc32Table = [256]uint32{ + 0x00000000, 0x77073096, 0xee0e612c, 0x990951ba, 0x076dc419, 0x706af48f, + 0xe963a535, 0x9e6495a3, 0x0edb8832, 0x79dcb8a4, 0xe0d5e91e, 0x97d2d988, + 0x09b64c2b, 0x7eb17cbd, 0xe7b82d07, 0x90bf1d91, 0x1db71064, 0x6ab020f2, + 0xf3b97148, 0x84be41de, 0x1adad47d, 0x6ddde4eb, 0xf4d4b551, 0x83d385c7, + 0x136c9856, 0x646ba8c0, 0xfd62f97a, 0x8a65c9ec, 0x14015c4f, 0x63066cd9, + 0xfa0f3d63, 0x8d080df5, 0x3b6e20c8, 0x4c69105e, 0xd56041e4, 0xa2677172, + 0x3c03e4d1, 0x4b04d447, 0xd20d85fd, 0xa50ab56b, 0x35b5a8fa, 0x42b2986c, + 0xdbbbc9d6, 0xacbcf940, 0x32d86ce3, 0x45df5c75, 0xdcd60dcf, 0xabd13d59, + 0x26d930ac, 0x51de003a, 0xc8d75180, 0xbfd06116, 0x21b4f4b5, 0x56b3c423, + 0xcfba9599, 0xb8bda50f, 0x2802b89e, 0x5f058808, 0xc60cd9b2, 0xb10be924, + 0x2f6f7c87, 0x58684c11, 0xc1611dab, 0xb6662d3d, 0x76dc4190, 0x01db7106, + 0x98d220bc, 0xefd5102a, 0x71b18589, 0x06b6b51f, 0x9fbfe4a5, 0xe8b8d433, + 0x7807c9a2, 0x0f00f934, 0x9609a88e, 0xe10e9818, 0x7f6a0dbb, 0x086d3d2d, + 0x91646c97, 0xe6635c01, 0x6b6b51f4, 0x1c6c6162, 0x856530d8, 0xf262004e, + 0x6c0695ed, 0x1b01a57b, 0x8208f4c1, 0xf50fc457, 0x65b0d9c6, 0x12b7e950, + 0x8bbeb8ea, 0xfcb9887c, 0x62dd1ddf, 0x15da2d49, 0x8cd37cf3, 0xfbd44c65, + 0x4db26158, 0x3ab551ce, 0xa3bc0074, 0xd4bb30e2, 0x4adfa541, 0x3dd895d7, + 0xa4d1c46d, 0xd3d6f4fb, 0x4369e96a, 0x346ed9fc, 0xad678846, 0xda60b8d0, + 0x44042d73, 0x33031de5, 0xaa0a4c5f, 0xdd0d7cc9, 0x5005713c, 0x270241aa, + 0xbe0b1010, 0xc90c2086, 0x5768b525, 0x206f85b3, 0xb966d409, 0xce61e49f, + 0x5edef90e, 0x29d9c998, 0xb0d09822, 0xc7d7a8b4, 0x59b33d17, 0x2eb40d81, + 0xb7bd5c3b, 0xc0ba6cad, 0xedb88320, 0x9abfb3b6, 0x03b6e20c, 0x74b1d29a, + 0xead54739, 0x9dd277af, 0x04db2615, 0x73dc1683, 0xe3630b12, 0x94643b84, + 0x0d6d6a3e, 0x7a6a5aa8, 0xe40ecf0b, 0x9309ff9d, 0x0a00ae27, 0x7d079eb1, + 0xf00f9344, 0x8708a3d2, 0x1e01f268, 0x6906c2fe, 0xf762575d, 0x806567cb, + 0x196c3671, 0x6e6b06e7, 0xfed41b76, 0x89d32be0, 0x10da7a5a, 0x67dd4acc, + 0xf9b9df6f, 0x8ebeeff9, 0x17b7be43, 0x60b08ed5, 0xd6d6a3e8, 0xa1d1937e, + 0x38d8c2c4, 0x4fdff252, 0xd1bb67f1, 0xa6bc5767, 0x3fb506dd, 0x48b2364b, + 0xd80d2bda, 0xaf0a1a4c, 0x36034af6, 0x41047a60, 0xdf60efc3, 0xa867df55, + 0x316e8eef, 0x4669be79, 0xcb61b38c, 0xbc66831a, 0x256fd2a0, 0x5268e236, + 0xcc0c7795, 0xbb0b4703, 0x220216b9, 0x5505262f, 0xc5ba3bbe, 0xb2bd0b28, + 0x2bb45a92, 0x5cb36a04, 0xc2d7ffa7, 0xb5d0cf31, 0x2cd99e8b, 0x5bdeae1d, + 0x9b64c2b0, 0xec63f226, 0x756aa39c, 0x026d930a, 0x9c0906a9, 0xeb0e363f, + 0x72076785, 0x05005713, 0x95bf4a82, 0xe2b87a14, 0x7bb12bae, 0x0cb61b38, + 0x92d28e9b, 0xe5d5be0d, 0x7cdcefb7, 0x0bdbdf21, 0x86d3d2d4, 0xf1d4e242, + 0x68ddb3f8, 0x1fda836e, 0x81be16cd, 0xf6b9265b, 0x6fb077e1, 0x18b74777, + 0x88085ae6, 0xff0f6a70, 0x66063bca, 0x11010b5c, 0x8f659eff, 0xf862ae69, + 0x616bffd3, 0x166ccf45, 0xa00ae278, 0xd70dd2ee, 0x4e048354, 0x3903b3c2, + 0xa7672661, 0xd06016f7, 0x4969474d, 0x3e6e77db, 0xaed16a4a, 0xd9d65adc, + 0x40df0b66, 0x37d83bf8, 0xa9bcae53, 0xdebb9ec5, 0x47b2cf7f, 0x30b5ffe9, + 0xbdbdf21c, 0xcabac28a, 0x53b39330, 0x24b4a3a6, 0xbad03605, 0xcdd70693, + 0x54de5729, 0x23d967bf, 0xb3667a2e, 0xc4614ab8, 0x5d681b02, 0x2a6f2b94, + 0xb40bbe37, 0xc30c8ea1, 0x5a05df1b, 0x2d02ef8d, +} + +func crc32(data []byte) uint32 { + crc := uint32(0xFFFFFFFF) + for _, b := range data { + crc = (crc >> 8) ^ crc32Table[(crc^uint32(b))&0xFF] + } + return crc ^ 0xFFFFFFFF +} + +// ============================================================================= +// BATCH MANAGER - замена TransactionManager +// ============================================================================= + +// BatchManager управляет batch-операциями и WAL. +type BatchManager struct { + nextBatchID atomic.Uint64 + wal *SegmentedWALManager + logger LoggerInterface + walPath string + stats *BatchStats + aof *AOFManager // append-only file для мутаций вне batch + + // activeBatches — реестр активных (незакоммиченных) batch'ей. + // Ключ — Batch.ID, значение — *Batch. + // Используется eviction-механизмом (runtime_limits.go), чтобы не + // вытеснять документы, участвующие в незакоммиченной batch-операции. + activeBatches sync.Map +} + +// BatchStats — статистика batch-операций. +type BatchStats struct { + TotalBatches atomic.Uint64 + TotalOperations atomic.Uint64 + TotalFailed atomic.Uint64 + StartTime time.Time +} + +// ============================================================================= +// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ +// ============================================================================= + +var ( + globalBatchManager *BatchManager + batchManagerOnce sync.Once + + globalStorage *Storage +) + +// InitBatchManager инициализирует batch-менеджер. +func InitBatchManager(walPath string) error { + return InitBatchManagerWithConfig(walPath, nil) +} + +// InitBatchManagerWithConfig инициализирует batch-менеджер с конфигурацией. +func InitBatchManagerWithConfig(walPath string, config map[string]interface{}) error { + var err error + batchManagerOnce.Do(func() { + fsyncEnabled := true + if config != nil { + if v, ok := config["fsync_enabled"].(bool); ok { + fsyncEnabled = v + } + } + + globalBatchManager = &BatchManager{ + walPath: walPath, + stats: &BatchStats{StartTime: time.Now()}, + } + globalBatchManager.nextBatchID.Store(1) + + globalBatchManager.wal, err = NewSegmentedWALManager( + filepath.Dir(walPath), fsyncEnabled, nil) + if err != nil { + return + } + + // Инициализируем AOF для мутаций вне batch. + globalBatchManager.aof, err = NewAOFManager( + filepath.Join(filepath.Dir(walPath), "aof"), + fsyncEnabled, nil) + if err != nil { + return + } + }) + return err +} + +// GetBatchManager возвращает глобальный batch-менеджер. +func GetBatchManager() *BatchManager { return globalBatchManager } + +// SetBatchLogger устанавливает логгер. +func SetBatchLogger(logger LoggerInterface) { + if globalBatchManager != nil { + globalBatchManager.logger = logger + if globalBatchManager.wal != nil { + globalBatchManager.wal.logger = logger + } + if globalBatchManager.aof != nil { + globalBatchManager.aof.SetLogger(logger) + } + } +} + +// SetGlobalStorage устанавливает глобальное хранилище. +func SetGlobalStorage(s *Storage) { globalStorage = s } + +// GetGlobalStorage возвращает глобальное хранилище. +func GetGlobalStorage() *Storage { return globalStorage } + +// ============================================================================= +// РЕЕСТР АКТИВНЫХ BATCH'ЕЙ +// ============================================================================= + +// RegisterBatch регистрирует batch как активный. +// Вызывается в NewBatch. +func (bm *BatchManager) RegisterBatch(b *Batch) { + if bm == nil || b == nil { + return + } + bm.activeBatches.Store(b.ID, b) +} + +// UnregisterBatch снимает регистрацию batch. +// Вызывается в Commit (успешном или неуспешном). +func (bm *BatchManager) UnregisterBatch(b *Batch) { + if bm == nil || b == nil { + return + } + bm.activeBatches.Delete(b.ID) +} + +// IsDocInActiveBatch проверяет, участвует ли документ в активной +// (незакоммиченной) batch-операции. +// +// Используется в runtime_limits.go (eviction), чтобы не вытеснять +// документы, которые в данный момент участвуют в batch. +func (bm *BatchManager) IsDocInActiveBatch(docID, dbName, collName string) bool { + if bm == nil { + return false + } + found := false + bm.activeBatches.Range(func(key, value interface{}) bool { + b, ok := value.(*Batch) + if !ok || b == nil { + return true + } + for _, op := range b.Operations { + if op.DocumentID == docID && + op.Database == dbName && + op.Collection == collName { + found = true + return false + } + } + return true + }) + return found +} + +// GetActiveBatchCount возвращает количество активных batch'ей. +// Полезно для метрик и отладки. +func (bm *BatchManager) GetActiveBatchCount() int { + if bm == nil { + return 0 + } + count := 0 + bm.activeBatches.Range(func(_, _ interface{}) bool { + count++ + return true + }) + return count +} + +// ============================================================================= +// BATCH API +// ============================================================================= + +// NewBatch создаёт новый batch. +func NewBatch(database, collection string) *Batch { + if globalBatchManager == nil { + _ = InitBatchManager("futriis.wal") + } + b := &Batch{ + ID: TransactionID(globalBatchManager.nextBatchID.Add(1) - 1), + Operations: make([]BatchOperation, 0, 16), + Timestamp: time.Now().UnixMilli(), + Database: database, + Collection: collection, + } + // Регистрируем batch как активный, чтобы eviction не вытеснил + // документы, участвующие в нём. + globalBatchManager.RegisterBatch(b) + return b +} + +// AddInsert добавляет операцию вставки. +func (b *Batch) AddInsert(doc *Document) { + b.Operations = append(b.Operations, BatchOperation{ + Type: "insert", + Database: b.Database, + Collection: b.Collection, + DocumentID: doc.ID, + Data: doc.GetFields(), + Version: doc.Version, + }) +} + +// AddUpdate добавляет операцию обновления. +func (b *Batch) AddUpdate(docID string, updates map[string]interface{}) { + b.Operations = append(b.Operations, BatchOperation{ + Type: "update", + Database: b.Database, + Collection: b.Collection, + DocumentID: docID, + Data: updates, + }) +} + +// AddDelete добавляет операцию удаления. +func (b *Batch) AddDelete(docID string) { + b.Operations = append(b.Operations, BatchOperation{ + Type: "delete", + Database: b.Database, + Collection: b.Collection, + DocumentID: docID, + }) +} + +// AddRestore добавляет операцию восстановления. +func (b *Batch) AddRestore(docID string) { + b.Operations = append(b.Operations, BatchOperation{ + Type: "restore", + Database: b.Database, + Collection: b.Collection, + DocumentID: docID, + }) +} + +// Commit применяет batch атомарно. +// +// Гарантия атомарности: +// 1. Сериализуем весь batch и пишем в WAL (fsync). +// 2. Только после успешной записи в WAL применяем операции. +// 3. Если применение падает на середине — мы можем восстановить +// состояние, проиграв WAL (операции идемпотентны или +// откатываются вручную). +// +// В любом случае (успех или ошибка) batch снимается с регистрации +// в activeBatches. +func (b *Batch) Commit() error { + if globalBatchManager == nil { + return fmt.Errorf("batch manager not initialized") + } + + // Снимаем регистрацию в любом случае — batch короткоживущий. + defer globalBatchManager.UnregisterBatch(b) + + if len(b.Operations) == 0 { + return nil + } + + // 1. Пишем batch в WAL. + data, err := json.Marshal(b) + if err != nil { + return fmt.Errorf("failed to marshal batch: %v", err) + } + + record := &WALRecord{Type: 1, Data: data} + if err := globalBatchManager.wal.WriteSync(record); err != nil { + globalBatchManager.stats.TotalFailed.Add(1) + return fmt.Errorf("failed to write batch to WAL: %v", err) + } + + // 2. Применяем операции. + if err := b.apply(); err != nil { + globalBatchManager.stats.TotalFailed.Add(1) + return fmt.Errorf("batch apply failed: %v", err) + } + + // 3. Также пишем в AOF для быстрого восстановления (опционально). + if globalBatchManager.aof != nil { + if err := globalBatchManager.aof.AppendBatch(b); err != nil { + if globalBatchManager.logger != nil { + globalBatchManager.logger.Warn(fmt.Sprintf("AOF append failed: %v", err)) + } + } + } + + globalBatchManager.stats.TotalBatches.Add(1) + globalBatchManager.stats.TotalOperations.Add(uint64(len(b.Operations))) + return nil +} + +// apply применяет все операции batch. +// +// Если операция падает, мы пытаемся откатить уже применённые +// (best-effort). Полная атомарность гарантируется только на уровне +// WAL: при восстановлении batch либо весь применится, либо нет. +func (b *Batch) apply() error { + if globalStorage == nil { + return fmt.Errorf("storage not initialized") + } + + // Собираем откат для уже применённых операций. + type appliedOp struct { + opType string + db string + coll string + docID string + oldDoc *Document + } + applied := make([]appliedOp, 0, len(b.Operations)) + + rollback := func() { + for i := len(applied) - 1; i >= 0; i-- { + ao := applied[i] + db, err := globalStorage.GetDatabase(ao.db) + if err != nil { + continue + } + coll, err := db.GetCollection(ao.coll) + if err != nil { + continue + } + switch ao.opType { + case "insert": + coll.PermanentDelete(ao.docID) + case "update": + if ao.oldDoc != nil { + coll.Update(ao.docID, ao.oldDoc.GetFields()) + } + case "delete": + coll.RestoreDeleted(ao.docID) + } + } + } + + for _, op := range b.Operations { + db, err := globalStorage.GetDatabase(op.Database) + if err != nil { + rollback() + return fmt.Errorf("database not found: %s", op.Database) + } + coll, err := db.GetCollection(op.Collection) + if err != nil { + rollback() + return fmt.Errorf("collection not found: %s", op.Collection) + } + + switch op.Type { + case "insert": + doc := NewDocumentWithID(op.DocumentID) + for k, v := range op.Data { + doc.SetField(k, v) + } + doc.Version = op.Version + if err := coll.Insert(doc); err != nil { + rollback() + return err + } + applied = append(applied, appliedOp{ + opType: "insert", db: op.Database, coll: op.Collection, docID: op.DocumentID, + }) + + case "update": + oldDoc, _ := coll.Find(op.DocumentID) + var oldCopy *Document + if oldDoc != nil { + oldCopy = oldDoc.Clone() + } + if err := coll.Update(op.DocumentID, op.Data); err != nil { + rollback() + return err + } + applied = append(applied, appliedOp{ + opType: "update", db: op.Database, coll: op.Collection, + docID: op.DocumentID, oldDoc: oldCopy, + }) + + case "delete": + if err := coll.Delete(op.DocumentID); err != nil { + rollback() + return err + } + applied = append(applied, appliedOp{ + opType: "delete", db: op.Database, coll: op.Collection, docID: op.DocumentID, + }) + + case "restore": + if err := coll.RestoreDeleted(op.DocumentID); err != nil { + rollback() + return err + } + applied = append(applied, appliedOp{ + opType: "restore", db: op.Database, coll: op.Collection, docID: op.DocumentID, + }) + } + } + return nil +} + +// ============================================================================= +// TIME-TRAVEL QUERIES (WAL-based) +// ============================================================================= + +// GetDocumentAtTimestamp возвращает документ на указанный момент времени, +// проигрывая WAL с начала до нужного LSN. +// +// Это O(N) по количеству записей WAL — для in-memory СУБД с небольшим WAL +// это приемлемо. Для больших WAL нужен индекс по времени. +func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) { + if globalBatchManager == nil { + return nil, fmt.Errorf("batch manager not initialized") + } + if globalBatchManager.wal == nil { + return nil, fmt.Errorf("WAL not initialized") + } + + records, err := globalBatchManager.wal.ReadAll() + if err != nil { + return nil, fmt.Errorf("failed to read WAL: %v", err) + } + + var current *Document + for _, rec := range records { + if rec.Timestamp > timestamp { + break + } + if rec.Type != 1 { + continue + } + var batch Batch + if err := json.Unmarshal(rec.Data, &batch); err != nil { + continue + } + for _, op := range batch.Operations { + if op.DocumentID != docID { + continue + } + switch op.Type { + case "insert": + doc := NewDocumentWithID(op.DocumentID) + for k, v := range op.Data { + doc.SetField(k, v) + } + doc.Version = op.Version + doc.CreatedAt = batch.Timestamp + doc.UpdatedAt = batch.Timestamp + current = doc + case "update": + if current != nil { + current.Update(op.Data) + current.UpdatedAt = batch.Timestamp + } + case "delete": + if current != nil { + current.SoftDelete() + current.DeletedAt = batch.Timestamp + } + case "restore": + if current != nil { + current.Restore() + current.UpdatedAt = batch.Timestamp + } + } + } + } + + if current == nil { + return nil, fmt.Errorf("document %s not found at timestamp %d", docID, timestamp) + } + return current, nil +} + +// GetDocumentVersionsSince возвращает все версии документа, +// созданные после указанного времени. +func GetDocumentVersionsSince(docID string, since int64) ([]*Document, error) { + if globalBatchManager == nil || globalBatchManager.wal == nil { + return nil, fmt.Errorf("batch manager not initialized") + } + records, err := globalBatchManager.wal.ReadAll() + if err != nil { + return nil, err + } + + versions := make([]*Document, 0) + for _, rec := range records { + if rec.Timestamp <= since || rec.Type != 1 { + continue + } + var batch Batch + if err := json.Unmarshal(rec.Data, &batch); err != nil { + continue + } + for _, op := range batch.Operations { + if op.DocumentID != docID { + continue + } + doc := NewDocumentWithID(op.DocumentID) + if op.Type == "insert" || op.Type == "update" { + for k, v := range op.Data { + doc.SetField(k, v) + } + doc.Version = op.Version + doc.UpdatedAt = batch.Timestamp + versions = append(versions, doc) + } + } + } + return versions, nil +} + +// ============================================================================= +// WAL MANAGER (Segmented) +// ============================================================================= + +// WALSegment представляет один сегмент WAL. +type WALSegment struct { + ID uint32 + File *os.File + Writer *bufio.Writer + Path string + StartLSN uint64 + EndLSN uint64 + Size int64 + mu sync.Mutex +} + +// WALIndexEntry — запись индекса WAL. +type WALIndexEntry struct { + LSN uint64 + SegmentID uint32 + Offset int64 + Length uint32 + Checksum uint32 +} + +// WALIndexManager — индекс WAL для быстрого доступа. +type WALIndexManager struct { + index map[uint64]*WALIndexEntry + segments map[uint32]*WALSegment + mu sync.RWMutex + indexPath string +} + +// SegmentedWALManager — сегментированный WAL. +type SegmentedWALManager struct { + segmentsDir string + segments map[uint32]*WALSegment + currentSegment *WALSegment + currentSegmentID uint32 + index *WALIndexManager + mu sync.RWMutex + writeChan chan *WALRecord + stopCh chan struct{} + wg sync.WaitGroup + batchSize int + logger LoggerInterface + fsyncEnabled bool + + // ackMap — map LSN -> chan struct{} для синхронного подтверждения. + // Ключ — LSN записи. После записи в WAL закрываем канал. + ackMu sync.Mutex + ackMap map[uint64]chan struct{} +} + +// NewSegmentedWALManager создаёт новый сегментированный WAL. +func NewSegmentedWALManager(segmentsDir string, fsyncEnabled bool, logger LoggerInterface) (*SegmentedWALManager, error) { + if err := os.MkdirAll(segmentsDir, 0755); err != nil { + return nil, fmt.Errorf("failed to create segments dir: %v", err) + } + wm := &SegmentedWALManager{ + segmentsDir: segmentsDir, + segments: make(map[uint32]*WALSegment), + index: &WALIndexManager{ + index: make(map[uint64]*WALIndexEntry), + segments: make(map[uint32]*WALSegment), + indexPath: filepath.Join(segmentsDir, WALIndexPrefix+"index.json"), + }, + writeChan: make(chan *WALRecord, 10000), + stopCh: make(chan struct{}), + batchSize: 100, + logger: logger, + fsyncEnabled: fsyncEnabled, + ackMap: make(map[uint64]chan struct{}), + } + if err := wm.loadExistingSegments(); err != nil { + return nil, err + } + if err := wm.index.load(); err != nil { + if logger != nil { + logger.Warn(fmt.Sprintf("Failed to load WAL index: %v", err)) + } + } + if wm.currentSegment == nil { + if err := wm.rotateSegmentInternal(); err != nil { + return nil, err + } + } + wm.wg.Add(1) + go wm.writerLoop() + return wm, nil +} + +func (wm *SegmentedWALManager) loadExistingSegments() error { + files, err := filepath.Glob(filepath.Join(wm.segmentsDir, WALSegmentPrefix+"*")) + if err != nil { + return err + } + for _, filePath := range files { + var segmentID uint32 + if _, err := fmt.Sscanf(filepath.Base(filePath), WALSegmentPrefix+"%d.log", &segmentID); err != nil { + continue + } + file, err := os.OpenFile(filePath, os.O_RDWR, 0644) + if err != nil { + continue + } + stat, _ := file.Stat() + segment := &WALSegment{ + ID: segmentID, + File: file, + Writer: bufio.NewWriterSize(file, 64*1024), + Path: filePath, + Size: stat.Size(), + StartLSN: uint64(segmentID) * WALSegmentSize / 100, + } + wm.segments[segmentID] = segment + if segmentID > wm.currentSegmentID { + wm.currentSegmentID = segmentID + wm.currentSegment = segment + } + } + return nil +} + +// rotateSegmentInternal — вызывается только при захваченном wm.mu. +func (wm *SegmentedWALManager) rotateSegmentInternal() error { + newSegmentID := wm.currentSegmentID + 1 + segmentPath := filepath.Join(wm.segmentsDir, fmt.Sprintf(WALSegmentPrefix+"%d.log", newSegmentID)) + file, err := os.OpenFile(segmentPath, os.O_CREATE|os.O_APPEND|os.O_RDWR, 0644) + if err != nil { + return fmt.Errorf("failed to create segment: %v", err) + } + newSegment := &WALSegment{ + ID: newSegmentID, + File: file, + Writer: bufio.NewWriterSize(file, 64*1024), + Path: segmentPath, + StartLSN: wm.getCurrentLSNLocked(), + } + if wm.currentSegment != nil { + wm.currentSegment.Writer.Flush() + if wm.fsyncEnabled { + RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond) + } + wm.currentSegment.EndLSN = wm.getCurrentLSNLocked() + wm.currentSegment.File.Close() + FsyncDir(wm.segmentsDir) + } + wm.currentSegment = newSegment + wm.currentSegmentID = newSegmentID + wm.segments[newSegmentID] = newSegment + if wm.logger != nil { + wm.logger.Info(fmt.Sprintf("Created new WAL segment: %d", newSegmentID)) + } + return nil +} + +func (wm *SegmentedWALManager) getCurrentLSNLocked() uint64 { + if wm.currentSegment == nil { + return 1 + } + return wm.currentSegment.StartLSN + uint64(wm.currentSegment.Size/100) +} + +// Write — асинхронная запись. +func (wm *SegmentedWALManager) Write(record *WALRecord) error { + record.Timestamp = time.Now().UnixMilli() + wm.writeChan <- record + return nil +} + +// WriteSync — синхронная запись с гарантией fsync. +// +// Используем ackMap по LSN, чтобы не было race condition +// с чужими подтверждениями. +func (wm *SegmentedWALManager) WriteSync(record *WALRecord) error { + record.Timestamp = time.Now().UnixMilli() + + // Присваиваем LSN под защитой mu. + wm.mu.Lock() + lsn := wm.getCurrentLSNLocked() + 1 + record.LSN = lsn + ackCh := make(chan struct{}) + wm.ackMu.Lock() + wm.ackMap[lsn] = ackCh + wm.ackMu.Unlock() + wm.mu.Unlock() + + wm.writeChan <- record + + select { + case <-ackCh: + return nil + case <-time.After(10 * time.Second): + wm.ackMu.Lock() + delete(wm.ackMap, lsn) + wm.ackMu.Unlock() + return fmt.Errorf("WAL write sync timeout (LSN %d)", lsn) + } +} + +func (wm *SegmentedWALManager) writerLoop() { + defer wm.wg.Done() + batch := make([]*WALRecord, 0, wm.batchSize) + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + for { + select { + case record, ok := <-wm.writeChan: + if !ok { + if len(batch) > 0 { + wm.flushBatch(batch) + } + return + } + batch = append(batch, record) + if len(batch) >= wm.batchSize { + wm.flushBatch(batch) + batch = batch[:0] + } + case <-ticker.C: + if len(batch) > 0 { + wm.flushBatch(batch) + batch = batch[:0] + } + case <-wm.stopCh: + if len(batch) > 0 { + wm.flushBatch(batch) + } + return + } + } +} + +// flushBatch — записывает batch записей и подтверждает через ackMap. +func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { + wm.mu.Lock() + + for _, record := range batch { + if wm.currentSegment.Size >= WALSegmentSize { + if err := wm.rotateSegmentInternal(); err != nil { + if wm.logger != nil { + wm.logger.Error(fmt.Sprintf("Failed to rotate segment: %v", err)) + } + continue + } + } + data, err := json.Marshal(record) + if err != nil { + continue + } + lsnBytes := make([]byte, 8) + binary.BigEndian.PutUint64(lsnBytes, record.LSN) + crcData := append(lsnBytes, data...) + record.CRC = crc32(crcData) + header := make([]byte, 8) + binary.BigEndian.PutUint32(header[0:4], uint32(len(data))) + binary.BigEndian.PutUint32(header[4:8], record.CRC) + if _, err := wm.currentSegment.Writer.Write(header); err != nil { + continue + } + if _, err := wm.currentSegment.Writer.Write(data); err != nil { + continue + } + wm.index.addEntry(&WALIndexEntry{ + LSN: record.LSN, + SegmentID: wm.currentSegment.ID, + Offset: wm.currentSegment.Size, + Length: uint32(len(data)), + Checksum: record.CRC, + }) + wm.currentSegment.Size += int64(8 + len(data)) + wm.currentSegment.EndLSN = record.LSN + } + wm.currentSegment.Writer.Flush() + if wm.fsyncEnabled { + RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond) + } + wm.index.save() + + wm.mu.Unlock() + + // Подтверждаем все записи из batch. + wm.ackMu.Lock() + for _, record := range batch { + if ch, ok := wm.ackMap[record.LSN]; ok { + close(ch) + delete(wm.ackMap, record.LSN) + } + } + wm.ackMu.Unlock() +} + +// Sync — синхронизация. +func (wm *SegmentedWALManager) Sync() error { + wm.mu.Lock() + defer wm.mu.Unlock() + if wm.currentSegment == nil { + return nil + } + if err := wm.currentSegment.Writer.Flush(); err != nil { + return err + } + if wm.fsyncEnabled { + return RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond) + } + return nil +} + +// ReadAll — читает все записи WAL. +func (wm *SegmentedWALManager) ReadAll() ([]*WALRecord, error) { + wm.mu.RLock() + segments := make([]*WALSegment, 0, len(wm.segments)) + for _, seg := range wm.segments { + segments = append(segments, seg) + } + wm.mu.RUnlock() + sort.Slice(segments, func(i, j int) bool { + return segments[i].ID < segments[j].ID + }) + records := make([]*WALRecord, 0) + for _, seg := range segments { + segRecords, err := wm.readSegmentRecords(seg) + if err != nil { + if wm.logger != nil { + wm.logger.Warn(fmt.Sprintf("Failed to read segment %d: %v", seg.ID, err)) + } + continue + } + records = append(records, segRecords...) + } + return records, nil +} + +// ReadSince — читает записи после указанного LSN. +func (wm *SegmentedWALManager) ReadSince(lsn uint64) ([]*WALRecord, error) { + allRecords, err := wm.ReadAll() + if err != nil { + return nil, err + } + result := make([]*WALRecord, 0) + for _, record := range allRecords { + if record.LSN > lsn { + result = append(result, record) + } + } + return result, nil +} + +// GetCurrentLSN — текущий LSN. +func (wm *SegmentedWALManager) GetCurrentLSN() uint64 { + wm.mu.RLock() + defer wm.mu.RUnlock() + if wm.currentSegment == nil { + return 1 + } + return wm.currentSegment.EndLSN +} + +func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord, error) { + seg.mu.Lock() + defer seg.mu.Unlock() + if seg.File == nil { + return nil, nil + } + seg.Writer.Flush() + seg.File.Seek(0, 0) + records := make([]*WALRecord, 0) + reader := bufio.NewReader(seg.File) + headerBuf := make([]byte, 8) + for { + _, err := io.ReadFull(reader, headerBuf) + if err != nil { + break + } + recordLen := binary.BigEndian.Uint32(headerBuf[0:4]) + expectedCRC := binary.BigEndian.Uint32(headerBuf[4:8]) + if recordLen == 0 || recordLen > 100*1024*1024 { + break + } + recordData := make([]byte, recordLen) + _, err = io.ReadFull(reader, recordData) + if err != nil { + break + } + var record WALRecord + if err := json.Unmarshal(recordData, &record); err != nil { + continue + } + lsnBytes := make([]byte, 8) + binary.BigEndian.PutUint64(lsnBytes, record.LSN) + crcData := append(lsnBytes, recordData...) + calculatedCRC := crc32(crcData) + if calculatedCRC != expectedCRC && calculatedCRC != record.CRC { + continue + } + records = append(records, &record) + } + return records, nil +} + +// Close — закрывает WAL. +func (wm *SegmentedWALManager) Close() error { + close(wm.stopCh) + wm.wg.Wait() + close(wm.writeChan) + + wm.mu.Lock() + defer wm.mu.Unlock() + if wm.currentSegment != nil { + wm.currentSegment.Writer.Flush() + if wm.fsyncEnabled { + RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond) + } + wm.currentSegment.File.Close() + } + wm.index.save() + return nil +} + +func (im *WALIndexManager) addEntry(entry *WALIndexEntry) { + im.mu.Lock() + defer im.mu.Unlock() + im.index[entry.LSN] = entry +} + +func (im *WALIndexManager) save() error { + im.mu.RLock() + defer im.mu.RUnlock() + data, err := json.Marshal(im.index) + if err != nil { + return err + } + return os.WriteFile(im.indexPath, data, 0644) +} + +func (im *WALIndexManager) load() error { + data, err := os.ReadFile(im.indexPath) + if err != nil { + if os.IsNotExist(err) { + return nil + } + return err + } + if len(data) == 0 { + return nil + } + return json.Unmarshal(data, &im.index) +} + +// ============================================================================= +// SAGA PERSISTENT STORAGE +// ============================================================================= + +// SagaState — состояние SAGA. +type SagaState struct { + ID string `json:"id"` + Status string `json:"status"` + CurrentStep int `json:"current_step"` + Steps []SagaStepState `json:"steps"` + Data map[string]interface{} `json:"data"` + CreatedAt int64 `json:"created_at"` + UpdatedAt int64 `json:"updated_at"` + CompletedAt int64 `json:"completed_at,omitempty"` + CompensationExecuted bool `json:"compensation_executed"` + NodeID string `json:"node_id"` + CoordinatorID string `json:"coordinator_id"` + Version uint64 `json:"version"` + RetryCount int `json:"retry_count"` + LastError string `json:"last_error,omitempty"` + ExecutionID string `json:"execution_id"` +} + +// SagaStepState — состояние шага SAGA. +type SagaStepState struct { + ID string `json:"id"` + Name string `json:"name"` + Status string `json:"status"` + Data map[string]interface{} `json:"data"` + StartedAt int64 `json:"started_at"` + CompletedAt int64 `json:"completed_at"` + ExecutionID string `json:"execution_id"` + RetryCount int `json:"retry_count"` + LastError string `json:"last_error,omitempty"` + CompensatedAt int64 `json:"compensated_at,omitempty"` +} + +// SagaPersistentStorage — персистентное хранилище SAGA. +type SagaPersistentStorage struct { + baseDir string + mu sync.RWMutex + cache map[string]*SagaState + maxCache int + logger LoggerInterface + fsyncEnabled bool + fsyncMaxRetries int + fsyncRetryDelay time.Duration +} + +// NewSagaPersistentStorageWithConfig создаёт хранилище SAGA. +func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInterface) (*SagaPersistentStorage, error) { + baseDir := SagaStateDir + maxCache := 10000 + fsyncEnabled := true + fsyncMaxRetries := 3 + fsyncRetryDelay := 100 * time.Millisecond + + if cfg != nil { + baseDir = cfg.GetStateDir() + maxCache = cfg.GetMaxCacheSize() + fsyncEnabled = cfg.IsFsyncEnabled() + fsyncMaxRetries = cfg.GetFsyncMaxRetries() + fsyncRetryDelay = cfg.GetFsyncRetryDelay() + } + + fullPath := filepath.Join(baseDir) + if err := os.MkdirAll(fullPath, 0755); err != nil { + return nil, fmt.Errorf("failed to create saga state directory: %v", err) + } + + return &SagaPersistentStorage{ + baseDir: fullPath, + cache: make(map[string]*SagaState), + maxCache: maxCache, + logger: logger, + fsyncEnabled: fsyncEnabled, + fsyncMaxRetries: fsyncMaxRetries, + fsyncRetryDelay: fsyncRetryDelay, + }, nil +} + +func (sps *SagaPersistentStorage) getStatePath(sagaID string) string { + return filepath.Join(sps.baseDir, fmt.Sprintf("%s%s%s", SagaStateFilePrefix, sagaID, SagaStateFileSuffix)) +} + +// Save сохраняет состояние SAGA. +func (sps *SagaPersistentStorage) Save(state *SagaState) error { + sps.mu.Lock() + defer sps.mu.Unlock() + + state.UpdatedAt = time.Now().UnixMilli() + state.Version++ + + if len(sps.cache) >= sps.maxCache { + var oldestKey string + var oldestTime int64 = time.Now().UnixMilli() + for k, v := range sps.cache { + if v.UpdatedAt < oldestTime { + oldestTime = v.UpdatedAt + oldestKey = k + } + } + if oldestKey != "" { + delete(sps.cache, oldestKey) + } + } + sps.cache[state.ID] = state + + path := sps.getStatePath(state.ID) + data, err := json.MarshalIndent(state, "", " ") + if err != nil { + return fmt.Errorf("failed to marshal saga state: %v", err) + } + + tmpPath := path + ".tmp" + if err := os.WriteFile(tmpPath, data, 0644); err != nil { + return fmt.Errorf("failed to write saga state: %v", err) + } + + if sps.fsyncEnabled { + if f, err := os.OpenFile(tmpPath, os.O_RDWR, 0644); err == nil { + RealFsyncWithRetry(f, sps.fsyncMaxRetries, sps.fsyncRetryDelay) + f.Close() + } + } + + if err := os.Rename(tmpPath, path); err != nil { + return fmt.Errorf("failed to rename saga state: %v", err) + } + + if sps.fsyncEnabled { + FsyncDir(sps.baseDir) + } + return nil +} + +// Load загружает состояние SAGA. +func (sps *SagaPersistentStorage) Load(sagaID string) (*SagaState, error) { + sps.mu.RLock() + if state, ok := sps.cache[sagaID]; ok { + sps.mu.RUnlock() + return state, nil + } + sps.mu.RUnlock() + + path := sps.getStatePath(sagaID) + data, err := os.ReadFile(path) + if err != nil { + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("failed to read saga state: %v", err) + } + + var state SagaState + if err := json.Unmarshal(data, &state); err != nil { + return nil, fmt.Errorf("failed to unmarshal saga state: %v", err) + } + + sps.mu.Lock() + sps.cache[sagaID] = &state + sps.mu.Unlock() + return &state, nil +} + +// Delete удаляет состояние SAGA. +func (sps *SagaPersistentStorage) Delete(sagaID string) error { + sps.mu.Lock() + delete(sps.cache, sagaID) + sps.mu.Unlock() + + path := sps.getStatePath(sagaID) + if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + return fmt.Errorf("failed to delete saga state: %v", err) + } + return nil +} + +// ListAll возвращает все состояния SAGA. +func (sps *SagaPersistentStorage) ListAll() ([]*SagaState, error) { + pattern := filepath.Join(sps.baseDir, fmt.Sprintf("%s*%s", SagaStateFilePrefix, SagaStateFileSuffix)) + files, err := filepath.Glob(pattern) + if err != nil { + return nil, err + } + states := make([]*SagaState, 0, len(files)) + for _, file := range files { + data, err := os.ReadFile(file) + if err != nil { + continue + } + var state SagaState + if err := json.Unmarshal(data, &state); err != nil { + continue + } + states = append(states, &state) + } + return states, nil +} + +// ListPending возвращает незавершённые SAGA. +func (sps *SagaPersistentStorage) ListPending() ([]*SagaState, error) { + all, err := sps.ListAll() + if err != nil { + return nil, err + } + pending := make([]*SagaState, 0) + for _, state := range all { + if state.Status == "pending" || state.Status == "running" || state.Status == "compensating" { + pending = append(pending, state) + } + } + return pending, nil +} + +// ============================================================================= +// SAGA ORCHESTRATOR (упрощённый, без мьютексов на горячем пути) +// ============================================================================= + +// SagaOrchestrator управляет SAGA. +type SagaOrchestrator struct { + storage *SagaPersistentStorage + mu sync.RWMutex + logger LoggerInterface + nodeID string + stopChan chan struct{} + wg sync.WaitGroup + isLeader atomic.Bool + leaderID string + config *config.SagaConfig + metrics *SagaMetrics + activeSagas sync.Map // map[string]*SagaTransaction +} + +// SagaMetrics — метрики SAGA. +type SagaMetrics struct { + TotalStarted atomic.Uint64 + TotalCompleted atomic.Uint64 + TotalAborted atomic.Uint64 + TotalFailed atomic.Uint64 + TotalCompensated atomic.Uint64 + ActiveCount atomic.Int64 + RecoveryCount atomic.Uint64 +} + +// NewSagaOrchestratorWithConfig создаёт оркестратор. +func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) { + if cfg == nil { + cfg = &config.SagaConfig{ + Enabled: true, + StateDir: "saga_states", + MaxRetries: 5, + RetryBackoffMs: 100, + SagaTimeoutSec: 300, + StuckCheckIntervalSec: 10, + RecoveryIntervalSec: 30, + CleanupPeriodHours: 24, + MaxCacheSize: 10000, + FsyncEnabled: true, + FsyncMaxRetries: 3, + FsyncRetryDelayMs: 100, + } + } + + storage, err := NewSagaPersistentStorageWithConfig(cfg, logger) + if err != nil { + return nil, err + } + + o := &SagaOrchestrator{ + storage: storage, + logger: logger, + nodeID: nodeID, + stopChan: make(chan struct{}), + metrics: &SagaMetrics{}, + config: cfg, + } + o.isLeader.Store(true) + + o.wg.Add(1) + go o.recoveryLoop() + o.wg.Add(1) + go o.orphanCleanupLoop() + + return o, nil +} + +// IsLeader — является ли узел лидером. +func (o *SagaOrchestrator) IsLeader() bool { return o.isLeader.Load() } + +// GetLeaderID — ID лидера. +func (o *SagaOrchestrator) GetLeaderID() string { return o.leaderID } + +func (o *SagaOrchestrator) recoveryLoop() { + defer o.wg.Done() + ticker := time.NewTicker(o.config.GetRecoveryInterval()) + defer ticker.Stop() + for { + select { + case <-ticker.C: + o.recoverPendingSagas() + case <-o.stopChan: + return + } + } +} + +func (o *SagaOrchestrator) orphanCleanupLoop() { + defer o.wg.Done() + ticker := time.NewTicker(SagaCleanupInterval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + o.cleanupOrphanSteps() + case <-o.stopChan: + return + } + } +} + +// cleanupOrphanSteps — очистка "осиротевших" SAGA. +func (o *SagaOrchestrator) cleanupOrphanSteps() { + pending, err := o.storage.ListPending() + if err != nil { + return + } + now := time.Now().UnixMilli() + for _, state := range pending { + if state.Status == "running" || state.Status == "compensating" { + elapsed := now - state.UpdatedAt + if elapsed > SagaOrphanTimeout.Milliseconds() { + if o.logger != nil { + o.logger.Warn(fmt.Sprintf("Found orphan saga %s, cleaning up", state.ID)) + } + state.Status = "orphaned" + state.LastError = fmt.Sprintf("orphaned after %d ms", elapsed) + o.storage.Save(state) + o.compensateOrphanSaga(state) + } + } + } +} + +// compensateOrphanSaga — компенсация "осиротевшей" SAGA. +// +// Проверяем nil на Compensate, сохраняем состояние после +// каждой компенсации, обрабатываем compensation_failed. +func (o *SagaOrchestrator) compensateOrphanSaga(state *SagaState) { + if state.CompensationExecuted { + return + } + if o.logger != nil { + o.logger.Info(fmt.Sprintf("Compensating orphan saga %s", state.ID)) + } + + saga := o.stateToTransaction(state) + if err := o.executeCompensation(saga, len(saga.Steps)-1); err != nil { + if o.logger != nil { + o.logger.Error(fmt.Sprintf("Failed to compensate orphan saga %s: %v", state.ID, err)) + } + state.Status = "compensation_failed" + state.LastError = err.Error() + o.storage.Save(state) + return + } + state.CompensationExecuted = true + state.Status = "orphaned_compensated" + o.storage.Save(state) +} + +// stateToTransaction конвертирует SagaState в SagaTransaction. +func (o *SagaOrchestrator) stateToTransaction(state *SagaState) *SagaTransaction { + saga := &SagaTransaction{ + ID: state.ID, + Status: state.Status, + CurrentStep: state.CurrentStep, + Data: state.Data, + CreatedAt: state.CreatedAt, + UpdatedAt: state.UpdatedAt, + CompletedAt: state.CompletedAt, + compensationExecuted: state.CompensationExecuted, + executedOps: make(map[string]bool), + } + for _, stepState := range state.Steps { + step := &SagaStep{ + ID: stepState.ID, + Name: stepState.Name, + Status: stepState.Status, + Data: stepState.Data, + StartedAt: stepState.StartedAt, + CompletedAt: stepState.CompletedAt, + ExecutionID: stepState.ExecutionID, + RetryCount: stepState.RetryCount, + LastError: stepState.LastError, + } + saga.Steps = append(saga.Steps, step) + if stepState.Status == "completed" { + saga.executedOps[stepState.ExecutionID] = true + } + } + return saga +} + +// recoverPendingSagas — восстановление незавершённых SAGA. +func (o *SagaOrchestrator) recoverPendingSagas() { + pending, err := o.storage.ListPending() + if err != nil { + return + } + if len(pending) == 0 { + return + } + if o.logger != nil { + o.logger.Info(fmt.Sprintf("Recovering %d pending sagas", len(pending))) + } + for _, state := range pending { + if err := o.resumeSaga(state); err != nil { + if o.logger != nil { + o.logger.Error(fmt.Sprintf("Failed to resume saga %s: %v", state.ID, err)) + } + } else { + o.metrics.RecoveryCount.Add(1) + } + } +} + +// resumeSaga — возобновление SAGA. +func (o *SagaOrchestrator) resumeSaga(state *SagaState) error { + saga := o.stateToTransaction(state) + o.activeSagas.Store(state.ID, saga) + return o.executeSagaInternal(saga) +} + +// BeginSaga — начать SAGA. +func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { + existing, err := o.storage.Load(id) + if err != nil { + return nil, err + } + if existing != nil { + return nil, fmt.Errorf("saga %s already exists", id) + } + now := time.Now().UnixMilli() + state := &SagaState{ + ID: id, + Status: "pending", + CurrentStep: 0, + Data: make(map[string]interface{}), + CreatedAt: now, + UpdatedAt: now, + NodeID: o.nodeID, + Version: 1, + ExecutionID: fmt.Sprintf("%s_%d", id, now), + } + if err := o.storage.Save(state); err != nil { + return nil, err + } + saga := &SagaTransaction{ + ID: id, + Steps: make([]*SagaStep, 0), + CurrentStep: 0, + Status: "pending", + CreatedAt: now, + UpdatedAt: now, + Data: make(map[string]interface{}), + executedOps: make(map[string]bool), + } + o.activeSagas.Store(id, saga) + o.metrics.TotalStarted.Add(1) + o.metrics.ActiveCount.Add(1) + return saga, nil +} + +// Execute — выполнить SAGA. +func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error { + return o.executeSagaInternal(saga) +} + +// executeSagaInternal — внутренняя логика выполнения SAGA. +// +// Компенсация идемпотентна и корректно обрабатывает ошибки. +func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { + saga.mu.Lock() + if saga.Status != "pending" && saga.Status != "running" { + saga.mu.Unlock() + return fmt.Errorf("saga %s is not in pending or running state", saga.ID) + } + saga.Status = "running" + saga.UpdatedAt = time.Now().UnixMilli() + saga.mu.Unlock() + + if err := o.saveSagaState(saga); err != nil { + return err + } + + for i := saga.CurrentStep; i < len(saga.Steps); i++ { + saga.mu.Lock() + step := saga.Steps[i] + saga.CurrentStep = i + if step.Status == "completed" { + saga.mu.Unlock() + continue + } + step.Status = "running" + step.StartedAt = time.Now().UnixMilli() + saga.mu.Unlock() + + if saga.HasExecuted(step.ExecutionID) { + saga.mu.Lock() + step.Status = "completed" + step.CompletedAt = time.Now().UnixMilli() + saga.mu.Unlock() + continue + } + + var err error + maxRetries := SagaMaxRetries + if o.config != nil { + maxRetries = o.config.GetMaxRetries() + } + for retry := 0; retry < maxRetries; retry++ { + saga.mu.Lock() + step.RetryCount = retry + 1 + saga.mu.Unlock() + if err = step.Execute(); err == nil { + break + } + saga.mu.Lock() + step.LastError = err.Error() + saga.mu.Unlock() + backoff := time.Duration(o.config.RetryBackoffMs*(1<= 0; j-- { + saga.mu.RLock() + step := saga.Steps[j] + status := step.Status + saga.mu.RUnlock() + + if status == "compensated" || status == "pending" { + continue + } + if status != "completed" && status != "failed" { + continue + } + + if step.Compensate == nil { + // Нет функции компенсации — считаем шаг скомпенсированным. + saga.mu.Lock() + step.Status = "compensated" + step.CompletedAt = time.Now().UnixMilli() + saga.mu.Unlock() + o.saveSagaState(saga) + continue + } + + var err error + maxRetries := SagaMaxRetries + if o.config != nil { + maxRetries = o.config.GetMaxRetries() + } + for retry := 0; retry < maxRetries; retry++ { + if err = step.Compensate(); err == nil { + break + } + backoff := time.Duration(o.config.RetryBackoffMs*(1<