2119 lines
63 KiB
Go
2119 lines
63 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.
|
|
|
|
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()
|
|
}
|