Files
futriix/internal/storage/transactions.go
T

2118 lines
63 KiB
Go
Raw Normal View History

2026-10-03 20:31:37 +00:00
/*
* Copyright 2026 Safronov Grigorii
*
* Licensed under the CDDL, Version 1.0 (the "License");
* you may not use this file except in compliance with the License.
*
* You may obtain a copy of the License at
* https://opensource.org/licenses/CDDL-1.0
*/
// Файл: internal/storage/transactions.go
// Назначение: Batch-операции (замена полноценных ACID-транзакций) +
// Time-travel queries на основе WAL.
//
// АРХИТЕКТУРНЫЕ ИЗМЕНЕНИЯ (2026):
// - Удалены все мьютексы и MVCC (MVCCManager, VisibilityMap,
// ReadTimestampCache, DocumentVersion). Они давали огромную сложность
// без реальной пользы для in-memory СУБД.
// - Полноценные ACID-транзакции заменены на batch-операции.
// Batch — это группа операций, которые применяются атомарно:
// либо все, либо ни одна. Гарантия атомарности обеспечивается
// через WAL: сначала пишем весь batch в WAL, потом применяем.
// - SAGA-шаги используют batch для гарантии атомарности шага.
// - Time-travel реализован через WAL: для получения версии документа
// на момент времени проигрываем WAL до нужного LSN.
// - Удалён AsyncRecoveryManager (был источником race conditions).
// Восстановление теперь синхронное при старте.
//
// ИСПРАВЛЕНО (2026-10): добавлен реестр activeBatches в BatchManager,
// который используется eviction-механизмом (runtime_limits.go), чтобы
// не вытеснять документы, участвующие в незакоммиченной batch-операции.
// Ранее runtime_limits.go ссылался на удалённые globalTxManager и
// Transaction — теперь эти ссылки заменены на GetBatchManager() и *Batch.
package storage
import (
"bufio"
"encoding/binary"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"sort"
"sync"
"sync/atomic"
"time"
"futriis/internal/config"
)
// =============================================================================
// БАЗОВЫЕ ТИПЫ
// =============================================================================
type TransactionID uint64
// BatchOperation — одна операция в batch.
type BatchOperation struct {
Type string `json:"type"` // insert, update, delete, restore
Database string `json:"database"`
Collection string `json:"collection"`
DocumentID string `json:"document_id"`
Data map[string]interface{} `json:"data,omitempty"`
Version uint64 `json:"version,omitempty"`
}
// Batch представляет группу операций, применяемых атомарно.
type Batch struct {
ID TransactionID `json:"id"`
Operations []BatchOperation `json:"operations"`
Timestamp int64 `json:"timestamp"`
Database string `json:"database,omitempty"`
Collection string `json:"collection,omitempty"`
LSN uint64 `json:"lsn,omitempty"`
}
// WALRecord — запись в WAL.
type WALRecord struct {
CRC uint32 `json:"crc"`
Length uint32 `json:"length"`
Type byte `json:"type"` // 1 = batch
Data []byte `json:"data"`
Timestamp int64 `json:"timestamp"`
LSN uint64 `json:"lsn"`
}
// =============================================================================
// КОНСТАНТЫ
// =============================================================================
const (
WALSegmentSize = 64 * 1024 * 1024
WALSegmentPrefix = "wal_segment_"
WALIndexPrefix = "wal_index_"
VersionPruneInterval = 5 * time.Minute
SagaStateDir = "saga_states"
SagaStateFilePrefix = "saga_state_"
SagaStateFileSuffix = ".json"
SagaCheckpointInterval = 30 * time.Second
SagaMaxRetries = 5
SagaRetryBackoff = 100 * time.Millisecond
SagaOrphanTimeout = 5 * time.Minute
SagaCleanupInterval = 1 * time.Minute
)
// =============================================================================
// CRC32
// =============================================================================
var crc32Table = [256]uint32{
0x00000000, 0x77073096, 0xee0e612c, 0x990951ba, 0x076dc419, 0x706af48f,
0xe963a535, 0x9e6495a3, 0x0edb8832, 0x79dcb8a4, 0xe0d5e91e, 0x97d2d988,
0x09b64c2b, 0x7eb17cbd, 0xe7b82d07, 0x90bf1d91, 0x1db71064, 0x6ab020f2,
0xf3b97148, 0x84be41de, 0x1adad47d, 0x6ddde4eb, 0xf4d4b551, 0x83d385c7,
0x136c9856, 0x646ba8c0, 0xfd62f97a, 0x8a65c9ec, 0x14015c4f, 0x63066cd9,
0xfa0f3d63, 0x8d080df5, 0x3b6e20c8, 0x4c69105e, 0xd56041e4, 0xa2677172,
0x3c03e4d1, 0x4b04d447, 0xd20d85fd, 0xa50ab56b, 0x35b5a8fa, 0x42b2986c,
0xdbbbc9d6, 0xacbcf940, 0x32d86ce3, 0x45df5c75, 0xdcd60dcf, 0xabd13d59,
0x26d930ac, 0x51de003a, 0xc8d75180, 0xbfd06116, 0x21b4f4b5, 0x56b3c423,
0xcfba9599, 0xb8bda50f, 0x2802b89e, 0x5f058808, 0xc60cd9b2, 0xb10be924,
0x2f6f7c87, 0x58684c11, 0xc1611dab, 0xb6662d3d, 0x76dc4190, 0x01db7106,
0x98d220bc, 0xefd5102a, 0x71b18589, 0x06b6b51f, 0x9fbfe4a5, 0xe8b8d433,
0x7807c9a2, 0x0f00f934, 0x9609a88e, 0xe10e9818, 0x7f6a0dbb, 0x086d3d2d,
0x91646c97, 0xe6635c01, 0x6b6b51f4, 0x1c6c6162, 0x856530d8, 0xf262004e,
0x6c0695ed, 0x1b01a57b, 0x8208f4c1, 0xf50fc457, 0x65b0d9c6, 0x12b7e950,
0x8bbeb8ea, 0xfcb9887c, 0x62dd1ddf, 0x15da2d49, 0x8cd37cf3, 0xfbd44c65,
0x4db26158, 0x3ab551ce, 0xa3bc0074, 0xd4bb30e2, 0x4adfa541, 0x3dd895d7,
0xa4d1c46d, 0xd3d6f4fb, 0x4369e96a, 0x346ed9fc, 0xad678846, 0xda60b8d0,
0x44042d73, 0x33031de5, 0xaa0a4c5f, 0xdd0d7cc9, 0x5005713c, 0x270241aa,
0xbe0b1010, 0xc90c2086, 0x5768b525, 0x206f85b3, 0xb966d409, 0xce61e49f,
0x5edef90e, 0x29d9c998, 0xb0d09822, 0xc7d7a8b4, 0x59b33d17, 0x2eb40d81,
0xb7bd5c3b, 0xc0ba6cad, 0xedb88320, 0x9abfb3b6, 0x03b6e20c, 0x74b1d29a,
0xead54739, 0x9dd277af, 0x04db2615, 0x73dc1683, 0xe3630b12, 0x94643b84,
0x0d6d6a3e, 0x7a6a5aa8, 0xe40ecf0b, 0x9309ff9d, 0x0a00ae27, 0x7d079eb1,
0xf00f9344, 0x8708a3d2, 0x1e01f268, 0x6906c2fe, 0xf762575d, 0x806567cb,
0x196c3671, 0x6e6b06e7, 0xfed41b76, 0x89d32be0, 0x10da7a5a, 0x67dd4acc,
0xf9b9df6f, 0x8ebeeff9, 0x17b7be43, 0x60b08ed5, 0xd6d6a3e8, 0xa1d1937e,
0x38d8c2c4, 0x4fdff252, 0xd1bb67f1, 0xa6bc5767, 0x3fb506dd, 0x48b2364b,
0xd80d2bda, 0xaf0a1a4c, 0x36034af6, 0x41047a60, 0xdf60efc3, 0xa867df55,
0x316e8eef, 0x4669be79, 0xcb61b38c, 0xbc66831a, 0x256fd2a0, 0x5268e236,
0xcc0c7795, 0xbb0b4703, 0x220216b9, 0x5505262f, 0xc5ba3bbe, 0xb2bd0b28,
0x2bb45a92, 0x5cb36a04, 0xc2d7ffa7, 0xb5d0cf31, 0x2cd99e8b, 0x5bdeae1d,
0x9b64c2b0, 0xec63f226, 0x756aa39c, 0x026d930a, 0x9c0906a9, 0xeb0e363f,
0x72076785, 0x05005713, 0x95bf4a82, 0xe2b87a14, 0x7bb12bae, 0x0cb61b38,
0x92d28e9b, 0xe5d5be0d, 0x7cdcefb7, 0x0bdbdf21, 0x86d3d2d4, 0xf1d4e242,
0x68ddb3f8, 0x1fda836e, 0x81be16cd, 0xf6b9265b, 0x6fb077e1, 0x18b74777,
0x88085ae6, 0xff0f6a70, 0x66063bca, 0x11010b5c, 0x8f659eff, 0xf862ae69,
0x616bffd3, 0x166ccf45, 0xa00ae278, 0xd70dd2ee, 0x4e048354, 0x3903b3c2,
0xa7672661, 0xd06016f7, 0x4969474d, 0x3e6e77db, 0xaed16a4a, 0xd9d65adc,
0x40df0b66, 0x37d83bf8, 0xa9bcae53, 0xdebb9ec5, 0x47b2cf7f, 0x30b5ffe9,
0xbdbdf21c, 0xcabac28a, 0x53b39330, 0x24b4a3a6, 0xbad03605, 0xcdd70693,
0x54de5729, 0x23d967bf, 0xb3667a2e, 0xc4614ab8, 0x5d681b02, 0x2a6f2b94,
0xb40bbe37, 0xc30c8ea1, 0x5a05df1b, 0x2d02ef8d,
}
func crc32(data []byte) uint32 {
crc := uint32(0xFFFFFFFF)
for _, b := range data {
crc = (crc >> 8) ^ crc32Table[(crc^uint32(b))&0xFF]
}
return crc ^ 0xFFFFFFFF
}
// =============================================================================
// BATCH MANAGER - замена TransactionManager
// =============================================================================
// BatchManager управляет batch-операциями и WAL.
type BatchManager struct {
nextBatchID atomic.Uint64
wal *SegmentedWALManager
logger LoggerInterface
walPath string
stats *BatchStats
aof *AOFManager // append-only file для мутаций вне batch
// activeBatches — реестр активных (незакоммиченных) batch'ей.
// Ключ — Batch.ID, значение — *Batch.
// Используется eviction-механизмом (runtime_limits.go), чтобы не
// вытеснять документы, участвующие в незакоммиченной batch-операции.
activeBatches sync.Map
}
// BatchStats — статистика batch-операций.
type BatchStats struct {
TotalBatches atomic.Uint64
TotalOperations atomic.Uint64
TotalFailed atomic.Uint64
StartTime time.Time
}
// =============================================================================
// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ
// =============================================================================
var (
globalBatchManager *BatchManager
batchManagerOnce sync.Once
globalStorage *Storage
)
// InitBatchManager инициализирует batch-менеджер.
func InitBatchManager(walPath string) error {
return InitBatchManagerWithConfig(walPath, nil)
}
// InitBatchManagerWithConfig инициализирует batch-менеджер с конфигурацией.
func InitBatchManagerWithConfig(walPath string, config map[string]interface{}) error {
var err error
batchManagerOnce.Do(func() {
fsyncEnabled := true
if config != nil {
if v, ok := config["fsync_enabled"].(bool); ok {
fsyncEnabled = v
}
}
globalBatchManager = &BatchManager{
walPath: walPath,
stats: &BatchStats{StartTime: time.Now()},
}
globalBatchManager.nextBatchID.Store(1)
globalBatchManager.wal, err = NewSegmentedWALManager(
filepath.Dir(walPath), fsyncEnabled, nil)
if err != nil {
return
}
// Инициализируем AOF для мутаций вне batch.
globalBatchManager.aof, err = NewAOFManager(
filepath.Join(filepath.Dir(walPath), "aof"),
fsyncEnabled, nil)
if err != nil {
return
}
})
return err
}
// GetBatchManager возвращает глобальный batch-менеджер.
func GetBatchManager() *BatchManager { return globalBatchManager }
// SetBatchLogger устанавливает логгер.
func SetBatchLogger(logger LoggerInterface) {
if globalBatchManager != nil {
globalBatchManager.logger = logger
if globalBatchManager.wal != nil {
globalBatchManager.wal.logger = logger
}
if globalBatchManager.aof != nil {
globalBatchManager.aof.SetLogger(logger)
}
}
}
// SetGlobalStorage устанавливает глобальное хранилище.
func SetGlobalStorage(s *Storage) { globalStorage = s }
// GetGlobalStorage возвращает глобальное хранилище.
func GetGlobalStorage() *Storage { return globalStorage }
// =============================================================================
// РЕЕСТР АКТИВНЫХ BATCH'ЕЙ
// =============================================================================
// RegisterBatch регистрирует batch как активный.
// Вызывается в NewBatch.
func (bm *BatchManager) RegisterBatch(b *Batch) {
if bm == nil || b == nil {
return
}
bm.activeBatches.Store(b.ID, b)
}
// UnregisterBatch снимает регистрацию batch.
// Вызывается в Commit (успешном или неуспешном).
func (bm *BatchManager) UnregisterBatch(b *Batch) {
if bm == nil || b == nil {
return
}
bm.activeBatches.Delete(b.ID)
}
// IsDocInActiveBatch проверяет, участвует ли документ в активной
// (незакоммиченной) batch-операции.
//
// Используется в runtime_limits.go (eviction), чтобы не вытеснять
// документы, которые в данный момент участвуют в batch.
func (bm *BatchManager) IsDocInActiveBatch(docID, dbName, collName string) bool {
if bm == nil {
return false
}
found := false
bm.activeBatches.Range(func(key, value interface{}) bool {
b, ok := value.(*Batch)
if !ok || b == nil {
return true
}
for _, op := range b.Operations {
if op.DocumentID == docID &&
op.Database == dbName &&
op.Collection == collName {
found = true
return false
}
}
return true
})
return found
}
// GetActiveBatchCount возвращает количество активных batch'ей.
// Полезно для метрик и отладки.
func (bm *BatchManager) GetActiveBatchCount() int {
if bm == nil {
return 0
}
count := 0
bm.activeBatches.Range(func(_, _ interface{}) bool {
count++
return true
})
return count
}
// =============================================================================
// BATCH API
// =============================================================================
// NewBatch создаёт новый batch.
func NewBatch(database, collection string) *Batch {
if globalBatchManager == nil {
_ = InitBatchManager("futriis.wal")
}
b := &Batch{
ID: TransactionID(globalBatchManager.nextBatchID.Add(1) - 1),
Operations: make([]BatchOperation, 0, 16),
Timestamp: time.Now().UnixMilli(),
Database: database,
Collection: collection,
}
// Регистрируем batch как активный, чтобы eviction не вытеснил
// документы, участвующие в нём.
globalBatchManager.RegisterBatch(b)
return b
}
// AddInsert добавляет операцию вставки.
func (b *Batch) AddInsert(doc *Document) {
b.Operations = append(b.Operations, BatchOperation{
Type: "insert",
Database: b.Database,
Collection: b.Collection,
DocumentID: doc.ID,
Data: doc.GetFields(),
Version: doc.Version,
})
}
// AddUpdate добавляет операцию обновления.
func (b *Batch) AddUpdate(docID string, updates map[string]interface{}) {
b.Operations = append(b.Operations, BatchOperation{
Type: "update",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
Data: updates,
})
}
// AddDelete добавляет операцию удаления.
func (b *Batch) AddDelete(docID string) {
b.Operations = append(b.Operations, BatchOperation{
Type: "delete",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
})
}
// AddRestore добавляет операцию восстановления.
func (b *Batch) AddRestore(docID string) {
b.Operations = append(b.Operations, BatchOperation{
Type: "restore",
Database: b.Database,
Collection: b.Collection,
DocumentID: docID,
})
}
// Commit применяет batch атомарно.
//
// Гарантия атомарности:
// 1. Сериализуем весь batch и пишем в WAL (fsync).
// 2. Только после успешной записи в WAL применяем операции.
// 3. Если применение падает на середине — мы можем восстановить
// состояние, проиграв WAL (операции идемпотентны или
// откатываются вручную).
//
// В любом случае (успех или ошибка) batch снимается с регистрации
// в activeBatches.
func (b *Batch) Commit() error {
if globalBatchManager == nil {
return fmt.Errorf("batch manager not initialized")
}
// Снимаем регистрацию в любом случае — batch короткоживущий.
defer globalBatchManager.UnregisterBatch(b)
if len(b.Operations) == 0 {
return nil
}
// 1. Пишем batch в WAL.
data, err := json.Marshal(b)
if err != nil {
return fmt.Errorf("failed to marshal batch: %v", err)
}
record := &WALRecord{Type: 1, Data: data}
if err := globalBatchManager.wal.WriteSync(record); err != nil {
globalBatchManager.stats.TotalFailed.Add(1)
return fmt.Errorf("failed to write batch to WAL: %v", err)
}
// 2. Применяем операции.
if err := b.apply(); err != nil {
globalBatchManager.stats.TotalFailed.Add(1)
return fmt.Errorf("batch apply failed: %v", err)
}
// 3. Также пишем в AOF для быстрого восстановления (опционально).
if globalBatchManager.aof != nil {
if err := globalBatchManager.aof.AppendBatch(b); err != nil {
if globalBatchManager.logger != nil {
globalBatchManager.logger.Warn(fmt.Sprintf("AOF append failed: %v", err))
}
}
}
globalBatchManager.stats.TotalBatches.Add(1)
globalBatchManager.stats.TotalOperations.Add(uint64(len(b.Operations)))
return nil
}
// apply применяет все операции batch.
//
// Если операция падает, мы пытаемся откатить уже применённые
// (best-effort). Полная атомарность гарантируется только на уровне
// WAL: при восстановлении batch либо весь применится, либо нет.
func (b *Batch) apply() error {
if globalStorage == nil {
return fmt.Errorf("storage not initialized")
}
// Собираем откат для уже применённых операций.
type appliedOp struct {
opType string
db string
coll string
docID string
oldDoc *Document
}
applied := make([]appliedOp, 0, len(b.Operations))
rollback := func() {
for i := len(applied) - 1; i >= 0; i-- {
ao := applied[i]
db, err := globalStorage.GetDatabase(ao.db)
if err != nil {
continue
}
coll, err := db.GetCollection(ao.coll)
if err != nil {
continue
}
switch ao.opType {
case "insert":
coll.PermanentDelete(ao.docID)
case "update":
if ao.oldDoc != nil {
coll.Update(ao.docID, ao.oldDoc.GetFields())
}
case "delete":
coll.RestoreDeleted(ao.docID)
}
}
}
for _, op := range b.Operations {
db, err := globalStorage.GetDatabase(op.Database)
if err != nil {
rollback()
return fmt.Errorf("database not found: %s", op.Database)
}
coll, err := db.GetCollection(op.Collection)
if err != nil {
rollback()
return fmt.Errorf("collection not found: %s", op.Collection)
}
switch op.Type {
case "insert":
doc := NewDocumentWithID(op.DocumentID)
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
if err := coll.Insert(doc); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "insert", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
case "update":
oldDoc, _ := coll.Find(op.DocumentID)
var oldCopy *Document
if oldDoc != nil {
oldCopy = oldDoc.Clone()
}
if err := coll.Update(op.DocumentID, op.Data); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "update", db: op.Database, coll: op.Collection,
docID: op.DocumentID, oldDoc: oldCopy,
})
case "delete":
if err := coll.Delete(op.DocumentID); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "delete", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
case "restore":
if err := coll.RestoreDeleted(op.DocumentID); err != nil {
rollback()
return err
}
applied = append(applied, appliedOp{
opType: "restore", db: op.Database, coll: op.Collection, docID: op.DocumentID,
})
}
}
return nil
}
// =============================================================================
// TIME-TRAVEL QUERIES (WAL-based)
// =============================================================================
// GetDocumentAtTimestamp возвращает документ на указанный момент времени,
// проигрывая WAL с начала до нужного LSN.
//
// Это O(N) по количеству записей WAL — для in-memory СУБД с небольшим WAL
// это приемлемо. Для больших WAL нужен индекс по времени.
func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) {
if globalBatchManager == nil {
return nil, fmt.Errorf("batch manager not initialized")
}
if globalBatchManager.wal == nil {
return nil, fmt.Errorf("WAL not initialized")
}
records, err := globalBatchManager.wal.ReadAll()
if err != nil {
return nil, fmt.Errorf("failed to read WAL: %v", err)
}
var current *Document
for _, rec := range records {
if rec.Timestamp > timestamp {
break
}
if rec.Type != 1 {
continue
}
var batch Batch
if err := json.Unmarshal(rec.Data, &batch); err != nil {
continue
}
for _, op := range batch.Operations {
if op.DocumentID != docID {
continue
}
switch op.Type {
case "insert":
doc := NewDocumentWithID(op.DocumentID)
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
doc.CreatedAt = batch.Timestamp
doc.UpdatedAt = batch.Timestamp
current = doc
case "update":
if current != nil {
current.Update(op.Data)
current.UpdatedAt = batch.Timestamp
}
case "delete":
if current != nil {
current.SoftDelete()
current.DeletedAt = batch.Timestamp
}
case "restore":
if current != nil {
current.Restore()
current.UpdatedAt = batch.Timestamp
}
}
}
}
if current == nil {
return nil, fmt.Errorf("document %s not found at timestamp %d", docID, timestamp)
}
return current, nil
}
// GetDocumentVersionsSince возвращает все версии документа,
// созданные после указанного времени.
func GetDocumentVersionsSince(docID string, since int64) ([]*Document, error) {
if globalBatchManager == nil || globalBatchManager.wal == nil {
return nil, fmt.Errorf("batch manager not initialized")
}
records, err := globalBatchManager.wal.ReadAll()
if err != nil {
return nil, err
}
versions := make([]*Document, 0)
for _, rec := range records {
if rec.Timestamp <= since || rec.Type != 1 {
continue
}
var batch Batch
if err := json.Unmarshal(rec.Data, &batch); err != nil {
continue
}
for _, op := range batch.Operations {
if op.DocumentID != docID {
continue
}
doc := NewDocumentWithID(op.DocumentID)
if op.Type == "insert" || op.Type == "update" {
for k, v := range op.Data {
doc.SetField(k, v)
}
doc.Version = op.Version
doc.UpdatedAt = batch.Timestamp
versions = append(versions, doc)
}
}
}
return versions, nil
}
// =============================================================================
// WAL MANAGER (Segmented)
// =============================================================================
// WALSegment представляет один сегмент WAL.
type WALSegment struct {
ID uint32
File *os.File
Writer *bufio.Writer
Path string
StartLSN uint64
EndLSN uint64
Size int64
mu sync.Mutex
}
// WALIndexEntry — запись индекса WAL.
type WALIndexEntry struct {
LSN uint64
SegmentID uint32
Offset int64
Length uint32
Checksum uint32
}
// WALIndexManager — индекс WAL для быстрого доступа.
type WALIndexManager struct {
index map[uint64]*WALIndexEntry
segments map[uint32]*WALSegment
mu sync.RWMutex
indexPath string
}
// SegmentedWALManager — сегментированный WAL.
type SegmentedWALManager struct {
segmentsDir string
segments map[uint32]*WALSegment
currentSegment *WALSegment
currentSegmentID uint32
index *WALIndexManager
mu sync.RWMutex
writeChan chan *WALRecord
stopCh chan struct{}
wg sync.WaitGroup
batchSize int
logger LoggerInterface
fsyncEnabled bool
// ackMap — map LSN -> chan struct{} для синхронного подтверждения.
// Ключ — LSN записи. После записи в WAL закрываем канал.
ackMu sync.Mutex
ackMap map[uint64]chan struct{}
}
// NewSegmentedWALManager создаёт новый сегментированный WAL.
func NewSegmentedWALManager(segmentsDir string, fsyncEnabled bool, logger LoggerInterface) (*SegmentedWALManager, error) {
if err := os.MkdirAll(segmentsDir, 0755); err != nil {
return nil, fmt.Errorf("failed to create segments dir: %v", err)
}
wm := &SegmentedWALManager{
segmentsDir: segmentsDir,
segments: make(map[uint32]*WALSegment),
index: &WALIndexManager{
index: make(map[uint64]*WALIndexEntry),
segments: make(map[uint32]*WALSegment),
indexPath: filepath.Join(segmentsDir, WALIndexPrefix+"index.json"),
},
writeChan: make(chan *WALRecord, 10000),
stopCh: make(chan struct{}),
batchSize: 100,
logger: logger,
fsyncEnabled: fsyncEnabled,
ackMap: make(map[uint64]chan struct{}),
}
if err := wm.loadExistingSegments(); err != nil {
return nil, err
}
if err := wm.index.load(); err != nil {
if logger != nil {
logger.Warn(fmt.Sprintf("Failed to load WAL index: %v", err))
}
}
if wm.currentSegment == nil {
if err := wm.rotateSegmentInternal(); err != nil {
return nil, err
}
}
wm.wg.Add(1)
go wm.writerLoop()
return wm, nil
}
func (wm *SegmentedWALManager) loadExistingSegments() error {
files, err := filepath.Glob(filepath.Join(wm.segmentsDir, WALSegmentPrefix+"*"))
if err != nil {
return err
}
for _, filePath := range files {
var segmentID uint32
if _, err := fmt.Sscanf(filepath.Base(filePath), WALSegmentPrefix+"%d.log", &segmentID); err != nil {
continue
}
file, err := os.OpenFile(filePath, os.O_RDWR, 0644)
if err != nil {
continue
}
stat, _ := file.Stat()
segment := &WALSegment{
ID: segmentID,
File: file,
Writer: bufio.NewWriterSize(file, 64*1024),
Path: filePath,
Size: stat.Size(),
StartLSN: uint64(segmentID) * WALSegmentSize / 100,
}
wm.segments[segmentID] = segment
if segmentID > wm.currentSegmentID {
wm.currentSegmentID = segmentID
wm.currentSegment = segment
}
}
return nil
}
// rotateSegmentInternal — вызывается только при захваченном wm.mu.
func (wm *SegmentedWALManager) rotateSegmentInternal() error {
newSegmentID := wm.currentSegmentID + 1
segmentPath := filepath.Join(wm.segmentsDir, fmt.Sprintf(WALSegmentPrefix+"%d.log", newSegmentID))
file, err := os.OpenFile(segmentPath, os.O_CREATE|os.O_APPEND|os.O_RDWR, 0644)
if err != nil {
return fmt.Errorf("failed to create segment: %v", err)
}
newSegment := &WALSegment{
ID: newSegmentID,
File: file,
Writer: bufio.NewWriterSize(file, 64*1024),
Path: segmentPath,
StartLSN: wm.getCurrentLSNLocked(),
}
if wm.currentSegment != nil {
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.currentSegment.EndLSN = wm.getCurrentLSNLocked()
wm.currentSegment.File.Close()
FsyncDir(wm.segmentsDir)
}
wm.currentSegment = newSegment
wm.currentSegmentID = newSegmentID
wm.segments[newSegmentID] = newSegment
if wm.logger != nil {
wm.logger.Info(fmt.Sprintf("Created new WAL segment: %d", newSegmentID))
}
return nil
}
func (wm *SegmentedWALManager) getCurrentLSNLocked() uint64 {
if wm.currentSegment == nil {
return 1
}
return wm.currentSegment.StartLSN + uint64(wm.currentSegment.Size/100)
}
// Write — асинхронная запись.
func (wm *SegmentedWALManager) Write(record *WALRecord) error {
record.Timestamp = time.Now().UnixMilli()
wm.writeChan <- record
return nil
}
// WriteSync — синхронная запись с гарантией fsync.
//
// Используем ackMap по LSN, чтобы не было race condition
// с чужими подтверждениями.
func (wm *SegmentedWALManager) WriteSync(record *WALRecord) error {
record.Timestamp = time.Now().UnixMilli()
// Присваиваем LSN под защитой mu.
wm.mu.Lock()
lsn := wm.getCurrentLSNLocked() + 1
record.LSN = lsn
ackCh := make(chan struct{})
wm.ackMu.Lock()
wm.ackMap[lsn] = ackCh
wm.ackMu.Unlock()
wm.mu.Unlock()
wm.writeChan <- record
select {
case <-ackCh:
return nil
case <-time.After(10 * time.Second):
wm.ackMu.Lock()
delete(wm.ackMap, lsn)
wm.ackMu.Unlock()
return fmt.Errorf("WAL write sync timeout (LSN %d)", lsn)
}
}
func (wm *SegmentedWALManager) writerLoop() {
defer wm.wg.Done()
batch := make([]*WALRecord, 0, wm.batchSize)
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for {
select {
case record, ok := <-wm.writeChan:
if !ok {
if len(batch) > 0 {
wm.flushBatch(batch)
}
return
}
batch = append(batch, record)
if len(batch) >= wm.batchSize {
wm.flushBatch(batch)
batch = batch[:0]
}
case <-ticker.C:
if len(batch) > 0 {
wm.flushBatch(batch)
batch = batch[:0]
}
case <-wm.stopCh:
if len(batch) > 0 {
wm.flushBatch(batch)
}
return
}
}
}
// flushBatch — записывает batch записей и подтверждает через ackMap.
func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) {
wm.mu.Lock()
for _, record := range batch {
if wm.currentSegment.Size >= WALSegmentSize {
if err := wm.rotateSegmentInternal(); err != nil {
if wm.logger != nil {
wm.logger.Error(fmt.Sprintf("Failed to rotate segment: %v", err))
}
continue
}
}
data, err := json.Marshal(record)
if err != nil {
continue
}
lsnBytes := make([]byte, 8)
binary.BigEndian.PutUint64(lsnBytes, record.LSN)
crcData := append(lsnBytes, data...)
record.CRC = crc32(crcData)
header := make([]byte, 8)
binary.BigEndian.PutUint32(header[0:4], uint32(len(data)))
binary.BigEndian.PutUint32(header[4:8], record.CRC)
if _, err := wm.currentSegment.Writer.Write(header); err != nil {
continue
}
if _, err := wm.currentSegment.Writer.Write(data); err != nil {
continue
}
wm.index.addEntry(&WALIndexEntry{
LSN: record.LSN,
SegmentID: wm.currentSegment.ID,
Offset: wm.currentSegment.Size,
Length: uint32(len(data)),
Checksum: record.CRC,
})
wm.currentSegment.Size += int64(8 + len(data))
wm.currentSegment.EndLSN = record.LSN
}
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.index.save()
wm.mu.Unlock()
// Подтверждаем все записи из batch.
wm.ackMu.Lock()
for _, record := range batch {
if ch, ok := wm.ackMap[record.LSN]; ok {
close(ch)
delete(wm.ackMap, record.LSN)
}
}
wm.ackMu.Unlock()
}
// Sync — синхронизация.
func (wm *SegmentedWALManager) Sync() error {
wm.mu.Lock()
defer wm.mu.Unlock()
if wm.currentSegment == nil {
return nil
}
if err := wm.currentSegment.Writer.Flush(); err != nil {
return err
}
if wm.fsyncEnabled {
return RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
return nil
}
// ReadAll — читает все записи WAL.
func (wm *SegmentedWALManager) ReadAll() ([]*WALRecord, error) {
wm.mu.RLock()
segments := make([]*WALSegment, 0, len(wm.segments))
for _, seg := range wm.segments {
segments = append(segments, seg)
}
wm.mu.RUnlock()
sort.Slice(segments, func(i, j int) bool {
return segments[i].ID < segments[j].ID
})
records := make([]*WALRecord, 0)
for _, seg := range segments {
segRecords, err := wm.readSegmentRecords(seg)
if err != nil {
if wm.logger != nil {
wm.logger.Warn(fmt.Sprintf("Failed to read segment %d: %v", seg.ID, err))
}
continue
}
records = append(records, segRecords...)
}
return records, nil
}
// ReadSince — читает записи после указанного LSN.
func (wm *SegmentedWALManager) ReadSince(lsn uint64) ([]*WALRecord, error) {
allRecords, err := wm.ReadAll()
if err != nil {
return nil, err
}
result := make([]*WALRecord, 0)
for _, record := range allRecords {
if record.LSN > lsn {
result = append(result, record)
}
}
return result, nil
}
// GetCurrentLSN — текущий LSN.
func (wm *SegmentedWALManager) GetCurrentLSN() uint64 {
wm.mu.RLock()
defer wm.mu.RUnlock()
if wm.currentSegment == nil {
return 1
}
return wm.currentSegment.EndLSN
}
func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord, error) {
seg.mu.Lock()
defer seg.mu.Unlock()
if seg.File == nil {
return nil, nil
}
seg.Writer.Flush()
seg.File.Seek(0, 0)
records := make([]*WALRecord, 0)
reader := bufio.NewReader(seg.File)
headerBuf := make([]byte, 8)
for {
_, err := io.ReadFull(reader, headerBuf)
if err != nil {
break
}
recordLen := binary.BigEndian.Uint32(headerBuf[0:4])
expectedCRC := binary.BigEndian.Uint32(headerBuf[4:8])
if recordLen == 0 || recordLen > 100*1024*1024 {
break
}
recordData := make([]byte, recordLen)
_, err = io.ReadFull(reader, recordData)
if err != nil {
break
}
var record WALRecord
if err := json.Unmarshal(recordData, &record); err != nil {
continue
}
lsnBytes := make([]byte, 8)
binary.BigEndian.PutUint64(lsnBytes, record.LSN)
crcData := append(lsnBytes, recordData...)
calculatedCRC := crc32(crcData)
if calculatedCRC != expectedCRC && calculatedCRC != record.CRC {
continue
}
records = append(records, &record)
}
return records, nil
}
// Close — закрывает WAL.
func (wm *SegmentedWALManager) Close() error {
close(wm.stopCh)
wm.wg.Wait()
close(wm.writeChan)
wm.mu.Lock()
defer wm.mu.Unlock()
if wm.currentSegment != nil {
wm.currentSegment.Writer.Flush()
if wm.fsyncEnabled {
RealFsyncWithRetry(wm.currentSegment.File, 3, 100*time.Millisecond)
}
wm.currentSegment.File.Close()
}
wm.index.save()
return nil
}
func (im *WALIndexManager) addEntry(entry *WALIndexEntry) {
im.mu.Lock()
defer im.mu.Unlock()
im.index[entry.LSN] = entry
}
func (im *WALIndexManager) save() error {
im.mu.RLock()
defer im.mu.RUnlock()
data, err := json.Marshal(im.index)
if err != nil {
return err
}
return os.WriteFile(im.indexPath, data, 0644)
}
func (im *WALIndexManager) load() error {
data, err := os.ReadFile(im.indexPath)
if err != nil {
if os.IsNotExist(err) {
return nil
}
return err
}
if len(data) == 0 {
return nil
}
return json.Unmarshal(data, &im.index)
}
// =============================================================================
// SAGA PERSISTENT STORAGE
// =============================================================================
// SagaState — состояние SAGA.
type SagaState struct {
ID string `json:"id"`
Status string `json:"status"`
CurrentStep int `json:"current_step"`
Steps []SagaStepState `json:"steps"`
Data map[string]interface{} `json:"data"`
CreatedAt int64 `json:"created_at"`
UpdatedAt int64 `json:"updated_at"`
CompletedAt int64 `json:"completed_at,omitempty"`
CompensationExecuted bool `json:"compensation_executed"`
NodeID string `json:"node_id"`
CoordinatorID string `json:"coordinator_id"`
Version uint64 `json:"version"`
RetryCount int `json:"retry_count"`
LastError string `json:"last_error,omitempty"`
ExecutionID string `json:"execution_id"`
}
// SagaStepState — состояние шага SAGA.
type SagaStepState struct {
ID string `json:"id"`
Name string `json:"name"`
Status string `json:"status"`
Data map[string]interface{} `json:"data"`
StartedAt int64 `json:"started_at"`
CompletedAt int64 `json:"completed_at"`
ExecutionID string `json:"execution_id"`
RetryCount int `json:"retry_count"`
LastError string `json:"last_error,omitempty"`
CompensatedAt int64 `json:"compensated_at,omitempty"`
}
// SagaPersistentStorage — персистентное хранилище SAGA.
type SagaPersistentStorage struct {
baseDir string
mu sync.RWMutex
cache map[string]*SagaState
maxCache int
logger LoggerInterface
fsyncEnabled bool
fsyncMaxRetries int
fsyncRetryDelay time.Duration
}
// NewSagaPersistentStorageWithConfig создаёт хранилище SAGA.
func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInterface) (*SagaPersistentStorage, error) {
baseDir := SagaStateDir
maxCache := 10000
fsyncEnabled := true
fsyncMaxRetries := 3
fsyncRetryDelay := 100 * time.Millisecond
if cfg != nil {
baseDir = cfg.GetStateDir()
maxCache = cfg.GetMaxCacheSize()
fsyncEnabled = cfg.IsFsyncEnabled()
fsyncMaxRetries = cfg.GetFsyncMaxRetries()
fsyncRetryDelay = cfg.GetFsyncRetryDelay()
}
fullPath := filepath.Join(baseDir)
if err := os.MkdirAll(fullPath, 0755); err != nil {
return nil, fmt.Errorf("failed to create saga state directory: %v", err)
}
return &SagaPersistentStorage{
baseDir: fullPath,
cache: make(map[string]*SagaState),
maxCache: maxCache,
logger: logger,
fsyncEnabled: fsyncEnabled,
fsyncMaxRetries: fsyncMaxRetries,
fsyncRetryDelay: fsyncRetryDelay,
}, nil
}
func (sps *SagaPersistentStorage) getStatePath(sagaID string) string {
return filepath.Join(sps.baseDir, fmt.Sprintf("%s%s%s", SagaStateFilePrefix, sagaID, SagaStateFileSuffix))
}
// Save сохраняет состояние SAGA.
func (sps *SagaPersistentStorage) Save(state *SagaState) error {
sps.mu.Lock()
defer sps.mu.Unlock()
state.UpdatedAt = time.Now().UnixMilli()
state.Version++
if len(sps.cache) >= sps.maxCache {
var oldestKey string
var oldestTime int64 = time.Now().UnixMilli()
for k, v := range sps.cache {
if v.UpdatedAt < oldestTime {
oldestTime = v.UpdatedAt
oldestKey = k
}
}
if oldestKey != "" {
delete(sps.cache, oldestKey)
}
}
sps.cache[state.ID] = state
path := sps.getStatePath(state.ID)
data, err := json.MarshalIndent(state, "", " ")
if err != nil {
return fmt.Errorf("failed to marshal saga state: %v", err)
}
tmpPath := path + ".tmp"
if err := os.WriteFile(tmpPath, data, 0644); err != nil {
return fmt.Errorf("failed to write saga state: %v", err)
}
if sps.fsyncEnabled {
if f, err := os.OpenFile(tmpPath, os.O_RDWR, 0644); err == nil {
RealFsyncWithRetry(f, sps.fsyncMaxRetries, sps.fsyncRetryDelay)
f.Close()
}
}
if err := os.Rename(tmpPath, path); err != nil {
return fmt.Errorf("failed to rename saga state: %v", err)
}
if sps.fsyncEnabled {
FsyncDir(sps.baseDir)
}
return nil
}
// Load загружает состояние SAGA.
func (sps *SagaPersistentStorage) Load(sagaID string) (*SagaState, error) {
sps.mu.RLock()
if state, ok := sps.cache[sagaID]; ok {
sps.mu.RUnlock()
return state, nil
}
sps.mu.RUnlock()
path := sps.getStatePath(sagaID)
data, err := os.ReadFile(path)
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, fmt.Errorf("failed to read saga state: %v", err)
}
var state SagaState
if err := json.Unmarshal(data, &state); err != nil {
return nil, fmt.Errorf("failed to unmarshal saga state: %v", err)
}
sps.mu.Lock()
sps.cache[sagaID] = &state
sps.mu.Unlock()
return &state, nil
}
// Delete удаляет состояние SAGA.
func (sps *SagaPersistentStorage) Delete(sagaID string) error {
sps.mu.Lock()
delete(sps.cache, sagaID)
sps.mu.Unlock()
path := sps.getStatePath(sagaID)
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("failed to delete saga state: %v", err)
}
return nil
}
// ListAll возвращает все состояния SAGA.
func (sps *SagaPersistentStorage) ListAll() ([]*SagaState, error) {
pattern := filepath.Join(sps.baseDir, fmt.Sprintf("%s*%s", SagaStateFilePrefix, SagaStateFileSuffix))
files, err := filepath.Glob(pattern)
if err != nil {
return nil, err
}
states := make([]*SagaState, 0, len(files))
for _, file := range files {
data, err := os.ReadFile(file)
if err != nil {
continue
}
var state SagaState
if err := json.Unmarshal(data, &state); err != nil {
continue
}
states = append(states, &state)
}
return states, nil
}
// ListPending возвращает незавершённые SAGA.
func (sps *SagaPersistentStorage) ListPending() ([]*SagaState, error) {
all, err := sps.ListAll()
if err != nil {
return nil, err
}
pending := make([]*SagaState, 0)
for _, state := range all {
if state.Status == "pending" || state.Status == "running" || state.Status == "compensating" {
pending = append(pending, state)
}
}
return pending, nil
}
// =============================================================================
// SAGA ORCHESTRATOR (упрощённый, без мьютексов на горячем пути)
// =============================================================================
// SagaOrchestrator управляет SAGA.
type SagaOrchestrator struct {
storage *SagaPersistentStorage
mu sync.RWMutex
logger LoggerInterface
nodeID string
stopChan chan struct{}
wg sync.WaitGroup
isLeader atomic.Bool
leaderID string
config *config.SagaConfig
metrics *SagaMetrics
activeSagas sync.Map // map[string]*SagaTransaction
}
// SagaMetrics — метрики SAGA.
type SagaMetrics struct {
TotalStarted atomic.Uint64
TotalCompleted atomic.Uint64
TotalAborted atomic.Uint64
TotalFailed atomic.Uint64
TotalCompensated atomic.Uint64
ActiveCount atomic.Int64
RecoveryCount atomic.Uint64
}
// NewSagaOrchestratorWithConfig создаёт оркестратор.
func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) {
if cfg == nil {
cfg = &config.SagaConfig{
Enabled: true,
StateDir: "saga_states",
MaxRetries: 5,
RetryBackoffMs: 100,
SagaTimeoutSec: 300,
StuckCheckIntervalSec: 10,
RecoveryIntervalSec: 30,
CleanupPeriodHours: 24,
MaxCacheSize: 10000,
FsyncEnabled: true,
FsyncMaxRetries: 3,
FsyncRetryDelayMs: 100,
}
}
storage, err := NewSagaPersistentStorageWithConfig(cfg, logger)
if err != nil {
return nil, err
}
o := &SagaOrchestrator{
storage: storage,
logger: logger,
nodeID: nodeID,
stopChan: make(chan struct{}),
metrics: &SagaMetrics{},
config: cfg,
}
o.isLeader.Store(true)
o.wg.Add(1)
go o.recoveryLoop()
o.wg.Add(1)
go o.orphanCleanupLoop()
return o, nil
}
// IsLeader — является ли узел лидером.
func (o *SagaOrchestrator) IsLeader() bool { return o.isLeader.Load() }
// GetLeaderID — ID лидера.
func (o *SagaOrchestrator) GetLeaderID() string { return o.leaderID }
func (o *SagaOrchestrator) recoveryLoop() {
defer o.wg.Done()
ticker := time.NewTicker(o.config.GetRecoveryInterval())
defer ticker.Stop()
for {
select {
case <-ticker.C:
o.recoverPendingSagas()
case <-o.stopChan:
return
}
}
}
func (o *SagaOrchestrator) orphanCleanupLoop() {
defer o.wg.Done()
ticker := time.NewTicker(SagaCleanupInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
o.cleanupOrphanSteps()
case <-o.stopChan:
return
}
}
}
// cleanupOrphanSteps — очистка "осиротевших" SAGA.
func (o *SagaOrchestrator) cleanupOrphanSteps() {
pending, err := o.storage.ListPending()
if err != nil {
return
}
now := time.Now().UnixMilli()
for _, state := range pending {
if state.Status == "running" || state.Status == "compensating" {
elapsed := now - state.UpdatedAt
if elapsed > SagaOrphanTimeout.Milliseconds() {
if o.logger != nil {
o.logger.Warn(fmt.Sprintf("Found orphan saga %s, cleaning up", state.ID))
}
state.Status = "orphaned"
state.LastError = fmt.Sprintf("orphaned after %d ms", elapsed)
o.storage.Save(state)
o.compensateOrphanSaga(state)
}
}
}
}
// compensateOrphanSaga — компенсация "осиротевшей" SAGA.
//
// Проверяем nil на Compensate, сохраняем состояние после
// каждой компенсации, обрабатываем compensation_failed.
func (o *SagaOrchestrator) compensateOrphanSaga(state *SagaState) {
if state.CompensationExecuted {
return
}
if o.logger != nil {
o.logger.Info(fmt.Sprintf("Compensating orphan saga %s", state.ID))
}
saga := o.stateToTransaction(state)
if err := o.executeCompensation(saga, len(saga.Steps)-1); err != nil {
if o.logger != nil {
o.logger.Error(fmt.Sprintf("Failed to compensate orphan saga %s: %v", state.ID, err))
}
state.Status = "compensation_failed"
state.LastError = err.Error()
o.storage.Save(state)
return
}
state.CompensationExecuted = true
state.Status = "orphaned_compensated"
o.storage.Save(state)
}
// stateToTransaction конвертирует SagaState в SagaTransaction.
func (o *SagaOrchestrator) stateToTransaction(state *SagaState) *SagaTransaction {
saga := &SagaTransaction{
ID: state.ID,
Status: state.Status,
CurrentStep: state.CurrentStep,
Data: state.Data,
CreatedAt: state.CreatedAt,
UpdatedAt: state.UpdatedAt,
CompletedAt: state.CompletedAt,
compensationExecuted: state.CompensationExecuted,
executedOps: make(map[string]bool),
}
for _, stepState := range state.Steps {
step := &SagaStep{
ID: stepState.ID,
Name: stepState.Name,
Status: stepState.Status,
Data: stepState.Data,
StartedAt: stepState.StartedAt,
CompletedAt: stepState.CompletedAt,
ExecutionID: stepState.ExecutionID,
RetryCount: stepState.RetryCount,
LastError: stepState.LastError,
}
saga.Steps = append(saga.Steps, step)
if stepState.Status == "completed" {
saga.executedOps[stepState.ExecutionID] = true
}
}
return saga
}
// recoverPendingSagas — восстановление незавершённых SAGA.
func (o *SagaOrchestrator) recoverPendingSagas() {
pending, err := o.storage.ListPending()
if err != nil {
return
}
if len(pending) == 0 {
return
}
if o.logger != nil {
o.logger.Info(fmt.Sprintf("Recovering %d pending sagas", len(pending)))
}
for _, state := range pending {
if err := o.resumeSaga(state); err != nil {
if o.logger != nil {
o.logger.Error(fmt.Sprintf("Failed to resume saga %s: %v", state.ID, err))
}
} else {
o.metrics.RecoveryCount.Add(1)
}
}
}
// resumeSaga — возобновление SAGA.
func (o *SagaOrchestrator) resumeSaga(state *SagaState) error {
saga := o.stateToTransaction(state)
o.activeSagas.Store(state.ID, saga)
return o.executeSagaInternal(saga)
}
// BeginSaga — начать SAGA.
func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) {
existing, err := o.storage.Load(id)
if err != nil {
return nil, err
}
if existing != nil {
return nil, fmt.Errorf("saga %s already exists", id)
}
now := time.Now().UnixMilli()
state := &SagaState{
ID: id,
Status: "pending",
CurrentStep: 0,
Data: make(map[string]interface{}),
CreatedAt: now,
UpdatedAt: now,
NodeID: o.nodeID,
Version: 1,
ExecutionID: fmt.Sprintf("%s_%d", id, now),
}
if err := o.storage.Save(state); err != nil {
return nil, err
}
saga := &SagaTransaction{
ID: id,
Steps: make([]*SagaStep, 0),
CurrentStep: 0,
Status: "pending",
CreatedAt: now,
UpdatedAt: now,
Data: make(map[string]interface{}),
executedOps: make(map[string]bool),
}
o.activeSagas.Store(id, saga)
o.metrics.TotalStarted.Add(1)
o.metrics.ActiveCount.Add(1)
return saga, nil
}
// Execute — выполнить SAGA.
func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error {
return o.executeSagaInternal(saga)
}
// executeSagaInternal — внутренняя логика выполнения SAGA.
//
// Компенсация идемпотентна и корректно обрабатывает ошибки.
func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error {
saga.mu.Lock()
if saga.Status != "pending" && saga.Status != "running" {
saga.mu.Unlock()
return fmt.Errorf("saga %s is not in pending or running state", saga.ID)
}
saga.Status = "running"
saga.UpdatedAt = time.Now().UnixMilli()
saga.mu.Unlock()
if err := o.saveSagaState(saga); err != nil {
return err
}
for i := saga.CurrentStep; i < len(saga.Steps); i++ {
saga.mu.Lock()
step := saga.Steps[i]
saga.CurrentStep = i
if step.Status == "completed" {
saga.mu.Unlock()
continue
}
step.Status = "running"
step.StartedAt = time.Now().UnixMilli()
saga.mu.Unlock()
if saga.HasExecuted(step.ExecutionID) {
saga.mu.Lock()
step.Status = "completed"
step.CompletedAt = time.Now().UnixMilli()
saga.mu.Unlock()
continue
}
var err error
maxRetries := SagaMaxRetries
if o.config != nil {
maxRetries = o.config.GetMaxRetries()
}
for retry := 0; retry < maxRetries; retry++ {
saga.mu.Lock()
step.RetryCount = retry + 1
saga.mu.Unlock()
if err = step.Execute(); err == nil {
break
}
saga.mu.Lock()
step.LastError = err.Error()
saga.mu.Unlock()
backoff := time.Duration(o.config.RetryBackoffMs*(1<<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)
}
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 корректно.
func (o *SagaOrchestrator) executeCompensation(saga *SagaTransaction, failedStep int) error {
saga.mu.Lock()
if saga.compensationExecuted {
saga.mu.Unlock()
return nil
}
saga.mu.Unlock()
for j := failedStep; j >= 0; j-- {
saga.mu.RLock()
step := saga.Steps[j]
status := step.Status
saga.mu.RUnlock()
if status == "compensated" || status == "pending" {
continue
}
if status != "completed" && status != "failed" {
continue
}
if step.Compensate == nil {
// Нет функции компенсации — считаем шаг скомпенсированным.
saga.mu.Lock()
step.Status = "compensated"
step.CompletedAt = time.Now().UnixMilli()
saga.mu.Unlock()
o.saveSagaState(saga)
continue
}
var err error
maxRetries := SagaMaxRetries
if o.config != nil {
maxRetries = o.config.GetMaxRetries()
}
for retry := 0; retry < maxRetries; retry++ {
if err = step.Compensate(); err == nil {
break
}
backoff := time.Duration(o.config.RetryBackoffMs*(1<<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)
}
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(),
"is_leader": o.IsLeader(),
"node_id": o.nodeID,
}
}
// Stop — остановка оркестратора.
func (o *SagaOrchestrator) Stop() {
close(o.stopChan)
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-менеджер.
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,
}
}
// 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
}
// 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()
}