2728 lines
86 KiB
Go
2728 lines
86 KiB
Go
/*
|
||
* 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()
|
||
}
|