Files
futriix/internal/storage/transactions.go
T

2728 lines
86 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 (это позволяет избежать циклических
// импортов и упрощает тестирование).
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.
// Нужна для идемпотентности шагов SAGA при повторе после падения.
SagaExecutedKeysFile = "saga_executed_keys.json"
// SagaLeaderStateFile — персистентное состояние лидера SAGA.
SagaLeaderStateFile = "saga_leader_state.json"
// SagaLeaderElectionInterval — как часто опрашиваем Raft на предмет
// лидерства. Меньше — быстрее реакция, но больше нагрузка.
SagaLeaderElectionInterval = 2 * time.Second
)
// =============================================================================
// CRC32
// =============================================================================
var crc32Table = [256]uint32{
0x00000000, 0x77073096, 0xee0e612c, 0x990951ba, 0x076dc419, 0x706af48f,
0xe963a535, 0x9e6495a3, 0x0edb8832, 0x79dcb8a4, 0xe0d5e91e, 0x97d2d988,
0x09b64c2b, 0x7eb17cbd, 0xe7b82d07, 0x90bf1d91, 0x1db71064, 0x6ab020f2,
0xf3b97148, 0x84be41de, 0x1adad47d, 0x6ddde4eb, 0xf4d4b551, 0x83d385c7,
0x136c9856, 0x646ba8c0, 0xfd62f97a, 0x8a65c9ec, 0x14015c4f, 0x63066cd9,
0xfa0f3d63, 0x8d080df5, 0x3b6e20c8, 0x4c69105e, 0xd56041e4, 0xa2677172,
0x3c03e4d1, 0x4b04d447, 0xd20d85fd, 0xa50ab56b, 0x35b5a8fa, 0x42b2986c,
0xdbbbc9d6, 0xacbcf940, 0x32d86ce3, 0x45df5c75, 0xdcd60dcf, 0xabd13d59,
0x26d930ac, 0x51de003a, 0xc8d75180, 0xbfd06116, 0x21b4f4b5, 0x56b3c423,
0xcfba9599, 0xb8bda50f, 0x2802b89e, 0x5f058808, 0xc60cd9b2, 0xb10be924,
0x2f6f7c87, 0x58684c11, 0xc1611dab, 0xb6662d3d, 0x76dc4190, 0x01db7106,
0x98d220bc, 0xefd5102a, 0x71b18589, 0x06b6b51f, 0x9fbfe4a5, 0xe8b8d433,
0x7807c9a2, 0x0f00f934, 0x9609a88e, 0xe10e9818, 0x7f6a0dbb, 0x086d3d2d,
0x91646c97, 0xe6635c01, 0x6b6b51f4, 0x1c6c6162, 0x856530d8, 0xf262004e,
0x6c0695ed, 0x1b01a57b, 0x8208f4c1, 0xf50fc457, 0x65b0d9c6, 0x12b7e950,
0x8bbeb8ea, 0xfcb9887c, 0x62dd1ddf, 0x15da2d49, 0x8cd37cf3, 0xfbd44c65,
0x4db26158, 0x3ab551ce, 0xa3bc0074, 0xd4bb30e2, 0x4adfa541, 0x3dd895d7,
0xa4d1c46d, 0xd3d6f4fb, 0x4369e96a, 0x346ed9fc, 0xad678846, 0xda60b8d0,
0x44042d73, 0x33031de5, 0xaa0a4c5f, 0xdd0d7cc9, 0x5005713c, 0x270241aa,
0xbe0b1010, 0xc90c2086, 0x5768b525, 0x206f85b3, 0xb966d409, 0xce61e49f,
0x5edef90e, 0x29d9c998, 0xb0d09822, 0xc7d7a8b4, 0x59b33d17, 0x2eb40d81,
0xb7bd5c3b, 0xc0ba6cad, 0xedb88320, 0x9abfb3b6, 0x03b6e20c, 0x74b1d29a,
0xead54739, 0x9dd277af, 0x04db2615, 0x73dc1683, 0xe3630b12, 0x94643b84,
0x0d6d6a3e, 0x7a6a5aa8, 0xe40ecf0b, 0x9309ff9d, 0x0a00ae27, 0x7d079eb1,
0xf00f9344, 0x8708a3d2, 0x1e01f268, 0x6906c2fe, 0xf762575d, 0x806567cb,
0x196c3671, 0x6e6b06e7, 0xfed41b76, 0x89d32be0, 0x10da7a5a, 0x67dd4acc,
0xf9b9df6f, 0x8ebeeff9, 0x17b7be43, 0x60b08ed5, 0xd6d6a3e8, 0xa1d1937e,
0x38d8c2c4, 0x4fdff252, 0xd1bb67f1, 0xa6bc5767, 0x3fb506dd, 0x48b2364b,
0xd80d2bda, 0xaf0a1a4c, 0x36034af6, 0x41047a60, 0xdf60efc3, 0xa867df55,
0x316e8eef, 0x4669be79, 0xcb61b38c, 0xbc66831a, 0x256fd2a0, 0x5268e236,
0xcc0c7795, 0xbb0b4703, 0x220216b9, 0x5505262f, 0xc5ba3bbe, 0xb2bd0b28,
0x2bb45a92, 0x5cb36a04, 0xc2d7ffa7, 0xb5d0cf31, 0x2cd99e8b, 0x5bdeae1d,
0x9b64c2b0, 0xec63f226, 0x756aa39c, 0x026d930a, 0x9c0906a9, 0xeb0e363f,
0x72076785, 0x05005713, 0x95bf4a82, 0xe2b87a14, 0x7bb12bae, 0x0cb61b38,
0x92d28e9b, 0xe5d5be0d, 0x7cdcefb7, 0x0bdbdf21, 0x86d3d2d4, 0xf1d4e242,
0x68ddb3f8, 0x1fda836e, 0x81be16cd, 0xf6b9265b, 0x6fb077e1, 0x18b74777,
0x88085ae6, 0xff0f6a70, 0x66063bca, 0x11010b5c, 0x8f659eff, 0xf862ae69,
0x616bffd3, 0x166ccf45, 0xa00ae278, 0xd70dd2ee, 0x4e048354, 0x3903b3c2,
0xa7672661, 0xd06016f7, 0x4969474d, 0x3e6e77db, 0xaed16a4a, 0xd9d65adc,
0x40df0b66, 0x37d83bf8, 0xa9bcae53, 0xdebb9ec5, 0x47b2cf7f, 0x30b5ffe9,
0xbdbdf21c, 0xcabac28a, 0x53b39330, 0x24b4a3a6, 0xbad03605, 0xcdd70693,
0x54de5729, 0x23d967bf, 0xb3667a2e, 0xc4614ab8, 0x5d681b02, 0x2a6f2b94,
0xb40bbe37, 0xc30c8ea1, 0x5a05df1b, 0x2d02ef8d,
}
func crc32(data []byte) uint32 {
crc := uint32(0xFFFFFFFF)
for _, b := range data {
crc = (crc >> 8) ^ crc32Table[(crc^uint32(b))&0xFF]
}
return crc ^ 0xFFFFFFFF
}
// =============================================================================
// BATCH MANAGER - замена TransactionManager
// =============================================================================
// BatchManager управляет batch-операциями и WAL.
type BatchManager struct {
nextBatchID atomic.Uint64
wal *SegmentedWALManager
logger LoggerInterface
walPath string
stats *BatchStats
aof *AOFManager // append-only file для мутаций вне batch
// activeBatches — реестр активных (незакоммиченных) batch'ей.
// Ключ — Batch.ID, значение — *Batch.
// Используется eviction-механизмом (runtime_limits.go), чтобы не
// вытеснять документы, участвующие в незакоммиченной batch-операции.
activeBatches sync.Map
}
// BatchStats — статистика batch-операций.
type BatchStats struct {
TotalBatches atomic.Uint64
TotalOperations atomic.Uint64
TotalFailed atomic.Uint64
StartTime time.Time
}
// =============================================================================
// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ
// =============================================================================
var (
globalBatchManager *BatchManager
batchManagerOnce sync.Once
globalStorage *Storage
)
// InitBatchManager инициализирует batch-менеджер.
func InitBatchManager(walPath string) error {
return InitBatchManagerWithConfig(walPath, nil)
}
// InitBatchManagerWithConfig инициализирует batch-менеджер с конфигурацией.
func InitBatchManagerWithConfig(walPath string, config map[string]interface{}) error {
var err error
batchManagerOnce.Do(func() {
fsyncEnabled := true
if config != nil {
if v, ok := config["fsync_enabled"].(bool); ok {
fsyncEnabled = v
}
}
globalBatchManager = &BatchManager{
walPath: walPath,
stats: &BatchStats{StartTime: time.Now()},
}
globalBatchManager.nextBatchID.Store(1)
globalBatchManager.wal, err = NewSegmentedWALManager(
filepath.Dir(walPath), fsyncEnabled, nil)
if err != nil {
return
}
// Инициализируем AOF для мутаций вне batch.
globalBatchManager.aof, err = NewAOFManager(
filepath.Join(filepath.Dir(walPath), "aof"),
fsyncEnabled, nil)
if err != nil {
return
}
})
return err
}
// GetBatchManager возвращает глобальный batch-менеджер.
func GetBatchManager() *BatchManager { return globalBatchManager }
// SetBatchLogger устанавливает логгер.
func SetBatchLogger(logger LoggerInterface) {
if globalBatchManager != nil {
globalBatchManager.logger = logger
if globalBatchManager.wal != nil {
globalBatchManager.wal.logger = logger
}
if globalBatchManager.aof != nil {
globalBatchManager.aof.SetLogger(logger)
}
}
}
// SetGlobalStorage устанавливает глобальное хранилище.
func SetGlobalStorage(s *Storage) { globalStorage = s }
// GetGlobalStorage возвращает глобальное хранилище.
func GetGlobalStorage() *Storage { return globalStorage }
// =============================================================================
// РЕЕСТР АКТИВНЫХ BATCH'ЕЙ
// =============================================================================
// RegisterBatch регистрирует batch как активный.
// Вызывается в NewBatch.
func (bm *BatchManager) RegisterBatch(b *Batch) {
if bm == nil || b == nil {
return
}
bm.activeBatches.Store(b.ID, b)
}
// UnregisterBatch снимает регистрацию batch.
// Вызывается в Commit (успешном или неуспешном).
func (bm *BatchManager) UnregisterBatch(b *Batch) {
if bm == nil || b == nil {
return
}
bm.activeBatches.Delete(b.ID)
}
// IsDocInActiveBatch проверяет, участвует ли документ в активной
// (незакоммиченной) batch-операции.
//
// Используется в runtime_limits.go (eviction), чтобы не вытеснять
// документы, которые в данный момент участвуют в batch.
func (bm *BatchManager) IsDocInActiveBatch(docID, dbName, collName string) bool {
if bm == nil {
return false
}
found := false
bm.activeBatches.Range(func(key, value interface{}) bool {
b, ok := value.(*Batch)
if !ok || b == nil {
return true
}
for _, op := range b.Operations {
if op.DocumentID == docID &&
op.Database == dbName &&
op.Collection == collName {
found = true
return false
}
}
return true
})
return found
}
// GetActiveBatchCount возвращает количество активных batch'ей.
// Полезно для метрик и отладки.
func (bm *BatchManager) GetActiveBatchCount() int {
if bm == nil {
return 0
}
count := 0
bm.activeBatches.Range(func(_, _ interface{}) bool {
count++
return true
})
return count
}
// =============================================================================
// BATCH API
// =============================================================================
// NewBatch создаёт новый batch.
func NewBatch(database, collection string) *Batch {
if globalBatchManager == nil {
_ = InitBatchManager("futriis.wal")
}
b := &Batch{
ID: TransactionID(globalBatchManager.nextBatchID.Add(1) - 1),
Operations: make([]BatchOperation, 0, 16),
Timestamp: time.Now().UnixMilli(),
Database: database,
Collection: collection,
}
// Регистрируем batch как активный, чтобы eviction не вытеснил
// документы, участвующие в нём.
globalBatchManager.RegisterBatch(b)
return b
}
// AddInsert добавляет операцию вставки.
func (b *Batch) AddInsert(doc *Document) {
b.Operations = append(b.Operations, BatchOperation{
Type: "insert",
Database: b.Database,
Collection: b.Collection,
DocumentID: doc.ID,
Data: doc.GetFields(),
Version: doc.Version,
})
}
// AddUpdate добавляет операцию обновления.
func (b *Batch) AddUpdate(docID string, updates map[string]interface{}) {
b.Operations = append(b.Operations, BatchOperation{
Type: "update",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
Data: updates,
})
}
// AddDelete добавляет операцию удаления.
func (b *Batch) AddDelete(docID string) {
b.Operations = append(b.Operations, BatchOperation{
Type: "delete",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
})
}
// AddRestore добавляет операцию восстановления.
func (b *Batch) AddRestore(docID string) {
b.Operations = append(b.Operations, BatchOperation{
Type: "restore",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
})
}
// Commit применяет batch атомарно.
//
// Гарантия атомарности:
// 1. Сериализуем весь batch и пишем в WAL (fsync).
// 2. Только после успешной записи в WAL применяем операции.
// 3. Если применение падает на середине — мы можем восстановить
// состояние, проиграв WAL (операции идемпотентны или
// откатываются вручную).
//
// В любом случае (успех или ошибка) batch снимается с регистрации
// в activeBatches.
func (b *Batch) Commit() error {
if globalBatchManager == nil {
return fmt.Errorf("batch manager not initialized")
}
// Снимаем регистрацию в любом случае — batch короткоживущий.
defer globalBatchManager.UnregisterBatch(b)
if len(b.Operations) == 0 {
return nil
}
// 1. Пишем batch в WAL.
data, err := json.Marshal(b)
if err != nil {
return fmt.Errorf("failed to marshal batch: %v", err)
}
record := &WALRecord{Type: 1, Data: data}
if err := globalBatchManager.wal.WriteSync(record); err != nil {
globalBatchManager.stats.TotalFailed.Add(1)
return fmt.Errorf("failed to write batch to WAL: %v", err)
}
// 2. Применяем операции.
if err := b.apply(); err != nil {
globalBatchManager.stats.TotalFailed.Add(1)
return fmt.Errorf("batch apply failed: %v", err)
}
// 3. Также пишем в AOF для быстрого восстановления (опционально).
if globalBatchManager.aof != nil {
if err := globalBatchManager.aof.AppendBatch(b); err != nil {
if globalBatchManager.logger != nil {
globalBatchManager.logger.Warn(fmt.Sprintf("AOF append failed: %v", err))
}
}
}
globalBatchManager.stats.TotalBatches.Add(1)
globalBatchManager.stats.TotalOperations.Add(uint64(len(b.Operations)))
return nil
}
// apply применяет все операции batch.
//
// Если операция падает, мы пытаемся откатить уже применённые
// (best-effort). Полная атомарность гарантируется только на уровне
// WAL: при восстановлении batch либо весь применится, либо нет.
func (b *Batch) apply() error {
if globalStorage == nil {
return fmt.Errorf("storage not initialized")
}
// Собираем откат для уже применённых операций.
type appliedOp struct {
opType string
db string
coll string
docID string
oldDoc *Document
}
applied := make([]appliedOp, 0, len(b.Operations))
rollback := func() {
for i := len(applied) - 1; i >= 0; i-- {
ao := applied[i]
db, err := globalStorage.GetDatabase(ao.db)
if err != nil {
continue
}
coll, err := db.GetCollection(ao.coll)
if err != nil {
continue
}
switch ao.opType {
case "insert":
coll.PermanentDelete(ao.docID)
case "update":
if ao.oldDoc != nil {
coll.Update(ao.docID, ao.oldDoc.GetFields())
}
case "delete":
coll.RestoreDeleted(ao.docID)
}
}
}
for _, op := range b.Operations {
db, err := globalStorage.GetDatabase(op.Database)
if err != nil {
rollback()
return fmt.Errorf("database not found: %s", op.Database)
}
coll, err := db.GetCollection(op.Collection)
if err != nil {
rollback()
return fmt.Errorf("collection not found: %s", op.Collection)
}
switch op.Type {
case "insert":
doc := NewDocumentWithID(op.DocumentID)
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
if err := coll.Insert(doc); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "insert", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
case "update":
oldDoc, _ := coll.Find(op.DocumentID)
var oldCopy *Document
if oldDoc != nil {
oldCopy = oldDoc.Clone()
}
if err := coll.Update(op.DocumentID, op.Data); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "update", db: op.Database, coll: op.Collection,
docID: op.DocumentID, oldDoc: oldCopy,
})
case "delete":
if err := coll.Delete(op.DocumentID); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "delete", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
case "restore":
if err := coll.RestoreDeleted(op.DocumentID); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "restore", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
}
}
return nil
}
// =============================================================================
// TIME-TRAVEL QUERIES (WAL-based)
// =============================================================================
// GetDocumentAtTimestamp возвращает документ на указанный момент времени,
// проигрывая WAL с начала до нужного LSN.
//
// Это O(N) по количеству записей WAL — для in-memory СУБД с небольшим WAL
// это приемлемо. Для больших WAL нужен индекс по времени.
func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) {
if globalBatchManager == nil {
return nil, fmt.Errorf("batch manager not initialized")
}
if globalBatchManager.wal == nil {
return nil, fmt.Errorf("WAL not initialized")
}
records, err := globalBatchManager.wal.ReadAll()
if err != nil {
return nil, fmt.Errorf("failed to read WAL: %v", err)
}
var current *Document
for _, rec := range records {
if rec.Timestamp > timestamp {
break
}
if rec.Type != 1 {
continue
}
var batch Batch
if err := json.Unmarshal(rec.Data, &batch); err != nil {
continue
}
for _, op := range batch.Operations {
if op.DocumentID != docID {
continue
}
switch op.Type {
case "insert":
doc := NewDocumentWithID(op.DocumentID)
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
doc.CreatedAt = batch.Timestamp
doc.UpdatedAt = batch.Timestamp
current = doc
case "update":
if current != nil {
current.Update(op.Data)
current.UpdatedAt = batch.Timestamp
}
case "delete":
if current != nil {
current.SoftDelete()
current.DeletedAt = batch.Timestamp
}
case "restore":
if current != nil {
current.Restore()
current.UpdatedAt = batch.Timestamp
}
}
}
}
if current == nil {
return nil, fmt.Errorf("document %s not found at timestamp %d", docID, timestamp)
}
return current, nil
}
// GetDocumentVersionsSince возвращает все версии документа,
// созданные после указанного времени.
func GetDocumentVersionsSince(docID string, since int64) ([]*Document, error) {
if globalBatchManager == nil || globalBatchManager.wal == nil {
return nil, fmt.Errorf("batch manager not initialized")
}
records, err := globalBatchManager.wal.ReadAll()
if err != nil {
return nil, err
}
versions := make([]*Document, 0)
for _, rec := range records {
if rec.Timestamp <= since || rec.Type != 1 {
continue
}
var batch Batch
if err := json.Unmarshal(rec.Data, &batch); err != nil {
continue
}
for _, op := range batch.Operations {
if op.DocumentID != docID {
continue
}
doc := NewDocumentWithID(op.DocumentID)
if op.Type == "insert" || op.Type == "update" {
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
doc.UpdatedAt = batch.Timestamp
versions = append(versions, doc)
}
}
}
return versions, nil
}
// =============================================================================
// WAL MANAGER (Segmented)
// =============================================================================
// WALSegment представляет один сегмент WAL.
type WALSegment struct {
ID uint32
File *os.File
Writer *bufio.Writer
Path string
StartLSN uint64
EndLSN uint64
Size int64
mu sync.Mutex
}
// WALIndexEntry — запись индекса WAL.
type WALIndexEntry struct {
LSN uint64
SegmentID uint32
Offset int64
Length uint32
Checksum uint32
}
// WALIndexManager — индекс WAL для быстрого доступа.
type WALIndexManager struct {
index map[uint64]*WALIndexEntry
segments map[uint32]*WALSegment
mu sync.RWMutex
indexPath string
}
// SegmentedWALManager — сегментированный WAL.
type SegmentedWALManager struct {
segmentsDir string
segments map[uint32]*WALSegment
currentSegment *WALSegment
currentSegmentID uint32
index *WALIndexManager
mu sync.RWMutex
writeChan chan *WALRecord
stopCh chan struct{}
wg sync.WaitGroup
batchSize int
logger LoggerInterface
fsyncEnabled bool
// ackMap — map LSN -> chan struct{} для синхронного подтверждения.
// Ключ — LSN записи. После записи в WAL закрываем канал.
ackMu sync.Mutex
ackMap map[uint64]chan struct{}
}
// NewSegmentedWALManager создаёт новый сегментированный WAL.
func NewSegmentedWALManager(segmentsDir string, fsyncEnabled bool, logger LoggerInterface) (*SegmentedWALManager, error) {
if err := os.MkdirAll(segmentsDir, 0755); err != nil {
return nil, fmt.Errorf("failed to create segments dir: %v", err)
}
wm := &SegmentedWALManager{
segmentsDir: segmentsDir,
segments: make(map[uint32]*WALSegment),
index: &WALIndexManager{
index: make(map[uint64]*WALIndexEntry),
segments: make(map[uint32]*WALSegment),
indexPath: filepath.Join(segmentsDir, WALIndexPrefix+"index.json"),
},
writeChan: make(chan *WALRecord, 10000),
stopCh: make(chan struct{}),
batchSize: 100,
logger: logger,
fsyncEnabled: fsyncEnabled,
ackMap: make(map[uint64]chan struct{}),
}
if err := wm.loadExistingSegments(); err != nil {
return nil, err
}
if err := wm.index.load(); err != nil {
if logger != nil {
logger.Warn(fmt.Sprintf("Failed to load WAL index: %v", err))
}
}
if wm.currentSegment == nil {
if err := wm.rotateSegmentInternal(); err != nil {
return nil, err
}
}
wm.wg.Add(1)
go wm.writerLoop()
return wm, nil
}
func (wm *SegmentedWALManager) loadExistingSegments() error {
files, err := filepath.Glob(filepath.Join(wm.segmentsDir, WALSegmentPrefix+"*"))
if err != nil {
return err
}
for _, filePath := range files {
var segmentID uint32
if _, err := fmt.Sscanf(filepath.Base(filePath), WALSegmentPrefix+"%d.log", &segmentID); err != nil {
continue
}
file, err := os.OpenFile(filePath, os.O_RDWR, 0644)
if err != nil {
continue
}
stat, _ := file.Stat()
segment := &WALSegment{
ID: segmentID,
File: file,
Writer: bufio.NewWriterSize(file, 64*1024),
Path: filePath,
Size: stat.Size(),
StartLSN: uint64(segmentID) * WALSegmentSize / 100,
}
wm.segments[segmentID] = segment
if segmentID > wm.currentSegmentID {
wm.currentSegmentID = segmentID
wm.currentSegment = segment
}
}
return nil
}
// rotateSegmentInternal — вызывается только при захваченном wm.mu.
func (wm *SegmentedWALManager) rotateSegmentInternal() error {
newSegmentID := wm.currentSegmentID + 1
segmentPath := filepath.Join(wm.segmentsDir, fmt.Sprintf(WALSegmentPrefix+"%d.log", newSegmentID))
file, err := os.OpenFile(segmentPath, os.O_CREATE|os.O_APPEND|os.O_RDWR, 0644)
if err != nil {
return fmt.Errorf("failed to create segment: %v", err)
}
newSegment := &WALSegment{
ID: newSegmentID,
File: file,
Writer: bufio.NewWriterSize(file, 64*1024),
Path: segmentPath,
StartLSN: wm.getCurrentLSNLocked(),
}
if wm.currentSegment != nil {
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.currentSegment.EndLSN = wm.getCurrentLSNLocked()
wm.currentSegment.File.Close()
FsyncDir(wm.segmentsDir)
}
wm.currentSegment = newSegment
wm.currentSegmentID = newSegmentID
wm.segments[newSegmentID] = newSegment
if wm.logger != nil {
wm.logger.Info(fmt.Sprintf("Created new WAL segment: %d", newSegmentID))
}
return nil
}
func (wm *SegmentedWALManager) getCurrentLSNLocked() uint64 {
if wm.currentSegment == nil {
return 1
}
return wm.currentSegment.StartLSN + uint64(wm.currentSegment.Size/100)
}
// Write — асинхронная запись.
func (wm *SegmentedWALManager) Write(record *WALRecord) error {
record.Timestamp = time.Now().UnixMilli()
wm.writeChan <- record
return nil
}
// WriteSync — синхронная запись с гарантией fsync.
//
// Используем ackMap по LSN, чтобы не было race condition
// с чужими подтверждениями.
func (wm *SegmentedWALManager) WriteSync(record *WALRecord) error {
record.Timestamp = time.Now().UnixMilli()
// Присваиваем LSN под защитой mu.
wm.mu.Lock()
lsn := wm.getCurrentLSNLocked() + 1
record.LSN = lsn
ackCh := make(chan struct{})
wm.ackMu.Lock()
wm.ackMap[lsn] = ackCh
wm.ackMu.Unlock()
wm.mu.Unlock()
wm.writeChan <- record
select {
case <-ackCh:
return nil
case <-time.After(10 * time.Second):
wm.ackMu.Lock()
delete(wm.ackMap, lsn)
wm.ackMu.Unlock()
return fmt.Errorf("WAL write sync timeout (LSN %d)", lsn)
}
}
func (wm *SegmentedWALManager) writerLoop() {
defer wm.wg.Done()
batch := make([]*WALRecord, 0, wm.batchSize)
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case record, ok := <-wm.writeChan:
if !ok {
if len(batch) > 0 {
wm.flushBatch(batch)
}
return
}
batch = append(batch, record)
if len(batch) >= wm.batchSize {
wm.flushBatch(batch)
batch = batch[:0]
}
case <-ticker.C:
if len(batch) > 0 {
wm.flushBatch(batch)
batch = batch[:0]
}
case <-wm.stopCh:
if len(batch) > 0 {
wm.flushBatch(batch)
}
return
}
}
}
// flushBatch — записывает batch записей и подтверждает через ackMap.
func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) {
wm.mu.Lock()
for _, record := range batch {
if wm.currentSegment.Size >= WALSegmentSize {
if err := wm.rotateSegmentInternal(); err != nil {
if wm.logger != nil {
wm.logger.Error(fmt.Sprintf("Failed to rotate segment: %v", err))
}
continue
}
}
data, err := json.Marshal(record)
if err != nil {
continue
}
lsnBytes := make([]byte, 8)
binary.BigEndian.PutUint64(lsnBytes, record.LSN)
crcData := append(lsnBytes, data...)
record.CRC = crc32(crcData)
header := make([]byte, 8)
binary.BigEndian.PutUint32(header[0:4], uint32(len(data)))
binary.BigEndian.PutUint32(header[4:8], record.CRC)
if _, err := wm.currentSegment.Writer.Write(header); err != nil {
continue
}
if _, err := wm.currentSegment.Writer.Write(data); err != nil {
continue
}
wm.index.addEntry(&WALIndexEntry{
LSN: record.LSN,
SegmentID: wm.currentSegment.ID,
Offset: wm.currentSegment.Size,
Length: uint32(len(data)),
Checksum: record.CRC,
})
wm.currentSegment.Size += int64(8 + len(data))
wm.currentSegment.EndLSN = record.LSN
}
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.index.save()
wm.mu.Unlock()
// Подтверждаем все записи из batch.
wm.ackMu.Lock()
for _, record := range batch {
if ch, ok := wm.ackMap[record.LSN]; ok {
close(ch)
delete(wm.ackMap, record.LSN)
}
}
wm.ackMu.Unlock()
}
// Sync — синхронизация.
func (wm *SegmentedWALManager) Sync() error {
wm.mu.Lock()
defer wm.mu.Unlock()
if wm.currentSegment == nil {
return nil
}
if err := wm.currentSegment.Writer.Flush(); err != nil {
return err
}
if wm.fsyncEnabled {
return RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
return nil
}
// ReadAll — читает все записи WAL.
func (wm *SegmentedWALManager) ReadAll() ([]*WALRecord, error) {
wm.mu.RLock()
segments := make([]*WALSegment, 0, len(wm.segments))
for _, seg := range wm.segments {
segments = append(segments, seg)
}
wm.mu.RUnlock()
sort.Slice(segments, func(i, j int) bool {
return segments[i].ID < segments[j].ID
})
records := make([]*WALRecord, 0)
for _, seg := range segments {
segRecords, err := wm.readSegmentRecords(seg)
if err != nil {
if wm.logger != nil {
wm.logger.Warn(fmt.Sprintf("Failed to read segment %d: %v", seg.ID, err))
}
continue
}
records = append(records, segRecords...)
}
return records, nil
}
// ReadSince — читает записи после указанного LSN.
func (wm *SegmentedWALManager) ReadSince(lsn uint64) ([]*WALRecord, error) {
allRecords, err := wm.ReadAll()
if err != nil {
return nil, err
}
result := make([]*WALRecord, 0)
for _, record := range allRecords {
if record.LSN > lsn {
result = append(result, record)
}
}
return result, nil
}
// GetCurrentLSN — текущий LSN.
func (wm *SegmentedWALManager) GetCurrentLSN() uint64 {
wm.mu.RLock()
defer wm.mu.RUnlock()
if wm.currentSegment == nil {
return 1
}
return wm.currentSegment.EndLSN
}
func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord, error) {
seg.mu.Lock()
defer seg.mu.Unlock()
if seg.File == nil {
return nil, nil
}
seg.Writer.Flush()
seg.File.Seek(0, 0)
records := make([]*WALRecord, 0)
reader := bufio.NewReader(seg.File)
headerBuf := make([]byte, 8)
for {
_, err := io.ReadFull(reader, headerBuf)
if err != nil {
break
}
recordLen := binary.BigEndian.Uint32(headerBuf[0:4])
expectedCRC := binary.BigEndian.Uint32(headerBuf[4:8])
if recordLen == 0 || recordLen > 100*1024*1024 {
break
}
recordData := make([]byte, recordLen)
_, err = io.ReadFull(reader, recordData)
if err != nil {
break
}
var record WALRecord
if err := json.Unmarshal(recordData, &record); err != nil {
continue
}
lsnBytes := make([]byte, 8)
binary.BigEndian.PutUint64(lsnBytes, record.LSN)
crcData := append(lsnBytes, recordData...)
calculatedCRC := crc32(crcData)
if calculatedCRC != expectedCRC && calculatedCRC != record.CRC {
continue
}
records = append(records, &record)
}
return records, nil
}
// Close — закрывает WAL.
func (wm *SegmentedWALManager) Close() error {
close(wm.stopCh)
wm.wg.Wait()
close(wm.writeChan)
wm.mu.Lock()
defer wm.mu.Unlock()
if wm.currentSegment != nil {
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.currentSegment.File.Close()
}
wm.index.save()
return nil
}
func (im *WALIndexManager) addEntry(entry *WALIndexEntry) {
im.mu.Lock()
defer im.mu.Unlock()
im.index[entry.LSN] = entry
}
func (im *WALIndexManager) save() error {
im.mu.RLock()
defer im.mu.RUnlock()
data, err := json.Marshal(im.index)
if err != nil {
return err
}
return os.WriteFile(im.indexPath, data, 0644)
}
func (im *WALIndexManager) load() error {
data, err := os.ReadFile(im.indexPath)
if err != nil {
if os.IsNotExist(err) {
return nil
}
return err
}
if len(data) == 0 {
return nil
}
return json.Unmarshal(data, &im.index)
}
// =============================================================================
// SAGA — ИНТЕРФЕЙСЫ ДЛЯ RAFT (без прямого импорта hashicorp/raft)
// =============================================================================
// SagaRaftAccessor — минимальный интерфейс к Raft, необходимый
// оркестратору SAGA для выборов лидера.
//
// Реализуется в internal/cluster (RaftCoordinator), чтобы избежать
// циклических импортов и позволить storage не зависеть от 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. Заполняется после успешного выполнения шага
// (в executeSagaInternal) и после успешной компенсации
// (в executeCompensation). Проверяется перед выполнением/компенсацией.
//
// Формат файла — JSON: { "execution_id": timestamp_ms, ... }.
// Это позволяет при рестарте быстро понять, какие шаги уже выполнены,
// и не дублировать их на удалённых узлах.
type SagaPersistentStorage struct {
baseDir string
mu sync.RWMutex
cache map[string]*SagaState
maxCache int
logger LoggerInterface
fsyncEnabled bool
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
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 сохраняет таблицу выполненных ключей на диск.
// Вызывается под executedMu.RLock или без удержания блокировки.
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, если ключ был добавлен впервые (т.е. операция
// действительно новая), и false, если ключ уже был (повтор).
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.
//
// Если 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
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)
}
}
// Если Raft доступен, регистрируем наблюдателя, чтобы оперативно
// получать уведомления о смене лидера.
if raft != nil {
raft.RegisterSagaLeaderObserver(nodeID, func(isLeader bool, leaderID string) {
e.updateLeadership(isLeader, leaderID)
})
}
// Запускаем фоновый опрос Raft (на случай, если наблюдатель не сработал).
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.
//
// ДОБАВЛЕНО: 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
}
// 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 // сколько SAGA было отклонено из-за отсутствия лидерства
}
// NewSagaOrchestratorWithConfig создаёт оркестратор.
//
// ДОБАВЛЕНО: второй аргумент raftAccessor может быть nil —
// тогда оркестратор работает в single-node режиме и всегда лидер.
//
// ДОБАВЛЕНО: если raftAccessor != nil, создаётся SagaLeaderElector,
// который опрашивает Raft и обновляет isLeader.
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)
// Начальное значение isLeader синхронизируем с выборщиком.
// Если Raft недоступен (nil), считаем себя лидером — single-node режим.
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 — является ли узел лидером.
//
// ДОБАВЛЕНО: теперь значение берётся из leaderElector, а не
// из жёстко выставленного флага.
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 уже выполнен.
// Делегируется в storage.
func (o *SagaOrchestrator) IsExecuted(executionID string) bool {
if o.storage == nil {
return false
}
return o.storage.IsExecuted(executionID)
}
// MarkExecuted помечает ExecutionID как выполненный.
// Возвращает true, если это первое выполнение.
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.
//
// Проверяем nil на Compensate, сохраняем состояние после
// каждой компенсации, обрабатываем compensation_failed.
//
// ДОБАВЛЕНО: проверяем лидерство — компенсацию выполняет только лидер,
// чтобы два узла не компенсировали одну и ту же 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.
//
// ДОБАВЛЕНО: только лидер выполняет 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.
//
// Компенсация идемпотентна и корректно обрабатывает ошибки.
//
// ДОБАВЛЕНО: идемпотентность через 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")
}
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 — выполнение компенсации.
//
// Идемпотентность: пропускаем уже скомпенсированные шаги.
// Сохраняем состояние после каждой компенсации.
// Обрабатываем 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")
}
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.
//
// ДОБАВЛЕНО: idempotent_skips, not_leader_rejects, leader_id.
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
}
// 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()
}
if globalBatchManager.aof != nil {
globalBatchManager.aof.Close()
}
return nil
}
// =============================================================================
// СОВМЕСТИМОСТЬ: псевдонимы для старого API
// =============================================================================
// GetTransactionStats — псевдоним для GetBatchStats (обратная совместимость).
func GetTransactionStats() map[string]interface{} {
return GetBatchStats()
}