/* * 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. // // ДОБАВЛЕНО (2026-10, идемпотентность SAGA + выборы лидера): // - SagaPersistentStorage ведёт персистентную таблицу выполненных // ExecutionID (saga_executed_keys.json). Перед выполнением шага // и перед его компенсацией оркестратор проверяет эту таблицу, // чтобы повторный запуск SAGA после падения не приводил к // дублированию операций на удалённых узлах. // - SagaOrchestrator получил SagaLeaderElector, который берёт лидера // из Raft (через SagaRaftAccessor). Если Raft недоступен — // используется резервный режим «единственный узел — лидер». // - isLeader больше не выставляется в true безусловно: значение // обновляется фоновой горутиной leaderElectionLoop, которая // опрашивает Raft.State() и публикует изменения в orchestrator. // - Добавлен SagaLeaderState (персистентный) — чтобы после рестарта // узел знал, был ли он лидером, и корректно себя вёл до первых // выборов. // - Добавлен интерфейс SagaRaftAccessor, чтобы storage не зависел // напрямую от hashicorp/raft. // // ДОБАВЛЕНО (2026-10, ревизия): // - GetAllSagas у SagaManager — читает все SAGA (включая завершённые) // с диска через SagaPersistentStorage.ListAll(). Используется // REPL-командой "saga list --all". 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 // SagaExecutedKeysFile — персистентная таблица выполненных ExecutionID. SagaExecutedKeysFile = "saga_executed_keys.json" // SagaLeaderStateFile — персистентное состояние лидера SAGA. SagaLeaderStateFile = "saga_leader_state.json" // SagaLeaderElectionInterval — как часто опрашиваем Raft на предмет лидерства. SagaLeaderElectionInterval = 2 * time.Second ) // ============================================================================= // 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 // activeBatches — реестр активных (незакоммиченных) 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 } 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 как активный. func (bm *BatchManager) RegisterBatch(b *Batch) { if bm == nil || b == nil { return } bm.activeBatches.Store(b.ID, b) } // UnregisterBatch снимает регистрацию batch. func (bm *BatchManager) UnregisterBatch(b *Batch) { if bm == nil || b == nil { return } bm.activeBatches.Delete(b.ID) } // IsDocInActiveBatch проверяет, участвует ли документ в активной 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, } 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 атомарно. func (b *Batch) Commit() error { if globalBatchManager == nil { return fmt.Errorf("batch manager not initialized") } defer globalBatchManager.UnregisterBatch(b) if len(b.Operations) == 0 { return nil } 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) } if err := b.apply(); err != nil { globalBatchManager.stats.TotalFailed.Add(1) return fmt.Errorf("batch apply failed: %v", err) } 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. 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 возвращает документ на указанный момент времени. 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 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. func (wm *SegmentedWALManager) WriteSync(record *WALRecord) error { record.Timestamp = time.Now().UnixMilli() 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() 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 — ИНТЕРФЕЙСЫ ДЛЯ RAFT // ============================================================================= // SagaRaftAccessor — минимальный интерфейс к Raft, необходимый // оркестратору SAGA для выборов лидера. // // Реализуется в internal/cluster (SagaLeaderBridge), чтобы storage // не зависел напрямую от hashicorp/raft. type SagaRaftAccessor interface { // IsSagaLeader возвращает true, если текущий узел является лидером // в Raft-группе. IsSagaLeader() bool // GetSagaLeaderID возвращает ID текущего лидера SAGA (может быть // пустым, если лидер ещё не выбран). GetSagaLeaderID() string // GetSagaCurrentTerm возвращает текущий терм Raft. GetSagaCurrentTerm() uint64 // RegisterSagaLeaderObserver регистрирует callback, который будет // вызван при смене лидера SAGA. RegisterSagaLeaderObserver(id string, cb func(isLeader bool, leaderID string)) } // SagaLeaderState — персистентное состояние лидера SAGA. type SagaLeaderState struct { LeaderID string `json:"leader_id"` IsLeader bool `json:"is_leader"` Term uint64 `json:"term"` UpdatedAt int64 `json:"updated_at"` NodeID string `json:"node_id"` } // ============================================================================= // 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. // // ДОБАВЛЕНО: executedKeys — персистентная таблица уже выполненных // ExecutionID. Заполняется после успешного выполнения шага и после // успешной компенсации. Проверяется перед выполнением/компенсацией. type SagaPersistentStorage struct { baseDir string mu sync.RWMutex cache map[string]*SagaState maxCache int logger LoggerInterface fsyncEnabled bool fsyncMaxRetries int fsyncRetryDelay time.Duration executedMu sync.RWMutex executedKeys map[string]int64 } // 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) } sps := &SagaPersistentStorage{ baseDir: fullPath, cache: make(map[string]*SagaState), maxCache: maxCache, logger: logger, fsyncEnabled: fsyncEnabled, fsyncMaxRetries: fsyncMaxRetries, fsyncRetryDelay: fsyncRetryDelay, executedKeys: make(map[string]int64), } if err := sps.loadExecutedKeys(); err != nil { if logger != nil { logger.Warn(fmt.Sprintf("Failed to load saga executed keys: %v", err)) } } return sps, nil } // getStatePath — путь к файлу состояния SAGA. func (sps *SagaPersistentStorage) getStatePath(sagaID string) string { return filepath.Join(sps.baseDir, fmt.Sprintf("%s%s%s", SagaStateFilePrefix, sagaID, SagaStateFileSuffix)) } // getExecutedKeysPath — путь к файлу таблицы выполненных ключей. func (sps *SagaPersistentStorage) getExecutedKeysPath() string { return filepath.Join(sps.baseDir, SagaExecutedKeysFile) } // loadExecutedKeys загружает таблицу выполненных ключей с диска. func (sps *SagaPersistentStorage) loadExecutedKeys() error { data, err := os.ReadFile(sps.getExecutedKeysPath()) if err != nil { if os.IsNotExist(err) { return nil } return err } if len(data) == 0 { return nil } var loaded map[string]int64 if err := json.Unmarshal(data, &loaded); err != nil { return err } sps.executedMu.Lock() sps.executedKeys = loaded sps.executedMu.Unlock() return nil } // saveExecutedKeys сохраняет таблицу выполненных ключей на диск. func (sps *SagaPersistentStorage) saveExecutedKeys() error { sps.executedMu.RLock() data, err := json.Marshal(sps.executedKeys) sps.executedMu.RUnlock() if err != nil { return err } path := sps.getExecutedKeysPath() tmpPath := path + ".tmp" if err := os.WriteFile(tmpPath, data, 0644); err != nil { return 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 err } if sps.fsyncEnabled { FsyncDir(sps.baseDir) } return nil } // IsExecuted проверяет, был ли ExecutionID уже выполнен. func (sps *SagaPersistentStorage) IsExecuted(executionID string) bool { sps.executedMu.RLock() defer sps.executedMu.RUnlock() _, ok := sps.executedKeys[executionID] return ok } // MarkExecuted атомарно отмечает ExecutionID как выполненный // и сохраняет таблицу на диск. // // Возвращает true, если ключ был добавлен впервые. func (sps *SagaPersistentStorage) MarkExecuted(executionID string) bool { sps.executedMu.Lock() if _, exists := sps.executedKeys[executionID]; exists { sps.executedMu.Unlock() return false } sps.executedKeys[executionID] = time.Now().UnixMilli() sps.executedMu.Unlock() if err := sps.saveExecutedKeys(); err != nil { if sps.logger != nil { sps.logger.Warn(fmt.Sprintf("Failed to persist executed keys: %v", err)) } } return true } // 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 } // SaveLeaderState сохраняет состояние лидера SAGA на диск. func (sps *SagaPersistentStorage) SaveLeaderState(state *SagaLeaderState) error { path := filepath.Join(sps.baseDir, SagaLeaderStateFile) data, err := json.MarshalIndent(state, "", " ") if err != nil { return err } tmpPath := path + ".tmp" if err := os.WriteFile(tmpPath, data, 0644); err != nil { return 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 err } if sps.fsyncEnabled { FsyncDir(sps.baseDir) } return nil } // LoadLeaderState загружает состояние лидера SAGA с диска. func (sps *SagaPersistentStorage) LoadLeaderState() (*SagaLeaderState, error) { path := filepath.Join(sps.baseDir, SagaLeaderStateFile) data, err := os.ReadFile(path) if err != nil { if os.IsNotExist(err) { return nil, nil } return nil, err } var state SagaLeaderState if err := json.Unmarshal(data, &state); err != nil { return nil, err } return &state, nil } // ============================================================================= // SAGA LEADER ELECTOR // ============================================================================= // SagaLeaderElector — выборы лидера SAGA через Raft. type SagaLeaderElector struct { nodeID string raft SagaRaftAccessor storage *SagaPersistentStorage logger LoggerInterface isLeader atomic.Bool leaderID atomic.Value currentTerm atomic.Uint64 observersMu sync.RWMutex observers map[string]func(isLeader bool, leaderID string) stopChan chan struct{} wg sync.WaitGroup } // NewSagaLeaderElector создаёт выборщик лидера SAGA. func NewSagaLeaderElector(nodeID string, raft SagaRaftAccessor, storage *SagaPersistentStorage, logger LoggerInterface) *SagaLeaderElector { e := &SagaLeaderElector{ nodeID: nodeID, raft: raft, storage: storage, logger: logger, observers: make(map[string]func(isLeader bool, leaderID string)), stopChan: make(chan struct{}), } e.leaderID.Store("") if storage != nil { if state, err := storage.LoadLeaderState(); err == nil && state != nil { e.isLeader.Store(state.IsLeader) e.currentTerm.Store(state.Term) e.leaderID.Store(state.LeaderID) } } if raft != nil { raft.RegisterSagaLeaderObserver(nodeID, func(isLeader bool, leaderID string) { e.updateLeadership(isLeader, leaderID) }) } e.wg.Add(1) go e.leaderElectionLoop() return e } // updateLeadership — центральная точка обновления лидерства. func (e *SagaLeaderElector) updateLeadership(isLeader bool, leaderID string) { oldLeader := e.isLeader.Load() oldLeaderID, _ := e.leaderID.Load().(string) e.isLeader.Store(isLeader) e.leaderID.Store(leaderID) term := uint64(0) if e.raft != nil { term = e.raft.GetSagaCurrentTerm() } e.currentTerm.Store(term) if e.storage != nil { _ = e.storage.SaveLeaderState(&SagaLeaderState{ LeaderID: leaderID, IsLeader: isLeader, Term: term, UpdatedAt: time.Now().UnixMilli(), NodeID: e.nodeID, }) } if oldLeader != isLeader || oldLeaderID != leaderID { e.observersMu.RLock() obs := make([]func(isLeader bool, leaderID string), 0, len(e.observers)) for _, cb := range e.observers { obs = append(obs, cb) } e.observersMu.RUnlock() for _, cb := range obs { func(cb func(isLeader bool, leaderID string)) { defer func() { if r := recover(); r != nil && e.logger != nil { e.logger.Error(fmt.Sprintf("Saga leader observer panicked: %v", r)) } }() cb(isLeader, leaderID) }(cb) } if e.logger != nil { e.logger.Info(fmt.Sprintf("Saga leadership changed: isLeader=%v, leaderID=%s, term=%d", isLeader, leaderID, term)) } } } // leaderElectionLoop — фоновый опрос Raft на предмет лидерства. func (e *SagaLeaderElector) leaderElectionLoop() { defer e.wg.Done() ticker := time.NewTicker(SagaLeaderElectionInterval) defer ticker.Stop() for { select { case <-ticker.C: if e.raft == nil { continue } isLeader := e.raft.IsSagaLeader() leaderID := e.raft.GetSagaLeaderID() e.updateLeadership(isLeader, leaderID) case <-e.stopChan: return } } } // IsLeader — является ли текущий узел лидером SAGA. func (e *SagaLeaderElector) IsLeader() bool { return e.isLeader.Load() } // GetLeaderID — ID текущего лидера SAGA. func (e *SagaLeaderElector) GetLeaderID() string { v, _ := e.leaderID.Load().(string) return v } // GetCurrentTerm — текущий терм Raft. func (e *SagaLeaderElector) GetCurrentTerm() uint64 { return e.currentTerm.Load() } // RegisterObserver регистрирует наблюдателя за сменой лидера. func (e *SagaLeaderElector) RegisterObserver(id string, cb func(isLeader bool, leaderID string)) { e.observersMu.Lock() defer e.observersMu.Unlock() e.observers[id] = cb } // UnregisterObserver снимает регистрацию наблюдателя. func (e *SagaLeaderElector) UnregisterObserver(id string) { e.observersMu.Lock() defer e.observersMu.Unlock() delete(e.observers, id) } // Stop останавливает выборщик. func (e *SagaLeaderElector) Stop() { close(e.stopChan) e.wg.Wait() } // ============================================================================= // 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 leaderElector *SagaLeaderElector raftAccessor SagaRaftAccessor } // 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 IdempotentSkips atomic.Uint64 NotLeaderRejects atomic.Uint64 } // NewSagaOrchestratorWithConfig создаёт оркестратор без Raft. func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) { return NewSagaOrchestratorWithRaft(cfg, nodeID, nil, logger) } // NewSagaOrchestratorWithRaft создаёт оркестратор с привязкой к Raft. func NewSagaOrchestratorWithRaft(cfg *config.SagaConfig, nodeID string, raftAccessor SagaRaftAccessor, 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, raftAccessor: raftAccessor, } o.leaderElector = NewSagaLeaderElector(nodeID, raftAccessor, storage, logger) if raftAccessor == nil { o.isLeader.Store(true) o.leaderID = nodeID o.leaderElector.updateLeadership(true, nodeID) } else { o.isLeader.Store(o.leaderElector.IsLeader()) o.leaderID = o.leaderElector.GetLeaderID() } o.leaderElector.RegisterObserver("orchestrator", func(isLeader bool, leaderID string) { o.mu.Lock() o.leaderID = leaderID o.mu.Unlock() }) o.wg.Add(1) go o.recoveryLoop() o.wg.Add(1) go o.orphanCleanupLoop() return o, nil } // IsLeader — является ли узел лидером. func (o *SagaOrchestrator) IsLeader() bool { if o.leaderElector == nil { return true } return o.leaderElector.IsLeader() } // GetLeaderID — ID лидера. func (o *SagaOrchestrator) GetLeaderID() string { if o.leaderElector == nil { return o.nodeID } if id := o.leaderElector.GetLeaderID(); id != "" { return id } o.mu.RLock() defer o.mu.RUnlock() return o.leaderID } // IsExecuted проверяет, был ли ExecutionID уже выполнен. func (o *SagaOrchestrator) IsExecuted(executionID string) bool { if o.storage == nil { return false } return o.storage.IsExecuted(executionID) } // MarkExecuted помечает ExecutionID как выполненный. func (o *SagaOrchestrator) MarkExecuted(executionID string) bool { if o.storage == nil { return true } first := o.storage.MarkExecuted(executionID) if !first { o.metrics.IdempotentSkips.Add(1) } return first } func (o *SagaOrchestrator) recoveryLoop() { defer o.wg.Done() ticker := time.NewTicker(o.config.GetRecoveryInterval()) defer ticker.Stop() for { select { case <-ticker.C: if !o.IsLeader() { continue } 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: if !o.IsLeader() { continue } 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. func (o *SagaOrchestrator) compensateOrphanSaga(state *SagaState) { if !o.IsLeader() { return } 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() { if !o.IsLeader() { return } 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) { if !o.IsLeader() { o.metrics.NotLeaderRejects.Add(1) return nil, fmt.Errorf("node is not the saga leader (leader: %s)", o.GetLeaderID()) } 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 { if !o.IsLeader() { o.metrics.NotLeaderRejects.Add(1) return fmt.Errorf("node is not the saga leader (leader: %s)", o.GetLeaderID()) } return o.executeSagaInternal(saga) } // executeSagaInternal — внутренняя логика выполнения SAGA. func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { if !o.IsLeader() { return fmt.Errorf("node is not the saga leader") } 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++ { if !o.IsLeader() { return fmt.Errorf("lost saga leadership during execution of %s", saga.ID) } 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 o.IsExecuted(step.ExecutionID) || saga.HasExecuted(step.ExecutionID) { saga.mu.Lock() step.Status = "completed" step.CompletedAt = time.Now().UnixMilli() saga.mu.Unlock() if o.logger != nil { o.logger.Debug(fmt.Sprintf("Saga step %s (execution_id=%s) already executed, skipping", step.Name, step.ExecutionID)) } 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-- { if !o.IsLeader() { return fmt.Errorf("lost saga leadership during compensation of %s", saga.ID) } 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 } compExecID := step.ExecutionID + ":compensate" if o.IsExecuted(compExecID) { saga.mu.Lock() step.Status = "compensated" step.CompletedAt = time.Now().UnixMilli() saga.mu.Unlock() o.saveSagaState(saga) continue } if step.Compensate == nil { saga.mu.Lock() step.Status = "compensated" step.CompletedAt = time.Now().UnixMilli() saga.mu.Unlock() o.MarkExecuted(compExecID) 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<