From f191e5b32896a6c35e30b3bf809dc4d11c09efd8 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: Mon, 5 Oct 2026 21:29:54 +0000 Subject: [PATCH] Update internal/storage/transactions.go --- internal/storage/transactions.go | 247 +++++++++---------------------- 1 file changed, 69 insertions(+), 178 deletions(-) diff --git a/internal/storage/transactions.go b/internal/storage/transactions.go index da8391e..ba3dfc7 100644 --- a/internal/storage/transactions.go +++ b/internal/storage/transactions.go @@ -33,11 +33,11 @@ // Transaction — теперь эти ссылки заменены на GetBatchManager() и *Batch. // // ДОБАВЛЕНО (2026-10, идемпотентность SAGA + выборы лидера): -// - SagaPersistentStorage теперь ведёт персистентную таблицу -// выполненных ExecutionID (saga_executed_keys.json). Перед выполнением -// шага и перед его компенсацией оркестратор проверяет эту таблицу, -// чтобы повторный запуск SAGA после падения не приводил к дублированию -// операций на удалённых узлах. +// - SagaPersistentStorage ведёт персистентную таблицу выполненных +// ExecutionID (saga_executed_keys.json). Перед выполнением шага +// и перед его компенсацией оркестратор проверяет эту таблицу, +// чтобы повторный запуск SAGA после падения не приводил к +// дублированию операций на удалённых узлах. // - SagaOrchestrator получил SagaLeaderElector, который берёт лидера // из Raft (через SagaRaftAccessor). Если Raft недоступен — // используется резервный режим «единственный узел — лидер». @@ -48,8 +48,12 @@ // узел знал, был ли он лидером, и корректно себя вёл до первых // выборов. // - Добавлен интерфейс SagaRaftAccessor, чтобы storage не зависел -// напрямую от hashicorp/raft (это позволяет избежать циклических -// импортов и упрощает тестирование). +// напрямую от hashicorp/raft. +// +// ДОБАВЛЕНО (2026-10, ревизия): +// - GetAllSagas у SagaManager — читает все SAGA (включая завершённые) +// с диска через SagaPersistentStorage.ListAll(). Используется +// REPL-командой "saga list --all". package storage @@ -127,14 +131,12 @@ const ( SagaCleanupInterval = 1 * time.Minute // SagaExecutedKeysFile — персистентная таблица выполненных ExecutionID. - // Нужна для идемпотентности шагов SAGA при повторе после падения. SagaExecutedKeysFile = "saga_executed_keys.json" // SagaLeaderStateFile — персистентное состояние лидера SAGA. SagaLeaderStateFile = "saga_leader_state.json" - // SagaLeaderElectionInterval — как часто опрашиваем Raft на предмет - // лидерства. Меньше — быстрее реакция, но больше нагрузка. + // SagaLeaderElectionInterval — как часто опрашиваем Raft на предмет лидерства. SagaLeaderElectionInterval = 2 * time.Second ) @@ -197,7 +199,7 @@ func crc32(data []byte) uint32 { } // ============================================================================= -// BATCH MANAGER - замена TransactionManager +// BATCH MANAGER — замена TransactionManager // ============================================================================= // BatchManager управляет batch-операциями и WAL. @@ -207,12 +209,9 @@ type BatchManager struct { logger LoggerInterface walPath string stats *BatchStats - aof *AOFManager // append-only file для мутаций вне batch + aof *AOFManager // activeBatches — реестр активных (незакоммиченных) batch'ей. - // Ключ — Batch.ID, значение — *Batch. - // Используется eviction-механизмом (runtime_limits.go), чтобы не - // вытеснять документы, участвующие в незакоммиченной batch-операции. activeBatches sync.Map } @@ -263,7 +262,6 @@ func InitBatchManagerWithConfig(walPath string, config map[string]interface{}) e return } - // Инициализируем AOF для мутаций вне batch. globalBatchManager.aof, err = NewAOFManager( filepath.Join(filepath.Dir(walPath), "aof"), fsyncEnabled, nil) @@ -301,7 +299,6 @@ func GetGlobalStorage() *Storage { return globalStorage } // ============================================================================= // RegisterBatch регистрирует batch как активный. -// Вызывается в NewBatch. func (bm *BatchManager) RegisterBatch(b *Batch) { if bm == nil || b == nil { return @@ -310,7 +307,6 @@ func (bm *BatchManager) RegisterBatch(b *Batch) { } // UnregisterBatch снимает регистрацию batch. -// Вызывается в Commit (успешном или неуспешном). func (bm *BatchManager) UnregisterBatch(b *Batch) { if bm == nil || b == nil { return @@ -318,11 +314,7 @@ func (bm *BatchManager) UnregisterBatch(b *Batch) { bm.activeBatches.Delete(b.ID) } -// IsDocInActiveBatch проверяет, участвует ли документ в активной -// (незакоммиченной) batch-операции. -// -// Используется в runtime_limits.go (eviction), чтобы не вытеснять -// документы, которые в данный момент участвуют в batch. +// IsDocInActiveBatch проверяет, участвует ли документ в активной batch-операции. func (bm *BatchManager) IsDocInActiveBatch(docID, dbName, collName string) bool { if bm == nil { return false @@ -347,7 +339,6 @@ func (bm *BatchManager) IsDocInActiveBatch(docID, dbName, collName string) bool } // GetActiveBatchCount возвращает количество активных batch'ей. -// Полезно для метрик и отладки. func (bm *BatchManager) GetActiveBatchCount() int { if bm == nil { return 0 @@ -376,8 +367,6 @@ func NewBatch(database, collection string) *Batch { Database: database, Collection: collection, } - // Регистрируем batch как активный, чтобы eviction не вытеснил - // документы, участвующие в нём. globalBatchManager.RegisterBatch(b) return b } @@ -426,29 +415,17 @@ func (b *Batch) AddRestore(docID string) { } // 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) @@ -460,13 +437,11 @@ func (b *Batch) Commit() error { 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 { @@ -481,16 +456,11 @@ func (b *Batch) Commit() error { } // 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 @@ -592,11 +562,7 @@ func (b *Batch) apply() error { // TIME-TRAVEL QUERIES (WAL-based) // ============================================================================= -// GetDocumentAtTimestamp возвращает документ на указанный момент времени, -// проигрывая WAL с начала до нужного LSN. -// -// Это O(N) по количеству записей WAL — для in-memory СУБД с небольшим WAL -// это приемлемо. Для больших WAL нужен индекс по времени. +// GetDocumentAtTimestamp возвращает документ на указанный момент времени. func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) { if globalBatchManager == nil { return nil, fmt.Errorf("batch manager not initialized") @@ -661,8 +627,7 @@ func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) { return current, nil } -// GetDocumentVersionsSince возвращает все версии документа, -// созданные после указанного времени. +// GetDocumentVersionsSince возвращает все версии документа после указанного времени. func GetDocumentVersionsSince(docID string, since int64) ([]*Document, error) { if globalBatchManager == nil || globalBatchManager.wal == nil { return nil, fmt.Errorf("batch manager not initialized") @@ -747,8 +712,6 @@ type SegmentedWALManager struct { logger LoggerInterface fsyncEnabled bool - // ackMap — map LSN -> chan struct{} для синхронного подтверждения. - // Ключ — LSN записи. После записи в WAL закрываем канал. ackMu sync.Mutex ackMap map[uint64]chan struct{} } @@ -871,13 +834,9 @@ func (wm *SegmentedWALManager) Write(record *WALRecord) error { } // 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 @@ -981,7 +940,6 @@ func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { wm.mu.Unlock() - // Подтверждаем все записи из batch. wm.ackMu.Lock() for _, record := range batch { if ch, ok := wm.ackMap[record.LSN]; ok { @@ -1148,16 +1106,15 @@ func (im *WALIndexManager) load() error { } return json.Unmarshal(data, &im.index) } - // ============================================================================= -// SAGA — ИНТЕРФЕЙСЫ ДЛЯ RAFT (без прямого импорта hashicorp/raft) +// SAGA — ИНТЕРФЕЙСЫ ДЛЯ RAFT // ============================================================================= // SagaRaftAccessor — минимальный интерфейс к Raft, необходимый // оркестратору SAGA для выборов лидера. // -// Реализуется в internal/cluster (RaftCoordinator), чтобы избежать -// циклических импортов и позволить storage не зависеть от raft напрямую. +// Реализуется в internal/cluster (SagaLeaderBridge), чтобы storage +// не зависел напрямую от hashicorp/raft. type SagaRaftAccessor interface { // IsSagaLeader возвращает true, если текущий узел является лидером // в Raft-группе. @@ -1224,13 +1181,8 @@ type SagaStepState struct { // SagaPersistentStorage — персистентное хранилище SAGA. // // ДОБАВЛЕНО: executedKeys — персистентная таблица уже выполненных -// ExecutionID. Заполняется после успешного выполнения шага -// (в executeSagaInternal) и после успешной компенсации -// (в executeCompensation). Проверяется перед выполнением/компенсацией. -// -// Формат файла — JSON: { "execution_id": timestamp_ms, ... }. -// Это позволяет при рестарте быстро понять, какие шаги уже выполнены, -// и не дублировать их на удалённых узлах. +// ExecutionID. Заполняется после успешного выполнения шага и после +// успешной компенсации. Проверяется перед выполнением/компенсацией. type SagaPersistentStorage struct { baseDir string mu sync.RWMutex @@ -1241,16 +1193,11 @@ type SagaPersistentStorage struct { fsyncMaxRetries int fsyncRetryDelay time.Duration - // executedKeys — таблица выполненных ExecutionID. - // Ключ — ExecutionID, значение — UnixMilli, когда он был выполнен. - // Хранится в памяти + на диске (saga_executed_keys.json). executedMu sync.RWMutex executedKeys map[string]int64 } // NewSagaPersistentStorageWithConfig создаёт хранилище SAGA. -// -// ДОБАВЛЕНО: при создании загружаем executedKeys с диска. func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInterface) (*SagaPersistentStorage, error) { baseDir := SagaStateDir maxCache := 10000 @@ -1282,7 +1229,6 @@ func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInt 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)) @@ -1325,7 +1271,6 @@ func (sps *SagaPersistentStorage) loadExecutedKeys() error { } // saveExecutedKeys сохраняет таблицу выполненных ключей на диск. -// Вызывается под executedMu.RLock или без удержания блокировки. func (sps *SagaPersistentStorage) saveExecutedKeys() error { sps.executedMu.RLock() data, err := json.Marshal(sps.executedKeys) @@ -1364,8 +1309,7 @@ func (sps *SagaPersistentStorage) IsExecuted(executionID string) bool { // MarkExecuted атомарно отмечает ExecutionID как выполненный // и сохраняет таблицу на диск. // -// Возвращает true, если ключ был добавлен впервые (т.е. операция -// действительно новая), и false, если ключ уже был (повтор). +// Возвращает true, если ключ был добавлен впервые. func (sps *SagaPersistentStorage) MarkExecuted(executionID string) bool { sps.executedMu.Lock() if _, exists := sps.executedKeys[executionID]; exists { @@ -1561,21 +1505,13 @@ func (sps *SagaPersistentStorage) LoadLeaderState() (*SagaLeaderState, error) { // ============================================================================= // SagaLeaderElector — выборы лидера SAGA через Raft. -// -// Если Raft доступен, IsLeader() опрашивает RaftAccessor. -// Если Raft недоступен, используется резервный режим: узел считается -// лидером, только если он один в кластере (single-node mode). -// -// Все изменения лидерства сохраняются в SagaPersistentStorage -// (файл saga_leader_state.json), чтобы после рестарта узел знал, -// был ли он лидером. type SagaLeaderElector struct { nodeID string raft SagaRaftAccessor storage *SagaPersistentStorage logger LoggerInterface isLeader atomic.Bool - leaderID atomic.Value // string + leaderID atomic.Value currentTerm atomic.Uint64 observersMu sync.RWMutex @@ -1597,7 +1533,6 @@ func NewSagaLeaderElector(nodeID string, raft SagaRaftAccessor, storage *SagaPer } e.leaderID.Store("") - // Пытаемся загрузить прошлое состояние лидера. if storage != nil { if state, err := storage.LoadLeaderState(); err == nil && state != nil { e.isLeader.Store(state.IsLeader) @@ -1606,15 +1541,12 @@ func NewSagaLeaderElector(nodeID string, raft SagaRaftAccessor, storage *SagaPer } } - // Если Raft доступен, регистрируем наблюдателя, чтобы оперативно - // получать уведомления о смене лидера. if raft != nil { raft.RegisterSagaLeaderObserver(nodeID, func(isLeader bool, leaderID string) { e.updateLeadership(isLeader, leaderID) }) } - // Запускаем фоновый опрос Raft (на случай, если наблюдатель не сработал). e.wg.Add(1) go e.leaderElectionLoop() @@ -1622,7 +1554,6 @@ func NewSagaLeaderElector(nodeID string, raft SagaRaftAccessor, storage *SagaPer } // updateLeadership — центральная точка обновления лидерства. -// Вызывается и наблюдателем, и фоновым опросом. func (e *SagaLeaderElector) updateLeadership(isLeader bool, leaderID string) { oldLeader := e.isLeader.Load() oldLeaderID, _ := e.leaderID.Load().(string) @@ -1636,7 +1567,6 @@ func (e *SagaLeaderElector) updateLeadership(isLeader bool, leaderID string) { } e.currentTerm.Store(term) - // Сохраняем на диск. if e.storage != nil { _ = e.storage.SaveLeaderState(&SagaLeaderState{ LeaderID: leaderID, @@ -1647,7 +1577,6 @@ func (e *SagaLeaderElector) updateLeadership(isLeader bool, leaderID string) { }) } - // Если что-то изменилось — уведомляем наблюдателей. if oldLeader != isLeader || oldLeaderID != leaderID { e.observersMu.RLock() obs := make([]func(isLeader bool, leaderID string), 0, len(e.observers)) @@ -1732,28 +1661,24 @@ func (e *SagaLeaderElector) Stop() { } // ============================================================================= -// SAGA ORCHESTRATOR (упрощённый, без мьютексов на горячем пути) +// SAGA ORCHESTRATOR // ============================================================================= // SagaOrchestrator управляет SAGA. -// -// ДОБАВЛЕНО: leaderElector — выборы лидера через Raft. -// isLeader больше не выставляется безусловно в true: значение -// синхронизируется с leaderElector. 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 - leaderElector *SagaLeaderElector - raftAccessor SagaRaftAccessor + 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. @@ -1765,17 +1690,11 @@ type SagaMetrics struct { TotalCompensated atomic.Uint64 ActiveCount atomic.Int64 RecoveryCount atomic.Uint64 - IdempotentSkips atomic.Uint64 // сколько шагов было пропущено по идемпотентности - NotLeaderRejects atomic.Uint64 // сколько SAGA было отклонено из-за отсутствия лидерства + IdempotentSkips atomic.Uint64 + NotLeaderRejects atomic.Uint64 } -// NewSagaOrchestratorWithConfig создаёт оркестратор. -// -// ДОБАВЛЕНО: второй аргумент raftAccessor может быть nil — -// тогда оркестратор работает в single-node режиме и всегда лидер. -// -// ДОБАВЛЕНО: если raftAccessor != nil, создаётся SagaLeaderElector, -// который опрашивает Raft и обновляет isLeader. +// NewSagaOrchestratorWithConfig создаёт оркестратор без Raft. func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) { return NewSagaOrchestratorWithRaft(cfg, nodeID, nil, logger) } @@ -1814,11 +1733,8 @@ func NewSagaOrchestratorWithRaft(cfg *config.SagaConfig, nodeID string, raftAcce raftAccessor: raftAccessor, } - // Создаём выборщик лидера. o.leaderElector = NewSagaLeaderElector(nodeID, raftAccessor, storage, logger) - // Начальное значение isLeader синхронизируем с выборщиком. - // Если Raft недоступен (nil), считаем себя лидером — single-node режим. if raftAccessor == nil { o.isLeader.Store(true) o.leaderID = nodeID @@ -1828,7 +1744,6 @@ func NewSagaOrchestratorWithRaft(cfg *config.SagaConfig, nodeID string, raftAcce o.leaderID = o.leaderElector.GetLeaderID() } - // Подписываемся на изменения лидерства. o.leaderElector.RegisterObserver("orchestrator", func(isLeader bool, leaderID string) { o.mu.Lock() o.leaderID = leaderID @@ -1844,9 +1759,6 @@ func NewSagaOrchestratorWithRaft(cfg *config.SagaConfig, nodeID string, raftAcce } // IsLeader — является ли узел лидером. -// -// ДОБАВЛЕНО: теперь значение берётся из leaderElector, а не -// из жёстко выставленного флага. func (o *SagaOrchestrator) IsLeader() bool { if o.leaderElector == nil { return true @@ -1868,7 +1780,6 @@ func (o *SagaOrchestrator) GetLeaderID() string { } // IsExecuted проверяет, был ли ExecutionID уже выполнен. -// Делегируется в storage. func (o *SagaOrchestrator) IsExecuted(executionID string) bool { if o.storage == nil { return false @@ -1877,7 +1788,6 @@ func (o *SagaOrchestrator) IsExecuted(executionID string) bool { } // MarkExecuted помечает ExecutionID как выполненный. -// Возвращает true, если это первое выполнение. func (o *SagaOrchestrator) MarkExecuted(executionID string) bool { if o.storage == nil { return true @@ -1896,7 +1806,6 @@ func (o *SagaOrchestrator) recoveryLoop() { for { select { case <-ticker.C: - // Восстановлением занимается только лидер. if !o.IsLeader() { continue } @@ -1914,7 +1823,6 @@ func (o *SagaOrchestrator) orphanCleanupLoop() { for { select { case <-ticker.C: - // Очисткой занимается только лидер. if !o.IsLeader() { continue } @@ -1925,7 +1833,7 @@ func (o *SagaOrchestrator) orphanCleanupLoop() { } } -// cleanupOrphanSteps — очистка "осиротевших" SAGA. +// cleanupOrphanSteps — очистка «осиротевших» SAGA. func (o *SagaOrchestrator) cleanupOrphanSteps() { pending, err := o.storage.ListPending() if err != nil { @@ -1948,13 +1856,7 @@ func (o *SagaOrchestrator) cleanupOrphanSteps() { } } -// compensateOrphanSaga — компенсация "осиротевшей" SAGA. -// -// Проверяем nil на Compensate, сохраняем состояние после -// каждой компенсации, обрабатываем compensation_failed. -// -// ДОБАВЛЕНО: проверяем лидерство — компенсацию выполняет только лидер, -// чтобы два узла не компенсировали одну и ту же SAGA. +// compensateOrphanSaga — компенсация «осиротевшей» SAGA. func (o *SagaOrchestrator) compensateOrphanSaga(state *SagaState) { if !o.IsLeader() { return @@ -2015,8 +1917,6 @@ func (o *SagaOrchestrator) stateToTransaction(state *SagaState) *SagaTransaction } // recoverPendingSagas — восстановление незавершённых SAGA. -// -// ДОБАВЛЕНО: выполняется только лидером. func (o *SagaOrchestrator) recoverPendingSagas() { if !o.IsLeader() { return @@ -2050,9 +1950,6 @@ func (o *SagaOrchestrator) resumeSaga(state *SagaState) error { } // BeginSaga — начать SAGA. -// -// ДОБАВЛЕНО: если узел не лидер — возвращаем ошибку, чтобы клиент -// перенаправил запрос на лидера. func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { if !o.IsLeader() { o.metrics.NotLeaderRejects.Add(1) @@ -2098,8 +1995,6 @@ func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { } // Execute — выполнить SAGA. -// -// ДОБАВЛЕНО: только лидер выполняет SAGA. func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error { if !o.IsLeader() { o.metrics.NotLeaderRejects.Add(1) @@ -2109,16 +2004,6 @@ func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error { } // executeSagaInternal — внутренняя логика выполнения SAGA. -// -// Компенсация идемпотентна и корректно обрабатывает ошибки. -// -// ДОБАВЛЕНО: идемпотентность через storage.IsExecuted / storage.MarkExecuted. -// Перед выполнением шага проверяем, не выполнялся ли он уже. Если да — -// пропускаем (IdempotentSkips++). После успешного выполнения отмечаем -// ExecutionID как выполненный на диске. -// -// ДОБАВЛЕНО: if !o.IsLeader() — прерываем выполнение, если узел потерял -// лидерство в процессе (защита от двух лидеров). func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { if !o.IsLeader() { return fmt.Errorf("node is not the saga leader") @@ -2138,8 +2023,6 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { } for i := saga.CurrentStep; i < len(saga.Steps); i++ { - // Проверяем лидерство на каждой итерации: если потеряли — - // прерываемся, чтобы не конфликтовать с новым лидером. if !o.IsLeader() { return fmt.Errorf("lost saga leadership during execution of %s", saga.ID) } @@ -2155,7 +2038,6 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { step.StartedAt = time.Now().UnixMilli() saga.mu.Unlock() - // ИДЕМПОТЕНТНОСТЬ: если шаг уже выполнялся (по журналу), пропускаем. if o.IsExecuted(step.ExecutionID) || saga.HasExecuted(step.ExecutionID) { saga.mu.Lock() step.Status = "completed" @@ -2219,8 +2101,6 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { return fmt.Errorf("saga %s aborted at step %s: %v", saga.ID, step.Name, err) } - // ИДЕМПОТЕНТНОСТЬ: отмечаем шаг как выполненный на диске, - // чтобы при рестарте не выполнить его повторно. o.MarkExecuted(step.ExecutionID) saga.mu.Lock() @@ -2276,14 +2156,6 @@ func (o *SagaOrchestrator) saveSagaState(saga *SagaTransaction) error { } // executeCompensation — выполнение компенсации. -// -// Идемпотентность: пропускаем уже скомпенсированные шаги. -// Сохраняем состояние после каждой компенсации. -// Обрабатываем compensation_failed корректно. -// -// ДОБАВЛЕНО: проверяем лидерство в начале и на каждой итерации. -// ДОБАВЛЕНО: идемпотентность через storage.IsExecuted на -// execution_id компенсации (используем суффикс ":compensate"). func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep int) error { if !o.IsLeader() { return fmt.Errorf("node is not the saga leader") @@ -2313,7 +2185,6 @@ func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep continue } - // ИДЕМПОТЕНТНОСТЬ: проверяем, не компенсировали ли уже. compExecID := step.ExecutionID + ":compensate" if o.IsExecuted(compExecID) { saga.mu.Lock() @@ -2325,7 +2196,6 @@ func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep } if step.Compensate == nil { - // Нет функции компенсации — считаем шаг скомпенсированным. saga.mu.Lock() step.Status = "compensated" step.CompletedAt = time.Now().UnixMilli() @@ -2355,7 +2225,6 @@ func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep o.saveSagaState(saga) return fmt.Errorf("compensation for step %s failed: %v", step.Name, err) } - // ИДЕМПОТЕНТНОСТЬ: отмечаем компенсацию как выполненную. o.MarkExecuted(compExecID) saga.mu.Lock() step.Status = "compensated" @@ -2387,8 +2256,6 @@ func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) { } // GetMetrics — метрики SAGA. -// -// ДОБАВЛЕНО: idempotent_skips, not_leader_rejects, leader_id. func (o *SagaOrchestrator) GetMetrics() map[string]interface{} { return map[string]interface{}{ "total_started": o.metrics.TotalStarted.Load(), @@ -2654,6 +2521,30 @@ func (sm *SagaManager) GetActiveSagas() []*SagaTransaction { return result } +// GetAllSagas возвращает все SAGA транзакции — и активные, и завершённые. +// +// ДОБАВЛЕНО (2026-10): используется REPL-командой "saga list --all". +// В отличие от GetActiveSagas, читает состояние с диска через +// SagaPersistentStorage.ListAll() и конвертирует каждую запись в +// SagaTransaction через stateToTransaction. +func (sm *SagaManager) GetAllSagas() []*SagaTransaction { + result := make([]*SagaTransaction, 0) + if sm.orchestrator == nil || sm.orchestrator.storage == nil { + return result + } + states, err := sm.orchestrator.storage.ListAll() + if err != nil { + return result + } + for _, st := range states { + if st == nil { + continue + } + result = append(result, sm.orchestrator.stateToTransaction(st)) + } + return result +} + // Stop — остановка. func (sm *SagaManager) Stop() { close(sm.stopChan) @@ -2721,7 +2612,7 @@ func StopBatchManager() error { // СОВМЕСТИМОСТЬ: псевдонимы для старого API // ============================================================================= -// GetTransactionStats — псевдоним для GetBatchStats (обратная совместимость). +// GetTransactionStats — псевдоним для GetBatchStats. func GetTransactionStats() map[string]interface{} { return GetBatchStats() -} +} \ No newline at end of file