From ae12264e27701c3ece90f76790e0033f1073b0be 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: Sun, 4 Oct 2026 17:47:19 +0000 Subject: [PATCH] Delete internal/storage/transactions.go --- internal/storage/transactions.go | 2118 ------------------------------ 1 file changed, 2118 deletions(-) delete mode 100644 internal/storage/transactions.go diff --git a/internal/storage/transactions.go b/internal/storage/transactions.go deleted file mode 100644 index 15a23d4..0000000 --- a/internal/storage/transactions.go +++ /dev/null @@ -1,2118 +0,0 @@ -/* - * Copyright 2026 Safronov Grigorii - * - * Licensed under the CDDL, Version 1.0 (the "License"); - * you may not use this file except in compliance with the License. - * - * You may obtain a copy of the License at - * https://opensource.org/licenses/CDDL-1.0 - */ - -// Файл: internal/storage/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<