Files
futriix/internal/storage/transactions.go
T

2600 lines
76 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/*
* 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".
//
// ИЗМЕНЕНО (2026-10): удалён AOF. BatchManager больше не содержит
// поле aof, и Commit не пишет в AOF. Единственный источник истины — WAL.
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
// 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
}
})
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
}
}
}
// 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)
}
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<<retry)) * time.Millisecond
time.Sleep(backoff)
}
saga.mu.Lock()
step.CompletedAt = time.Now().UnixMilli()
saga.mu.Unlock()
if err != nil {
saga.mu.Lock()
step.Status = "failed"
saga.Status = "compensating"
saga.UpdatedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
compensationErr := o.executeCompensation(saga, i)
if compensationErr != nil {
saga.mu.Lock()
saga.Status = "compensation_failed"
saga.UpdatedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
o.metrics.TotalFailed.Add(1)
o.metrics.ActiveCount.Add(-1)
return fmt.Errorf("saga %s compensation failed: %v", saga.ID, compensationErr)
}
saga.mu.Lock()
saga.Status = "aborted"
saga.UpdatedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
o.metrics.TotalAborted.Add(1)
o.metrics.ActiveCount.Add(-1)
return fmt.Errorf("saga %s aborted at step %s: %v", saga.ID, step.Name, err)
}
o.MarkExecuted(step.ExecutionID)
saga.mu.Lock()
step.Status = "completed"
saga.MarkExecuted(step.ExecutionID)
saga.UpdatedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
}
saga.mu.Lock()
saga.Status = "completed"
saga.CompletedAt = time.Now().UnixMilli()
saga.UpdatedAt = saga.CompletedAt
saga.mu.Unlock()
o.saveSagaState(saga)
o.metrics.TotalCompleted.Add(1)
o.metrics.ActiveCount.Add(-1)
return nil
}
// saveSagaState сохраняет состояние SAGA.
func (o *SagaOrchestrator) saveSagaState(saga *SagaTransaction) error {
saga.mu.RLock()
state := &SagaState{
ID: saga.ID,
Status: saga.Status,
CurrentStep: saga.CurrentStep,
Data: saga.Data,
CreatedAt: saga.CreatedAt,
UpdatedAt: saga.UpdatedAt,
CompletedAt: saga.CompletedAt,
CompensationExecuted: saga.compensationExecuted,
NodeID: o.nodeID,
ExecutionID: fmt.Sprintf("%s_%d", saga.ID, saga.CreatedAt),
}
state.Steps = make([]SagaStepState, len(saga.Steps))
for i, step := range saga.Steps {
state.Steps[i] = SagaStepState{
ID: step.ID,
Name: step.Name,
Status: step.Status,
Data: step.Data,
StartedAt: step.StartedAt,
CompletedAt: step.CompletedAt,
ExecutionID: step.ExecutionID,
RetryCount: step.RetryCount,
LastError: step.LastError,
}
}
saga.mu.RUnlock()
return o.storage.Save(state)
}
// executeCompensation — выполнение компенсации.
func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep int) error {
if !o.IsLeader() {
return fmt.Errorf("node is not the saga leader")
}
saga.mu.Lock()
if saga.compensationExecuted {
saga.mu.Unlock()
return nil
}
saga.mu.Unlock()
for j := failedStep; j >= 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<<retry)) * time.Millisecond
time.Sleep(backoff)
}
if err != nil {
saga.mu.Lock()
step.Status = "compensation_failed"
step.LastError = err.Error()
saga.mu.Unlock()
o.saveSagaState(saga)
return fmt.Errorf("compensation for step %s failed: %v", step.Name, err)
}
o.MarkExecuted(compExecID)
saga.mu.Lock()
step.Status = "compensated"
step.CompletedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
}
saga.mu.Lock()
saga.compensationExecuted = true
saga.mu.Unlock()
o.metrics.TotalCompensated.Add(1)
return nil
}
// GetSaga — получить SAGA.
func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) {
if val, ok := o.activeSagas.Load(id); ok {
return val.(*SagaTransaction), nil
}
state, err := o.storage.Load(id)
if err != nil {
return nil, err
}
if state == nil {
return nil, fmt.Errorf("saga %s not found", id)
}
return o.stateToTransaction(state), nil
}
// GetMetrics — метрики SAGA.
func (o *SagaOrchestrator) GetMetrics() map[string]interface{} {
return map[string]interface{}{
"total_started": o.metrics.TotalStarted.Load(),
"total_completed": o.metrics.TotalCompleted.Load(),
"total_aborted": o.metrics.TotalAborted.Load(),
"total_failed": o.metrics.TotalFailed.Load(),
"total_compensated": o.metrics.TotalCompensated.Load(),
"active_count": o.metrics.ActiveCount.Load(),
"recovery_count": o.metrics.RecoveryCount.Load(),
"idempotent_skips": o.metrics.IdempotentSkips.Load(),
"not_leader_rejects": o.metrics.NotLeaderRejects.Load(),
"is_leader": o.IsLeader(),
"leader_id": o.GetLeaderID(),
"node_id": o.nodeID,
}
}
// Stop — остановка оркестратора.
func (o *SagaOrchestrator) Stop() {
close(o.stopChan)
if o.leaderElector != nil {
o.leaderElector.Stop()
}
o.wg.Wait()
if o.logger != nil {
o.logger.Info("Saga orchestrator stopped")
}
}
// =============================================================================
// SAGA TRANSACTION
// =============================================================================
// SagaStep — шаг SAGA.
type SagaStep struct {
ID string `json:"id"`
Name string `json:"name"`
Execute func() error `json:"-"`
Compensate func() error `json:"-"`
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"`
}
// SagaTransaction — SAGA-транзакция.
type SagaTransaction struct {
ID string `json:"id"`
Steps []*SagaStep `json:"steps"`
CurrentStep int `json:"current_step"`
Status string `json:"status"`
CreatedAt int64 `json:"created_at"`
UpdatedAt int64 `json:"updated_at"`
CompletedAt int64 `json:"completed_at,omitempty"`
Data map[string]interface{} `json:"data"`
mu sync.RWMutex
executedOps map[string]bool `json:"-"`
compensationExecuted bool `json:"-"`
}
// AddStep добавляет шаг.
func (s *SagaTransaction) AddStep(name string, execute, compensate func() error, data map[string]interface{}) *SagaTransaction {
s.mu.Lock()
defer s.mu.Unlock()
stepID := fmt.Sprintf("%s_step_%d_%d", s.ID, len(s.Steps), time.Now().UnixNano())
step := &SagaStep{
ID: stepID,
Name: name,
Execute: execute,
Compensate: compensate,
Status: "pending",
Data: data,
StartedAt: time.Now().UnixMilli(),
ExecutionID: stepID,
RetryCount: 0,
}
s.Steps = append(s.Steps, step)
return s
}
// SetData устанавливает данные.
func (s *SagaTransaction) SetData(key string, value interface{}) {
s.mu.Lock()
defer s.mu.Unlock()
s.Data[key] = value
}
// GetData получает данные.
func (s *SagaTransaction) GetData(key string) (interface{}, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
val, ok := s.Data[key]
return val, ok
}
// HasExecuted — был ли шаг выполнен.
func (s *SagaTransaction) HasExecuted(operationID string) bool {
s.mu.RLock()
defer s.mu.RUnlock()
_, ok := s.executedOps[operationID]
return ok
}
// MarkExecuted — отметить шаг как выполненный.
func (s *SagaTransaction) MarkExecuted(operationID string) {
s.executedOps[operationID] = true
}
// =============================================================================
// SAGA MANAGER (обёртка)
// =============================================================================
// SagaManager — обёртка над оркестратором.
type SagaManager struct {
orchestrator *SagaOrchestrator
logger LoggerInterface
stopChan chan struct{}
wg sync.WaitGroup
maxRetries int
}
// NewSagaManager создаёт SAGA-менеджер (single-node, без Raft).
func NewSagaManager(logger LoggerInterface) *SagaManager {
orchestrator, err := NewSagaOrchestratorWithConfig(nil, "default-node", logger)
if err != nil {
if logger != nil {
logger.Error(fmt.Sprintf("Failed to create saga orchestrator: %v", err))
}
return &SagaManager{
logger: logger,
stopChan: make(chan struct{}),
maxRetries: 3,
}
}
return &SagaManager{
orchestrator: orchestrator,
logger: logger,
stopChan: make(chan struct{}),
maxRetries: 3,
}
}
// NewSagaManagerWithRaft создаёт SAGA-менеджер с привязкой к Raft.
func NewSagaManagerWithRaft(nodeID string, raftAccessor SagaRaftAccessor, logger LoggerInterface) *SagaManager {
orchestrator, err := NewSagaOrchestratorWithRaft(nil, nodeID, raftAccessor, logger)
if err != nil {
if logger != nil {
logger.Error(fmt.Sprintf("Failed to create saga orchestrator with raft: %v", err))
}
return &SagaManager{
logger: logger,
stopChan: make(chan struct{}),
maxRetries: 3,
}
}
return &SagaManager{
orchestrator: orchestrator,
logger: logger,
stopChan: make(chan struct{}),
maxRetries: 3,
}
}
// BeginSaga — начать SAGA.
func (sm *SagaManager) BeginSaga(id string) *SagaTransaction {
if sm.orchestrator != nil {
saga, err := sm.orchestrator.BeginSaga(id)
if err != nil {
if sm.logger != nil {
sm.logger.Error(fmt.Sprintf("Failed to begin saga: %v", err))
}
return &SagaTransaction{
ID: id,
Status: "pending",
CreatedAt: time.Now().UnixMilli(),
UpdatedAt: time.Now().UnixMilli(),
Data: make(map[string]interface{}),
executedOps: make(map[string]bool),
}
}
return saga
}
return &SagaTransaction{
ID: id,
Steps: make([]*SagaStep, 0),
CurrentStep: 0,
Status: "pending",
CreatedAt: time.Now().UnixMilli(),
UpdatedAt: time.Now().UnixMilli(),
Data: make(map[string]interface{}),
executedOps: make(map[string]bool),
}
}
// Execute — выполнить SAGA.
func (sm *SagaManager) Execute(saga *SagaTransaction) error {
if sm.orchestrator != nil {
return sm.orchestrator.Execute(saga)
}
return fmt.Errorf("orchestrator not available")
}
// GetSaga — получить SAGA.
func (sm *SagaManager) GetSaga(id string) (*SagaTransaction, error) {
if sm.orchestrator != nil {
return sm.orchestrator.GetSaga(id)
}
return nil, fmt.Errorf("saga %s not found", id)
}
// GetSagaStatus — статус SAGA.
func (sm *SagaManager) GetSagaStatus(id string) (string, error) {
saga, err := sm.GetSaga(id)
if err != nil {
return "", err
}
saga.mu.RLock()
defer saga.mu.RUnlock()
return saga.Status, nil
}
// GetOrchestrator — получить оркестратор.
func (sm *SagaManager) GetOrchestrator() *SagaOrchestrator {
return sm.orchestrator
}
// GetMetrics — метрики SAGA.
func (sm *SagaManager) GetMetrics() map[string]interface{} {
if sm.orchestrator != nil {
return sm.orchestrator.GetMetrics()
}
return map[string]interface{}{
"total_started": 0,
"total_completed": 0,
"total_aborted": 0,
"total_failed": 0,
"active_count": 0,
"is_leader": false,
"leader_id": "",
"node_id": "",
}
}
// GetActiveSagas — активные SAGA.
func (sm *SagaManager) GetActiveSagas() []*SagaTransaction {
result := make([]*SagaTransaction, 0)
if sm.orchestrator == nil {
return result
}
sm.orchestrator.activeSagas.Range(func(key, value interface{}) bool {
saga := value.(*SagaTransaction)
saga.mu.RLock()
status := saga.Status
saga.mu.RUnlock()
if status == "pending" || status == "running" {
result = append(result, saga)
}
return true
})
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)
sm.wg.Wait()
if sm.orchestrator != nil {
sm.orchestrator.Stop()
}
}
// =============================================================================
// ГЛОБАЛЬНЫЕ SAGA
// =============================================================================
var globalSagaOrchestrator *SagaOrchestrator
var sagaOrchestratorMu sync.RWMutex
// SetGlobalSagaOrchestrator — установить глобальный оркестратор.
func SetGlobalSagaOrchestrator(o *SagaOrchestrator) {
sagaOrchestratorMu.Lock()
defer sagaOrchestratorMu.Unlock()
globalSagaOrchestrator = o
}
// GetGlobalSagaOrchestrator — получить глобальный оркестратор.
func GetGlobalSagaOrchestrator() *SagaOrchestrator {
sagaOrchestratorMu.RLock()
defer sagaOrchestratorMu.RUnlock()
return globalSagaOrchestrator
}
// =============================================================================
// СТАТИСТИКА
// =============================================================================
// GetBatchStats возвращает статистику batch-операций.
func GetBatchStats() map[string]interface{} {
if globalBatchManager == nil {
return map[string]interface{}{"error": "batch manager not initialized"}
}
return map[string]interface{}{
"total_batches": globalBatchManager.stats.TotalBatches.Load(),
"total_operations": globalBatchManager.stats.TotalOperations.Load(),
"total_failed": globalBatchManager.stats.TotalFailed.Load(),
"active_batches": globalBatchManager.GetActiveBatchCount(),
"uptime_seconds": time.Since(globalBatchManager.stats.StartTime).Seconds(),
}
}
// StopBatchManager останавливает batch-менеджер.
func StopBatchManager() error {
if globalBatchManager == nil {
return nil
}
if globalBatchManager.wal != nil {
globalBatchManager.wal.Sync()
globalBatchManager.wal.Close()
}
return nil
}
// =============================================================================
// СОВМЕСТИМОСТЬ: псевдонимы для старого API
// =============================================================================
// GetTransactionStats — псевдоним для GetBatchStats.
func GetTransactionStats() map[string]interface{} {
return GetBatchStats()
}