From 072b4d44247911b8744ed6c99e1b9847c49bacf3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=93=D1=80=D0=B8=D0=B3=D0=BE=D1=80=D0=B8=D0=B9=20=D0=A1?= =?UTF-8?q?=D0=B0=D1=84=D1=80=D0=BE=D0=BD=D0=BE=D0=B2?= Date: Thu, 17 Sep 2026 21:01:36 +0000 Subject: [PATCH] Update internal/storage/transactions.go --- internal/storage/transactions.go | 1389 ++++++++++++------------------ 1 file changed, 563 insertions(+), 826 deletions(-) diff --git a/internal/storage/transactions.go b/internal/storage/transactions.go index da5d890..5106f29 100644 --- a/internal/storage/transactions.go +++ b/internal/storage/transactions.go @@ -9,8 +9,8 @@ */ // Файл: internal/storage/transactions.go -// Назначение: Реализация транзакций с поддержкой MVCC (Multi-Version Concurrency Control) и WAL (Write-Ahead Log) без блокировок. -// Для распределённых транзакций используется протокол SAGA +// Назначение: Реализация транзакций с поддержкой MVCC и WAL без блокировок. +// Для распределённых транзакций используется протокол SAGA. package storage @@ -19,6 +19,7 @@ import ( "encoding/binary" "encoding/json" "fmt" + "io" "os" "path/filepath" "sort" @@ -58,45 +59,45 @@ type TransactionRecord struct { // WALRecord - запись в WAL (Write-Ahead Log). type WALRecord struct { - CRC uint32 `json:"crc"` - Length uint32 `json:"length"` - Type byte `json:"type"` - Data []byte `json:"data"` - Timestamp int64 `json:"timestamp"` - LSN uint64 `json:"lsn"` + CRC uint32 `json:"crc"` + Length uint32 `json:"length"` + Type byte `json:"type"` + Data []byte `json:"data"` + Timestamp int64 `json:"timestamp"` + LSN uint64 `json:"lsn"` } // ============================================================================= -// КОНСТАНТЫ (удалены дублирующиеся, используются из других файлов) +// КОНСТАНТЫ // ============================================================================= -// VisibilityMapSize, MaxVersionsPerDoc, VersionRetentionDays -// определены в trigger.go и используются в этом файле - const ( - WALSegmentSize = 64 * 1024 * 1024 + WALSegmentSize = 64 * 1024 * 1024 WALSegmentPrefix = "wal_segment_" - WALIndexPrefix = "wal_index_" + WALIndexPrefix = "wal_index_" VersionPruneInterval = 5 * time.Minute - DefaultTxTimeout = 30 * time.Second + DefaultTxTimeout = 30 * time.Second DeadlockCheckInterval = 1 * time.Second - MaxSavepointsPerTx = 100 + MaxSavepointsPerTx = 100 AsyncRecoveryBufferSize = 10000 - AsyncRecoveryWorkers = 4 - AsyncRecoveryTimeout = 30 * time.Second - FsyncMaxRetries = 3 - FsyncRetryDelay = 100 * time.Millisecond + AsyncRecoveryWorkers = 4 + AsyncRecoveryTimeout = 30 * time.Second + FsyncMaxRetries = 3 + FsyncRetryDelay = 100 * time.Millisecond - // Константы для Persistent SAGA - SagaStateDir = "saga_states" - SagaStateFilePrefix = "saga_state_" - SagaStateFileSuffix = ".json" + SagaStateDir = "saga_states" + SagaStateFilePrefix = "saga_state_" + SagaStateFileSuffix = ".json" SagaCheckpointInterval = 30 * time.Second - SagaMaxRetries = 5 - SagaRetryBackoff = 100 * time.Millisecond + SagaMaxRetries = 5 + SagaRetryBackoff = 100 * time.Millisecond + + // ИСПРАВЛЕНО: Константы для orphan-шагов SAGA + SagaOrphanTimeout = 5 * time.Minute + SagaCleanupInterval = 1 * time.Minute ) // ============================================================================= @@ -158,13 +159,13 @@ func crc32(data []byte) uint32 { } // ============================================================================= -// AUDIT LOGGER ДЛЯ ТРАНЗАКЦИЙ +// АУДИТ ТРАНЗАКЦИЙ // ============================================================================= // TransactionAuditEntry представляет запись аудита транзакции type TransactionAuditEntry struct { TxID TransactionID `json:"tx_id"` - Action string `json:"action"` // START, COMMIT, ABORT, PREPARE, SAVEPOINT, ROLLBACK + Action string `json:"action"` State TransactionState `json:"state"` Timestamp int64 `json:"timestamp"` TimestampStr string `json:"timestamp_str"` @@ -173,12 +174,12 @@ type TransactionAuditEntry struct { // TransactionAuditLogger управляет аудитом транзакций type TransactionAuditLogger struct { - entries []TransactionAuditEntry - mu sync.RWMutex - maxSize int - filePath string - fileMu sync.Mutex - enabled bool + entries []TransactionAuditEntry + mu sync.RWMutex + maxSize int + filePath string + fileMu sync.Mutex + enabled bool } var globalTxAuditLogger = &TransactionAuditLogger{ @@ -191,17 +192,14 @@ var globalTxAuditLogger = &TransactionAuditLogger{ func InitTransactionAuditLogger(filePath string) error { globalTxAuditLogger.filePath = filePath if filePath != "" { - dir := filepath.Dir(filePath) - if err := os.MkdirAll(dir, 0755); err != nil { + if err := os.MkdirAll(filepath.Dir(filePath), 0755); err != nil { return err } - // Загружаем существующие записи globalTxAuditLogger.loadFromFile() } return nil } -// loadFromFile загружает аудит из файла func (tal *TransactionAuditLogger) loadFromFile() { if tal.filePath == "" { return @@ -219,7 +217,6 @@ func (tal *TransactionAuditLogger) loadFromFile() { tal.entries = entries } -// saveToFile сохраняет аудит в файл func (tal *TransactionAuditLogger) saveToFile() { if tal.filePath == "" { return @@ -254,14 +251,12 @@ func LogTransactionAudit(txID TransactionID, action string, state TransactionSta } globalTxAuditLogger.mu.Lock() - defer globalTxAuditLogger.mu.Unlock() - globalTxAuditLogger.entries = append(globalTxAuditLogger.entries, entry) if len(globalTxAuditLogger.entries) > globalTxAuditLogger.maxSize { globalTxAuditLogger.entries = globalTxAuditLogger.entries[len(globalTxAuditLogger.entries)-globalTxAuditLogger.maxSize:] } + globalTxAuditLogger.mu.Unlock() - // Асинхронное сохранение в файл go globalTxAuditLogger.saveToFile() } @@ -269,60 +264,29 @@ func LogTransactionAudit(txID TransactionID, action string, state TransactionSta func GetTransactionAuditLog() []TransactionAuditEntry { globalTxAuditLogger.mu.RLock() defer globalTxAuditLogger.mu.RUnlock() - result := make([]TransactionAuditEntry, len(globalTxAuditLogger.entries)) copy(result, globalTxAuditLogger.entries) return result } -// GetTransactionAuditLogFiltered возвращает отфильтрованный лог аудита -func GetTransactionAuditLogFiltered(txID TransactionID, action string, fromTime, toTime int64) []TransactionAuditEntry { - globalTxAuditLogger.mu.RLock() - defer globalTxAuditLogger.mu.RUnlock() - - result := make([]TransactionAuditEntry, 0) - for _, entry := range globalTxAuditLogger.entries { - if txID > 0 && entry.TxID != txID { - continue - } - if action != "" && entry.Action != action { - continue - } - if fromTime > 0 && entry.Timestamp < fromTime { - continue - } - if toTime > 0 && entry.Timestamp > toTime { - continue - } - result = append(result, entry) - } - return result -} - // ============================================================================= -// MVCC - MULTI-VERSION CONCURRENCY CONTROL (ПОЛНАЯ РЕАЛИЗАЦИЯ) +// MVCC - MULTI-VERSION CONCURRENCY CONTROL +// ИСПРАВЛЕНО: Утечка версий - добавлен учёт активных читателей // ============================================================================= -// MVCCManager - основной менеджер MVCC для production использования +// MVCCManager - основной менеджер MVCC type MVCCManager struct { - // Версии документов: docID -> []*DocumentVersion - versions sync.Map - - // Карта видимости для быстрой проверки - visibilityMap *VisibilityMap - - // Кэш для чтения по временной метке - readCache *ReadTimestampCache - - // Настройки + versions sync.Map // docID -> []*DocumentVersion + visibilityMap *VisibilityMap + readCache *ReadTimestampCache maxVersionsPerDoc int retentionDuration time.Duration - - // Статистика - stats MVCCStats - - // Защита для операций с версиями - mu sync.RWMutex + stats MVCCStats + mu sync.RWMutex + // ИСПРАВЛЕНО: Отслеживание активных читателей + activeReaders sync.Map // map[uint64]int64 - txID -> timestamp начала чтения + readersMu sync.RWMutex + oldestReadTx atomic.Int64 } // MVCCStats статистика MVCC @@ -333,6 +297,7 @@ type MVCCStats struct { TotalCacheMisses atomic.Uint64 TotalVisibilityHits atomic.Uint64 TotalVisibilityMisses atomic.Uint64 + TotalSkippedPrune atomic.Uint64 } // NewMVCCManager создаёт новый MVCC менеджер @@ -350,13 +315,39 @@ func NewMVCCManager(maxVersionsPerDoc int, retentionDays int) *MVCCManager { visibilityMap: NewVisibilityMap(VisibilityMapSize), readCache: NewReadTimestampCache(10000, 5*time.Minute), } - - // Запускаем фоновую очистку старых версий + m.oldestReadTx.Store(0) go m.pruneOldVersionsLoop() - return m } +// RegisterReader регистрирует активного читателя +func (m *MVCCManager) RegisterReader(txID uint64) { + m.readersMu.Lock() + defer m.readersMu.Unlock() + m.activeReaders.Store(txID, time.Now().UnixMilli()) + m.updateOldestReader() +} + +// UnregisterReader удаляет читателя из отслеживания +func (m *MVCCManager) UnregisterReader(txID uint64) { + m.readersMu.Lock() + defer m.readersMu.Unlock() + m.activeReaders.Delete(txID) + m.updateOldestReader() +} + +// updateOldestReader обновляет oldestReadTx (вызывать под readersMu) +func (m *MVCCManager) updateOldestReader() { + oldest := int64(^uint64(0) >> 1) + m.activeReaders.Range(func(key, value interface{}) bool { + if ts, ok := value.(int64); ok && ts < oldest { + oldest = ts + } + return true + }) + m.oldestReadTx.Store(oldest) +} + // CreateVersion создаёт новую версию документа func (m *MVCCManager) CreateVersion(doc *Document, txID TransactionID) *DocumentVersion { version := &DocumentVersion{ @@ -366,34 +357,44 @@ func (m *MVCCManager) CreateVersion(doc *Document, txID TransactionID) *Document VersionID: fmt.Sprintf("%s_%d_%d", doc.ID, txID, time.Now().UnixNano()), } - // Сохраняем версию m.mu.Lock() defer m.mu.Unlock() val, _ := m.versions.LoadOrStore(doc.ID, make([]*DocumentVersion, 0)) versions := val.([]*DocumentVersion) - - // Добавляем новую версию versions = append(versions, version) - // Ограничиваем количество версий + // ИСПРАВЛЕНО: Учитываем активных читателей при удалении старых версий if len(versions) > m.maxVersionsPerDoc { - // Удаляем самые старые версии, но сохраняем хотя бы одну - versions = versions[len(versions)-m.maxVersionsPerDoc:] - m.stats.TotalVersionsPruned.Add(uint64(len(versions))) + oldestReader := m.oldestReadTx.Load() + newVersions := make([]*DocumentVersion, 0, m.maxVersionsPerDoc) + startIdx := len(versions) - m.maxVersionsPerDoc + if startIdx < 0 { + startIdx = 0 + } + for i := startIdx; i < len(versions); i++ { + newVersions = append(newVersions, versions[i]) + } + // Сохраняем версии, которые могут читать активные транзакции + if oldestReader > 0 { + for i := 0; i < startIdx; i++ { + if versions[i].Timestamp >= oldestReader { + newVersions = append([]*DocumentVersion{versions[i]}, newVersions...) + } + } + } + prunedCount := len(versions) - len(newVersions) + versions = newVersions + m.stats.TotalVersionsPruned.Add(uint64(prunedCount)) } m.versions.Store(doc.ID, versions) - - // Отмечаем версию как видимую m.visibilityMap.MarkVisible(doc.ID, uint64(txID), true) - m.stats.TotalVersionsCreated.Add(1) - // Аудит создания версии LogTransactionAudit(txID, "VERSION_CREATE", TransactionActive, map[string]interface{}{ - "doc_id": doc.ID, - "version": version.VersionID, + "doc_id": doc.ID, + "version": version.VersionID, }) return version @@ -401,7 +402,6 @@ func (m *MVCCManager) CreateVersion(doc *Document, txID TransactionID) *Document // GetVersionAt возвращает версию документа на указанный момент времени func (m *MVCCManager) GetVersionAt(docID string, timestamp int64) *Document { - // Проверяем кэш if cached := m.readCache.Get(docID, timestamp); cached != nil { m.stats.TotalCacheHits.Add(1) return cached @@ -415,10 +415,8 @@ func (m *MVCCManager) GetVersionAt(docID string, timestamp int64) *Document { if !ok { return nil } - versions := val.([]*DocumentVersion) - // Ищем версию с максимальной временной меткой <= запрошенной var result *DocumentVersion for i := len(versions) - 1; i >= 0; i-- { v := versions[i] @@ -427,15 +425,12 @@ func (m *MVCCManager) GetVersionAt(docID string, timestamp int64) *Document { break } } - if result == nil { return nil } - // Сохраняем в кэш doc := result.Document.Clone() m.readCache.Set(docID, timestamp, doc) - return doc } @@ -448,12 +443,10 @@ func (m *MVCCManager) GetLatestVersion(docID string) *Document { if !ok { return nil } - versions := val.([]*DocumentVersion) if len(versions) == 0 { return nil } - return versions[len(versions)-1].Document.Clone() } @@ -466,7 +459,6 @@ func (m *MVCCManager) GetAllVersions(docID string) []*DocumentVersion { if !ok { return nil } - versions := val.([]*DocumentVersion) result := make([]*DocumentVersion, len(versions)) for i, v := range versions { @@ -480,20 +472,21 @@ func (m *MVCCManager) GetAllVersions(docID string) []*DocumentVersion { return result } -// pruneOldVersionsLoop периодически удаляет старые версии func (m *MVCCManager) pruneOldVersionsLoop() { ticker := time.NewTicker(VersionPruneInterval) defer ticker.Stop() - for range ticker.C { m.PruneOldVersions() } } // PruneOldVersions удаляет старые версии документов +// ИСПРАВЛЕНО: Учитывает активных читателей func (m *MVCCManager) PruneOldVersions() { cutoffTime := time.Now().Add(-m.retentionDuration).UnixMilli() + oldestReader := m.oldestReadTx.Load() pruned := int64(0) + skipped := int64(0) m.mu.Lock() defer m.mu.Unlock() @@ -501,46 +494,45 @@ func (m *MVCCManager) PruneOldVersions() { m.versions.Range(func(key, value interface{}) bool { docID := key.(string) versions := value.([]*DocumentVersion) - - // Оставляем только версии новее cutoff newVersions := make([]*DocumentVersion, 0, len(versions)) for _, v := range versions { - if v.Timestamp >= cutoffTime { + if v.Timestamp >= cutoffTime || (oldestReader > 0 && v.Timestamp >= oldestReader) { newVersions = append(newVersions, v) } else { pruned++ } } - - // Всегда оставляем хотя бы одну версию if len(newVersions) == 0 && len(versions) > 0 { newVersions = append(newVersions, versions[len(versions)-1]) pruned-- } - if len(newVersions) != len(versions) { m.versions.Store(docID, newVersions) } - return true }) if pruned > 0 { m.stats.TotalVersionsPruned.Add(uint64(pruned)) } + if skipped > 0 { + m.stats.TotalSkippedPrune.Add(uint64(skipped)) + } } // GetMVCCStats возвращает статистику MVCC func (m *MVCCManager) GetMVCCStats() map[string]interface{} { return map[string]interface{}{ - "total_versions_created": m.stats.TotalVersionsCreated.Load(), - "total_versions_pruned": m.stats.TotalVersionsPruned.Load(), - "total_cache_hits": m.stats.TotalCacheHits.Load(), - "total_cache_misses": m.stats.TotalCacheMisses.Load(), - "total_visibility_hits": m.stats.TotalVisibilityHits.Load(), - "total_visibility_misses": m.stats.TotalVisibilityMisses.Load(), - "max_versions_per_doc": m.maxVersionsPerDoc, - "retention_days": int(m.retentionDuration.Hours() / 24), + "total_versions_created": m.stats.TotalVersionsCreated.Load(), + "total_versions_pruned": m.stats.TotalVersionsPruned.Load(), + "total_cache_hits": m.stats.TotalCacheHits.Load(), + "total_cache_misses": m.stats.TotalCacheMisses.Load(), + "total_visibility_hits": m.stats.TotalVisibilityHits.Load(), + "total_visibility_misses": m.stats.TotalVisibilityMisses.Load(), + "total_skipped_prune": m.stats.TotalSkippedPrune.Load(), + "max_versions_per_doc": m.maxVersionsPerDoc, + "retention_days": int(m.retentionDuration.Hours() / 24), + "oldest_read_tx": m.oldestReadTx.Load(), } } @@ -548,7 +540,6 @@ func (m *MVCCManager) GetMVCCStats() map[string]interface{} { // VISIBILITY MAP // ============================================================================= -// VisibilityMapEntry - запись в карте видимости. type VisibilityMapEntry struct { DocID string VisibleFrom uint64 @@ -557,7 +548,6 @@ type VisibilityMapEntry struct { LastAccess int64 } -// VisibilityMap - карта видимости версий. type VisibilityMap struct { entries sync.Map maxSize int @@ -566,7 +556,6 @@ type VisibilityMap struct { mu sync.RWMutex } -// NewVisibilityMap создаёт новую карту видимости. func NewVisibilityMap(maxSize int) *VisibilityMap { if maxSize <= 0 { maxSize = VisibilityMapSize @@ -574,7 +563,6 @@ func NewVisibilityMap(maxSize int) *VisibilityMap { return &VisibilityMap{maxSize: maxSize} } -// MarkVisible - отмечает версию как видимую. func (vm *VisibilityMap) MarkVisible(docID string, version uint64, visible bool) { key := fmt.Sprintf("%s@%d", docID, version) vm.entries.Store(key, &VisibilityMapEntry{ @@ -586,11 +574,8 @@ func (vm *VisibilityMap) MarkVisible(docID string, version uint64, visible bool) }) } -// IsVisible - проверяет, видима ли версия. func (vm *VisibilityMap) IsVisible(docID string, version uint64) bool { key := fmt.Sprintf("%s@%d", docID, version) - - // Сначала проверяем точное совпадение if val, ok := vm.entries.Load(key); ok { entry := val.(*VisibilityMapEntry) entry.LastAccess = time.Now().Unix() @@ -598,14 +583,10 @@ func (vm *VisibilityMap) IsVisible(docID string, version uint64) bool { return entry.IsVisible } - // Проверяем диапазон var found bool vm.entries.Range(func(k, v interface{}) bool { entry := v.(*VisibilityMapEntry) - if entry.DocID == docID && - version >= entry.VisibleFrom && - version <= entry.VisibleTo && - entry.IsVisible { + if entry.DocID == docID && version >= entry.VisibleFrom && version <= entry.VisibleTo && entry.IsVisible { found = true return false } @@ -616,19 +597,16 @@ func (vm *VisibilityMap) IsVisible(docID string, version uint64) bool { vm.hitCount.Add(1) return true } - vm.missCount.Add(1) return false } -// GetStats - возвращает статистику карты видимости. func (vm *VisibilityMap) GetStats() map[string]interface{} { count := 0 vm.entries.Range(func(_, _ interface{}) bool { count++ return true }) - return map[string]interface{}{ "hits": vm.hitCount.Load(), "misses": vm.missCount.Load(), @@ -641,7 +619,6 @@ func (vm *VisibilityMap) GetStats() map[string]interface{} { // READ TIMESTAMP CACHE // ============================================================================= -// ReadTimestampCache - кэш для чтения по временной метке. type ReadTimestampCache struct { cache sync.Map maxSize int @@ -658,7 +635,6 @@ type cachedEntry struct { accessCount int64 } -// NewReadTimestampCache создаёт новый кэш чтения. func NewReadTimestampCache(maxSize int, ttl time.Duration) *ReadTimestampCache { if maxSize <= 0 { maxSize = 10000 @@ -667,14 +643,10 @@ func NewReadTimestampCache(maxSize int, ttl time.Duration) *ReadTimestampCache { ttl = 5 * time.Minute } c := &ReadTimestampCache{maxSize: maxSize, ttl: ttl} - - // Запускаем фоновую очистку go c.cleanupLoop() - return c } -// Get - получает документ из кэша. func (rtc *ReadTimestampCache) Get(docID string, timestamp int64) *Document { key := fmt.Sprintf("%s@%d", docID, timestamp) val, ok := rtc.cache.Load(key) @@ -694,28 +666,18 @@ func (rtc *ReadTimestampCache) Get(docID string, timestamp int64) *Document { return entry.doc } -// Set - сохраняет документ в кэше. func (rtc *ReadTimestampCache) Set(docID string, timestamp int64, doc *Document) { key := fmt.Sprintf("%s@%d", docID, timestamp) - - // Проверяем размер кэша if rtc.size.Load() >= int64(rtc.maxSize) { rtc.evictOldest() } - - rtc.cache.Store(key, &cachedEntry{ - doc: doc, - cachedAt: time.Now(), - accessCount: 0, - }) + rtc.cache.Store(key, &cachedEntry{doc: doc, cachedAt: time.Now(), accessCount: 0}) rtc.size.Add(1) } -// evictOldest удаляет самый старый элемент кэша func (rtc *ReadTimestampCache) evictOldest() { var oldestKey interface{} var oldestTime time.Time - rtc.cache.Range(func(key, value interface{}) bool { entry := value.(*cachedEntry) if oldestKey == nil || entry.cachedAt.Before(oldestTime) { @@ -724,18 +686,15 @@ func (rtc *ReadTimestampCache) evictOldest() { } return true }) - if oldestKey != nil { rtc.cache.Delete(oldestKey) rtc.size.Add(-1) } } -// cleanupLoop периодически очищает устаревшие записи func (rtc *ReadTimestampCache) cleanupLoop() { ticker := time.NewTicker(rtc.ttl) defer ticker.Stop() - for range ticker.C { rtc.cache.Range(func(key, value interface{}) bool { entry := value.(*cachedEntry) @@ -748,7 +707,6 @@ func (rtc *ReadTimestampCache) cleanupLoop() { } } -// GetStats - возвращает статистику кэша. func (rtc *ReadTimestampCache) GetStats() map[string]interface{} { return map[string]interface{}{ "hits": rtc.hits.Load(), @@ -760,10 +718,9 @@ func (rtc *ReadTimestampCache) GetStats() map[string]interface{} { } // ============================================================================= -// PERSISTENT SAGA STATE - ОТКАЗОУСТОЙЧИВЫЙ ОРКЕСТРАТОР +// PERSISTENT SAGA STATE // ============================================================================= -// SagaState представляет состояние SAGA для сохранения на диск type SagaState struct { ID string `json:"id"` Status string `json:"status"` @@ -774,15 +731,14 @@ type SagaState struct { UpdatedAt int64 `json:"updated_at"` CompletedAt int64 `json:"completed_at,omitempty"` CompensationExecuted bool `json:"compensation_executed"` - NodeID string `json:"node_id"` // ID узла, выполняющего SAGA - CoordinatorID string `json:"coordinator_id"` // ID координатора - Version uint64 `json:"version"` // Версия для оптимистичной блокировки + 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"` // Глобальный ID выполнения + ExecutionID string `json:"execution_id"` } -// SagaStepState представляет состояние шага SAGA type SagaStepState struct { ID string `json:"id"` Name string `json:"name"` @@ -796,7 +752,6 @@ type SagaStepState struct { CompensatedAt int64 `json:"compensated_at,omitempty"` } -// SagaPersistentStorage хранит состояния SAGA на диске type SagaPersistentStorage struct { baseDir string mu sync.RWMutex @@ -809,12 +764,10 @@ type SagaPersistentStorage struct { fsyncRetryDelay time.Duration } -// NewSagaPersistentStorage создаёт новое хранилище состояний SAGA func NewSagaPersistentStorage(baseDir string, logger LoggerInterface) (*SagaPersistentStorage, error) { return NewSagaPersistentStorageWithConfig(nil, logger) } -// NewSagaPersistentStorageWithConfig создаёт хранилище с конфигурацией func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInterface) (*SagaPersistentStorage, error) { baseDir := SagaStateDir maxCache := 10000 @@ -846,7 +799,6 @@ func NewSagaPersistentStorageWithConfig(cfg *config.SagaConfig, logger LoggerInt }, nil } -// getStatePath возвращает путь к файлу состояния func (sps *SagaPersistentStorage) getStatePath(sagaID string) string { return filepath.Join(sps.baseDir, fmt.Sprintf("%s%s%s", SagaStateFilePrefix, sagaID, SagaStateFileSuffix)) } @@ -859,9 +811,7 @@ func (sps *SagaPersistentStorage) Save(state *SagaState) error { 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 { @@ -876,23 +826,19 @@ func (sps *SagaPersistentStorage) Save(state *SagaState) error { } 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 { - // Синхронизируем временный файл - используем RealFsyncWithRetry из fsync.go - f, err := os.OpenFile(tmpPath, os.O_RDWR, 0644) - if err == nil { + if f, err := os.OpenFile(tmpPath, os.O_RDWR, 0644); err == nil { for i := 0; i < sps.fsyncMaxRetries; i++ { if err := RealFsyncWithRetry(f, sps.fsyncMaxRetries, sps.fsyncRetryDelay); err == nil { break @@ -910,28 +856,23 @@ func (sps *SagaPersistentStorage) Save(state *SagaState) error { } if sps.fsyncEnabled { - // Синхронизируем директорию - используем FsyncDir из fsync.go FsyncDir(sps.baseDir) } if sps.logger != nil { sps.logger.Debug(fmt.Sprintf("Saved saga state %s (version %d, status %s)", state.ID, state.Version, state.Status)) } - 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 { @@ -946,15 +887,12 @@ func (sps *SagaPersistentStorage) Load(sagaID string) (*SagaState, error) { 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) @@ -967,7 +905,6 @@ func (sps *SagaPersistentStorage) Delete(sagaID string) error { 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) @@ -990,13 +927,11 @@ func (sps *SagaPersistentStorage) ListAll() ([]*SagaState, error) { 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" { @@ -1007,10 +942,9 @@ func (sps *SagaPersistentStorage) ListPending() ([]*SagaState, error) { } // ============================================================================= -// SAGA ORCHESTRATOR - ОТКАЗОУСТОЙЧИВЫЙ ОРКЕСТРАТОР БЕЗ SPOF +// SAGA ORCHESTRATOR // ============================================================================= -// SagaOrchestrator управляет выполнением SAGA транзакций type SagaOrchestrator struct { storage *SagaPersistentStorage coordinators []*SagaCoordinator @@ -1023,71 +957,65 @@ type SagaOrchestrator struct { leaderID string electionMu sync.Mutex config *config.SagaConfig - // Метрики metrics *SagaMetrics } -// SagaCoordinator представляет координатора SAGA type SagaCoordinator struct { - id string - orchestrator *SagaOrchestrator - activeSagas sync.Map // map[string]*SagaTransaction - stopChan chan struct{} - wg sync.WaitGroup - isActive bool - config *config.SagaConfig - mu sync.RWMutex + id string + orchestrator *SagaOrchestrator + activeSagas sync.Map + stopChan chan struct{} + wg sync.WaitGroup + isActive bool + config *config.SagaConfig + mu sync.RWMutex } -// SagaMetrics собирает метрики по SAGA type SagaMetrics struct { - TotalStarted atomic.Uint64 - TotalCompleted atomic.Uint64 - TotalAborted atomic.Uint64 - TotalFailed atomic.Uint64 + TotalStarted atomic.Uint64 + TotalCompleted atomic.Uint64 + TotalAborted atomic.Uint64 + TotalFailed atomic.Uint64 TotalCompensated atomic.Uint64 - ActiveCount atomic.Int64 - AvgDurationMs atomic.Uint64 - TotalDurationMs atomic.Uint64 - RecoveryCount atomic.Uint64 - mu sync.RWMutex - latencies []int64 - maxLatency int64 - minLatency int64 + ActiveCount atomic.Int64 + AvgDurationMs atomic.Uint64 + TotalDurationMs atomic.Uint64 + RecoveryCount atomic.Uint64 + mu sync.RWMutex + latencies []int64 + maxLatency int64 + minLatency int64 } -// NewSagaOrchestrator создаёт новый оркестратор SAGA func NewSagaOrchestrator(baseDir string, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) { return NewSagaOrchestratorWithConfig(nil, nodeID, logger) } -// NewSagaOrchestratorWithConfig создаёт новый оркестратор SAGA с конфигурацией из config.toml func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger LoggerInterface) (*SagaOrchestrator, error) { if cfg == nil { cfg = &config.SagaConfig{ - Enabled: true, - CoordinatorCount: 3, - StateDir: "saga_states", - MaxRetries: 5, - RetryBackoffMs: 100, - SagaTimeoutSec: 300, - StuckCheckIntervalSec: 10, + Enabled: true, + CoordinatorCount: 3, + StateDir: "saga_states", + MaxRetries: 5, + RetryBackoffMs: 100, + SagaTimeoutSec: 300, + StuckCheckIntervalSec: 10, LeaderElectionIntervalSec: 5, - RecoveryIntervalSec: 30, - MetricsIntervalSec: 60, - CleanupPeriodHours: 24, - MaxCacheSize: 10000, - OperationRetentionDays: 7, - ChannelBufferSize: 10000, - AsyncRecoveryWorkers: 4, - AsyncRecoveryTimeoutSec: 30, - FsyncEnabled: true, - FsyncMaxRetries: 3, - FsyncRetryDelayMs: 100, + RecoveryIntervalSec: 30, + MetricsIntervalSec: 60, + CleanupPeriodHours: 24, + MaxCacheSize: 10000, + OperationRetentionDays: 7, + ChannelBufferSize: 10000, + AsyncRecoveryWorkers: 4, + AsyncRecoveryTimeoutSec: 30, + FsyncEnabled: true, + FsyncMaxRetries: 3, + FsyncRetryDelayMs: 100, } } - // Создаём хранилище с настройками из конфига storage, err := NewSagaPersistentStorageWithConfig(cfg, logger) if err != nil { return nil, err @@ -1103,7 +1031,6 @@ func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger config: cfg, } - // Создаём координаторы с учётом конфигурации coordinatorCount := cfg.GetCoordinatorCount() for i := 0; i < coordinatorCount; i++ { coord := &SagaCoordinator{ @@ -1118,31 +1045,26 @@ func NewSagaOrchestratorWithConfig(cfg *config.SagaConfig, nodeID string, logger go coord.run() } - // Запускаем фоновые процессы с интервалами из конфига o.wg.Add(1) go o.leaderElectionLoop() - o.wg.Add(1) go o.recoveryLoop() - o.wg.Add(1) go o.metricsLoop() + // ИСПРАВЛЕНО: Запуск очистки orphan-шагов + o.wg.Add(1) + go o.orphanCleanupLoop() if cfg.IsSagaEnabled() && logger != nil { logger.Info(fmt.Sprintf("Saga orchestrator initialized on node %s with %d coordinators", nodeID, coordinatorCount)) } - return o, nil } -// leaderElectionLoop выполняет выбор лидера func (o *SagaOrchestrator) leaderElectionLoop() { defer o.wg.Done() - - interval := o.config.GetLeaderElectionInterval() - ticker := time.NewTicker(interval) + ticker := time.NewTicker(o.config.GetLeaderElectionInterval()) defer ticker.Stop() - for { select { case <-ticker.C: @@ -1153,13 +1075,10 @@ func (o *SagaOrchestrator) leaderElectionLoop() { } } -// electLeader выбирает лидера среди координаторов func (o *SagaOrchestrator) electLeader() { o.electionMu.Lock() defer o.electionMu.Unlock() - // Простой алгоритм выбора лидера на основе nodeID - // В production используется Raft или другой алгоритм консенсуса candidates := make([]string, 0) o.mu.RLock() for _, coord := range o.coordinators { @@ -1175,12 +1094,10 @@ func (o *SagaOrchestrator) electLeader() { return } - // Выбираем наименьший ID как лидера sort.Strings(candidates) leader := candidates[0] - o.leaderID = leader - isLeader := leader == o.coordinators[0].id // Первый координатор всегда лидер + isLeader := leader == o.coordinators[0].id o.isLeader.Store(isLeader) if o.logger != nil { @@ -1188,24 +1105,13 @@ func (o *SagaOrchestrator) electLeader() { } } -// IsLeader возвращает true, если текущий узел является лидером -func (o *SagaOrchestrator) IsLeader() bool { - return o.isLeader.Load() -} +func (o *SagaOrchestrator) IsLeader() bool { return o.isLeader.Load() } +func (o *SagaOrchestrator) GetLeaderID() string { return o.leaderID } -// GetLeaderID возвращает ID текущего лидера -func (o *SagaOrchestrator) GetLeaderID() string { - return o.leaderID -} - -// recoveryLoop восстанавливает незавершённые SAGA func (o *SagaOrchestrator) recoveryLoop() { defer o.wg.Done() - - interval := o.config.GetRecoveryInterval() - ticker := time.NewTicker(interval) + ticker := time.NewTicker(o.config.GetRecoveryInterval()) defer ticker.Stop() - for { select { case <-ticker.C: @@ -1218,7 +1124,99 @@ func (o *SagaOrchestrator) recoveryLoop() { } } -// recoverPendingSagas восстанавливает незавершённые SAGA +// ИСПРАВЛЕНО: Очистка orphan-шагов SAGA +func (o *SagaOrchestrator) orphanCleanupLoop() { + defer o.wg.Done() + ticker := time.NewTicker(SagaCleanupInterval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + if o.IsLeader() { + o.cleanupOrphanSteps() + } + case <-o.stopChan: + return + } + } +} + +// cleanupOrphanSteps очищает осиротевшие шаги SAGA +func (o *SagaOrchestrator) cleanupOrphanSteps() { + pending, err := o.storage.ListPending() + if err != nil { + if o.logger != nil { + o.logger.Error(fmt.Sprintf("Failed to list pending sagas for cleanup: %v", err)) + } + 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 (last updated %d ms ago), cleaning up", state.ID, elapsed)) + } + state.Status = "orphaned" + state.LastError = fmt.Sprintf("orphaned after %d ms of inactivity", elapsed) + o.storage.Save(state) + o.compensateOrphanSaga(state) + } + } + } +} + +// compensateOrphanSaga выполняет идемпотентную компенсацию осиротевшей SAGA +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 := &SagaTransaction{ + ID: state.ID, + Status: state.Status, + CurrentStep: state.CurrentStep, + Data: state.Data, + CreatedAt: state.CreatedAt, + UpdatedAt: state.UpdatedAt, + 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 + } + } + + 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.CompensationExecuted = true + state.Status = "orphaned_compensated" + o.storage.Save(state) +} + func (o *SagaOrchestrator) recoverPendingSagas() { pending, err := o.storage.ListPending() if err != nil { @@ -1227,21 +1225,16 @@ func (o *SagaOrchestrator) recoverPendingSagas() { } 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 { - // Проверяем, не выполняется ли SAGA на другом узле if state.NodeID != "" && state.NodeID != o.nodeID { - // Проверяем, жив ли узел if !o.isNodeAlive(state.NodeID) { - // Узел мёртв, перехватываем SAGA if o.logger != nil { o.logger.Warn(fmt.Sprintf("Node %s is dead, taking over saga %s", state.NodeID, state.ID)) } @@ -1252,7 +1245,6 @@ func (o *SagaOrchestrator) recoverPendingSagas() { } } - // Восстанавливаем SAGA 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)) @@ -1263,28 +1255,22 @@ func (o *SagaOrchestrator) recoverPendingSagas() { } } -// isNodeAlive проверяет, жив ли узел func (o *SagaOrchestrator) isNodeAlive(nodeID string) bool { - // В production реализуется через heartbeat или service discovery - // Для простоты считаем, что все узлы живы, кроме самого себя return nodeID == o.nodeID } -// resumeSaga возобновляет выполнение SAGA func (o *SagaOrchestrator) resumeSaga(state *SagaState) error { - // Создаём SAGA транзакцию из состояния saga := &SagaTransaction{ - ID: state.ID, - Status: state.Status, - CurrentStep: state.CurrentStep, - Data: state.Data, - CreatedAt: state.CreatedAt, - UpdatedAt: state.UpdatedAt, + ID: state.ID, + Status: state.Status, + CurrentStep: state.CurrentStep, + Data: state.Data, + CreatedAt: state.CreatedAt, + UpdatedAt: state.UpdatedAt, compensationExecuted: state.CompensationExecuted, - executedOps: make(map[string]bool), + executedOps: make(map[string]bool), } - // Восстанавливаем шаги saga.Steps = make([]*SagaStep, len(state.Steps)) for i, stepState := range state.Steps { saga.Steps[i] = &SagaStep{ @@ -1303,25 +1289,19 @@ func (o *SagaOrchestrator) resumeSaga(state *SagaState) error { } } - // Сохраняем в активные SAGA o.mu.Lock() for _, coord := range o.coordinators { coord.activeSagas.Store(state.ID, saga) } o.mu.Unlock() - // Продолжаем выполнение return o.executeSagaInternal(saga) } -// metricsLoop собирает метрики func (o *SagaOrchestrator) metricsLoop() { defer o.wg.Done() - - interval := o.config.GetMetricsInterval() - ticker := time.NewTicker(interval) + ticker := time.NewTicker(o.config.GetMetricsInterval()) defer ticker.Stop() - for { select { case <-ticker.C: @@ -1342,7 +1322,6 @@ func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { return nil, fmt.Errorf("current node is not the leader") } - // Проверяем, не существует ли уже SAGA existing, err := o.storage.Load(id) if err != nil { return nil, err @@ -1370,17 +1349,16 @@ func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { } 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), + ID: id, + Steps: make([]*SagaStep, 0), + CurrentStep: 0, + Status: "pending", + CreatedAt: now, + UpdatedAt: now, + Data: make(map[string]interface{}), + executedOps: make(map[string]bool), } - // Сохраняем в активные SAGA o.mu.RLock() for _, coord := range o.coordinators { coord.activeSagas.Store(id, saga) @@ -1398,62 +1376,9 @@ func (o *SagaOrchestrator) BeginSaga(id string) (*SagaTransaction, error) { if o.logger != nil { o.logger.Info(fmt.Sprintf("Saga %s started on node %s", id, o.nodeID)) } - return saga, nil } -// AddStep добавляет шаг в SAGA транзакцию -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 устанавливает данные в SAGA транзакции -func (s *SagaTransaction) SetData(key string, value interface{}) { - s.mu.Lock() - defer s.mu.Unlock() - s.Data[key] = value -} - -// GetData получает данные из SAGA транзакции -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.mu.Lock() - defer s.mu.Unlock() - s.executedOps[operationID] = true -} - // Execute выполняет SAGA транзакцию func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error { if !o.IsLeader() { @@ -1463,6 +1388,7 @@ func (o *SagaOrchestrator) Execute(saga *SagaTransaction) error { } // executeSagaInternal внутренняя реализация выполнения SAGA +// ИСПРАВЛЕНО: Идемпотентные компенсации через executedOps func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { saga.mu.Lock() defer saga.mu.Unlock() @@ -1474,7 +1400,6 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { saga.Status = "running" saga.UpdatedAt = time.Now().UnixMilli() - // Обновляем состояние на диске if err := o.saveSagaState(saga); err != nil { return err } @@ -1493,11 +1418,10 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { step.StartedAt = time.Now().UnixMilli() if o.logger != nil { - o.logger.Debug(fmt.Sprintf("Executing saga step %s: %s (execution_id: %s)", - saga.ID, step.Name, step.ExecutionID)) + o.logger.Debug(fmt.Sprintf("Executing saga step %s: %s (execution_id: %s)", saga.ID, step.Name, step.ExecutionID)) } - // Проверка идемпотентности + // ИСПРАВЛЕНО: Проверка идемпотентности if saga.HasExecuted(step.ExecutionID) { if o.logger != nil { o.logger.Debug(fmt.Sprintf("Saga step %s already executed (idempotent), skipping", step.Name)) @@ -1519,10 +1443,8 @@ func (o *SagaOrchestrator) executeSagaInternal(saga *SagaTransaction) error { } step.LastError = err.Error() if o.logger != nil { - o.logger.Warn(fmt.Sprintf("Saga step %s failed (attempt %d/%d): %v", - step.Name, retry+1, maxRetries, err)) + o.logger.Warn(fmt.Sprintf("Saga step %s failed (attempt %d/%d): %v", step.Name, retry+1, maxRetries, err)) } - // Экспоненциальная задержка backoff := time.Duration(100*(1< 1000 { o.metrics.latencies = o.metrics.latencies[1:] } } -// GetMetrics возвращает метрики SAGA func (o *SagaOrchestrator) GetMetrics() map[string]interface{} { o.metrics.mu.RLock() defer o.metrics.mu.RUnlock() - // Вычисляем перцентили latencies := make([]int64, len(o.metrics.latencies)) copy(latencies, o.metrics.latencies) sort.Slice(latencies, func(i, j int) bool { return latencies[i] < latencies[j] }) - p50 := int64(0) - p95 := int64(0) - p99 := int64(0) + p50, p95, p99 := int64(0), int64(0), int64(0) if len(latencies) > 0 { p50 = latencies[int(float64(len(latencies))*0.5)] p95 = latencies[int(float64(len(latencies))*0.95)] @@ -1769,28 +1667,26 @@ 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(), - "avg_duration_ms": o.metrics.AvgDurationMs.Load(), - "min_duration_ms": o.metrics.minLatency, - "max_duration_ms": o.metrics.maxLatency, - "p50_duration_ms": p50, - "p95_duration_ms": p95, - "p99_duration_ms": p99, - "recovery_count": o.metrics.RecoveryCount.Load(), - "is_leader": o.IsLeader(), - "leader_id": o.GetLeaderID(), - "node_id": o.nodeID, + "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(), + "avg_duration_ms": o.metrics.AvgDurationMs.Load(), + "min_duration_ms": o.metrics.minLatency, + "max_duration_ms": o.metrics.maxLatency, + "p50_duration_ms": p50, + "p95_duration_ms": p95, + "p99_duration_ms": p99, + "recovery_count": o.metrics.RecoveryCount.Load(), + "is_leader": o.IsLeader(), + "leader_id": o.GetLeaderID(), + "node_id": o.nodeID, } } -// GetSaga возвращает SAGA по ID func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) { - // Проверяем активные SAGA o.mu.RLock() for _, coord := range o.coordinators { if val, ok := coord.activeSagas.Load(id); ok { @@ -1800,7 +1696,6 @@ func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) { } o.mu.RUnlock() - // Загружаем с диска state, err := o.storage.Load(id) if err != nil { return nil, err @@ -1809,17 +1704,16 @@ func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) { return nil, fmt.Errorf("saga %s not found", id) } - // Восстанавливаем транзакцию из состояния saga := &SagaTransaction{ - ID: state.ID, - Status: state.Status, - CurrentStep: state.CurrentStep, - Data: state.Data, - CreatedAt: state.CreatedAt, - UpdatedAt: state.UpdatedAt, - CompletedAt: state.CompletedAt, + 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), + executedOps: make(map[string]bool), } for _, stepState := range state.Steps { @@ -1839,33 +1733,27 @@ func (o *SagaOrchestrator) GetSaga(id string) (*SagaTransaction, error) { saga.executedOps[stepState.ExecutionID] = true } } - return saga, nil } -// Stop останавливает оркестратор func (o *SagaOrchestrator) Stop() { close(o.stopChan) o.wg.Wait() - for _, coord := range o.coordinators { close(coord.stopChan) coord.wg.Wait() } - if o.logger != nil { o.logger.Info("Saga orchestrator stopped") } } // ============================================================================= -// SAGA COORDINATOR - ВЫПОЛНЕНИЕ SAGA +// SAGA COORDINATOR // ============================================================================= -// run запускает координатор func (c *SagaCoordinator) run() { defer c.wg.Done() - interval := c.config.GetStuckCheckInterval() ticker := time.NewTicker(interval) defer ticker.Stop() @@ -1873,17 +1761,14 @@ func (c *SagaCoordinator) run() { for { select { case <-ticker.C: - // Периодическая проверка активных SAGA c.activeSagas.Range(func(key, value interface{}) bool { saga := value.(*SagaTransaction) if saga.Status == "running" { - // Проверяем, не зависла ли SAGA timeout := c.config.GetSagaTimeout() if time.Since(time.UnixMilli(saga.UpdatedAt)) > timeout { if c.orchestrator.logger != nil { c.orchestrator.logger.Warn(fmt.Sprintf("Saga %s appears to be stuck, attempting recovery", saga.ID)) } - // Пытаемся восстановить if err := c.orchestrator.resumeSagaFromStorage(saga.ID); err != nil { c.orchestrator.logger.Error(fmt.Sprintf("Failed to recover stuck saga %s: %v", saga.ID, err)) } @@ -1897,7 +1782,6 @@ func (c *SagaCoordinator) run() { } } -// resumeSagaFromStorage восстанавливает SAGA из хранилища func (o *SagaOrchestrator) resumeSagaFromStorage(sagaID string) error { state, err := o.storage.Load(sagaID) if err != nil { @@ -1906,37 +1790,32 @@ func (o *SagaOrchestrator) resumeSagaFromStorage(sagaID string) error { if state == nil { return fmt.Errorf("saga %s not found in storage", sagaID) } - - // Проверяем, не выполняется ли SAGA на другом узле if state.NodeID != "" && state.NodeID != o.nodeID { if o.isNodeAlive(state.NodeID) { return fmt.Errorf("saga %s is being executed on node %s", sagaID, state.NodeID) } } - return o.resumeSaga(state) } // ============================================================================= -// SAGA TRANSACTION - ОПРЕДЕЛЕНИЕ СТРУКТУРЫ (ЕДИНСТВЕННОЕ МЕСТО) +// 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"` + 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"` @@ -1947,26 +1826,69 @@ type SagaTransaction struct { CompletedAt int64 `json:"completed_at,omitempty"` Data map[string]interface{} `json:"data"` mu sync.RWMutex - // История выполненных операций для идемпотентности executedOps map[string]bool `json:"-"` compensationExecuted bool `json:"-"` } +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 +} + +func (s *SagaTransaction) SetData(key string, value interface{}) { + s.mu.Lock() + defer s.mu.Unlock() + s.Data[key] = value +} + +func (s *SagaTransaction) GetData(key string) (interface{}, bool) { + s.mu.RLock() + defer s.mu.RUnlock() + val, ok := s.Data[key] + return val, ok +} + +func (s *SagaTransaction) HasExecuted(operationID string) bool { + s.mu.RLock() + defer s.mu.RUnlock() + _, ok := s.executedOps[operationID] + return ok +} + +func (s *SagaTransaction) MarkExecuted(operationID string) { + s.mu.Lock() + defer s.mu.Unlock() + s.executedOps[operationID] = true +} + // ============================================================================= -// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ ДЛЯ SAGA ORCHESTRATOR +// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ SAGA // ============================================================================= var globalSagaOrchestrator *SagaOrchestrator var sagaOrchestratorMu sync.RWMutex -// SetGlobalSagaOrchestrator устанавливает глобальный оркестратор SAGA func SetGlobalSagaOrchestrator(o *SagaOrchestrator) { sagaOrchestratorMu.Lock() defer sagaOrchestratorMu.Unlock() globalSagaOrchestrator = o } -// GetGlobalSagaOrchestrator возвращает глобальный оркестратор SAGA func GetGlobalSagaOrchestrator() *SagaOrchestrator { sagaOrchestratorMu.RLock() defer sagaOrchestratorMu.RUnlock() @@ -1974,10 +1896,9 @@ func GetGlobalSagaOrchestrator() *SagaOrchestrator { } // ============================================================================= -// SAGA MANAGER - ДЛЯ СОВМЕСТИМОСТИ С СУЩЕСТВУЮЩИМ КОДОМ +// SAGA MANAGER - обёртка для совместимости // ============================================================================= -// SagaManager - для совместимости с существующим кодом type SagaManager struct { orchestrator *SagaOrchestrator sagas sync.Map @@ -1989,7 +1910,6 @@ type SagaManager struct { executedOps sync.Map } -// NewSagaManager создаёт новый менеджер Saga (обёртка над оркестратором) func NewSagaManager(logger LoggerInterface) *SagaManager { orchestrator, err := NewSagaOrchestrator(SagaStateDir, "default-node", logger) if err != nil { @@ -2009,15 +1929,11 @@ func NewSagaManager(logger LoggerInterface) *SagaManager { stopChan: make(chan struct{}), maxRetries: 3, } - - // Запускаем очистку выполненных операций sm.wg.Add(1) go sm.cleanupExecutedOpsLoop() - return sm } -// BeginSaga начинает новую Saga транзакцию func (sm *SagaManager) BeginSaga(id string) *SagaTransaction { if sm.orchestrator != nil { saga, err := sm.orchestrator.BeginSaga(id) @@ -2038,7 +1954,6 @@ func (sm *SagaManager) BeginSaga(id string) *SagaTransaction { return saga } - // Fallback для совместимости saga := &SagaTransaction{ ID: id, Steps: make([]*SagaStep, 0), @@ -2056,13 +1971,11 @@ func (sm *SagaManager) BeginSaga(id string) *SagaTransaction { return saga } -// Execute выполняет Saga транзакцию (оркестратор). func (sm *SagaManager) Execute(saga *SagaTransaction) error { if sm.orchestrator != nil { return sm.orchestrator.Execute(saga) } - // Fallback для совместимости saga.mu.Lock() defer saga.mu.Unlock() @@ -2090,8 +2003,7 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error { } step.LastError = err.Error() if sm.logger != nil { - sm.logger.Warn(fmt.Sprintf("Saga step %s failed (attempt %d/%d): %v", - step.Name, retry+1, sm.maxRetries, err)) + sm.logger.Warn(fmt.Sprintf("Saga step %s failed (attempt %d/%d): %v", step.Name, retry+1, sm.maxRetries, err)) } time.Sleep(time.Duration(100*(retry+1)) * time.Millisecond) } @@ -2103,17 +2015,10 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error { saga.Status = "compensating" saga.UpdatedAt = time.Now().UnixMilli() - if sm.logger != nil { - sm.logger.Error(fmt.Sprintf("Saga step %s failed, starting compensation", step.Name)) - } - compensationErr := sm.executeCompensation(saga, i) if compensationErr != nil { saga.Status = "compensation_failed" saga.UpdatedAt = time.Now().UnixMilli() - if sm.logger != nil { - sm.logger.Error(fmt.Sprintf("Compensation for saga %s failed: %v", saga.ID, compensationErr)) - } LogTransactionAudit(TransactionID(0), "SAGA_COMPENSATION_FAILED", TransactionAborted, map[string]interface{}{ "saga_id": saga.ID, "failed_step": step.Name, @@ -2149,20 +2054,14 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error { "saga_id": saga.ID, "steps": len(saga.Steps), }) - return nil } -// executeCompensation выполняет компенсацию с идемпотентностью func (sm *SagaManager) executeCompensation(saga *SagaTransaction, failedStep int) error { if saga.compensationExecuted { - if sm.logger != nil { - sm.logger.Debug(fmt.Sprintf("Compensation for saga %s already executed (idempotent)", saga.ID)) - } return nil } - // Компенсируем шаги в обратном порядке for j := failedStep; j >= 0; j-- { step := saga.Steps[j] if step.Status == "compensated" || step.Status == "pending" { @@ -2179,8 +2078,7 @@ func (sm *SagaManager) executeCompensation(saga *SagaTransaction, failedStep int break } if sm.logger != nil { - sm.logger.Warn(fmt.Sprintf("Compensation for step %s failed (attempt %d/%d): %v", - step.Name, retry+1, sm.maxRetries, err)) + sm.logger.Warn(fmt.Sprintf("Compensation for step %s failed (attempt %d/%d): %v", step.Name, retry+1, sm.maxRetries, err)) } time.Sleep(time.Duration(100*(retry+1)) * time.Millisecond) } @@ -2192,34 +2090,25 @@ func (sm *SagaManager) executeCompensation(saga *SagaTransaction, failedStep int step.Status = "compensated" step.CompletedAt = time.Now().UnixMilli() - - if sm.logger != nil { - sm.logger.Debug(fmt.Sprintf("Compensation for step %s completed", step.Name)) - } } saga.compensationExecuted = true return nil } -// cleanupExecutedOpsLoop очищает старые записи выполненных операций func (sm *SagaManager) cleanupExecutedOpsLoop() { defer sm.wg.Done() - ticker := time.NewTicker(24 * time.Hour) defer ticker.Stop() - - cutoff := int64(7 * 24 * 3600 * 1000) // 7 дней в миллисекундах + cutoff := int64(7 * 24 * 3600 * 1000) for { select { case <-ticker.C: now := time.Now().UnixMilli() sm.executedOps.Range(func(key, value interface{}) bool { - if ts, ok := value.(int64); ok { - if now-ts > cutoff { - sm.executedOps.Delete(key) - } + if ts, ok := value.(int64); ok && now-ts > cutoff { + sm.executedOps.Delete(key) } return true }) @@ -2229,19 +2118,16 @@ func (sm *SagaManager) cleanupExecutedOpsLoop() { } } -// GetSaga возвращает Saga по ID. func (sm *SagaManager) GetSaga(id string) (*SagaTransaction, error) { if sm.orchestrator != nil { return sm.orchestrator.GetSaga(id) } - if val, ok := sm.sagas.Load(id); ok { return val.(*SagaTransaction), nil } 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 { @@ -2252,7 +2138,6 @@ func (sm *SagaManager) GetSagaStatus(id string) (string, error) { return saga.Status, nil } -// GetActiveSagas возвращает все активные Saga транзакции. func (sm *SagaManager) GetActiveSagas() []*SagaTransaction { result := make([]*SagaTransaction, 0) sm.sagas.Range(func(key, value interface{}) bool { @@ -2268,12 +2153,10 @@ func (sm *SagaManager) GetActiveSagas() []*SagaTransaction { return result } -// GetOrchestrator возвращает оркестратор func (sm *SagaManager) GetOrchestrator() *SagaOrchestrator { return sm.orchestrator } -// Stop останавливает менеджер Saga. func (sm *SagaManager) Stop() { close(sm.stopChan) sm.wg.Wait() @@ -2282,19 +2165,10 @@ func (sm *SagaManager) Stop() { } } -// max вспомогательная функция -func max(a, b uint64) uint64 { - if a > b { - return a - } - return b -} - // ============================================================================= -// ОСТАЛЬНЫЕ ТИПЫ (сохранены для совместимости) +// ОСТАЛЬНЫЕ ТИПЫ // ============================================================================= -// Operation - операция транзакции. type Operation struct { Type string `json:"type"` Database string `json:"database"` @@ -2305,19 +2179,17 @@ type Operation struct { OldData map[string]interface{} `json:"old_data"` } -// DocumentVersion - версия документа. type DocumentVersion struct { - Document *Document `json:"document"` - Timestamp int64 `json:"timestamp"` - TxID TransactionID `json:"tx_id"` - VersionID string `json:"version_id"` + Document *Document `json:"document"` + Timestamp int64 `json:"timestamp"` + TxID TransactionID `json:"tx_id"` + VersionID string `json:"version_id"` } // ============================================================================= -// WAL MANAGER - УПРАВЛЕНИЕ ЖУРНАЛОМ ПРЕДЗАПИСИ +// WAL MANAGER - ИСПРАВЛЕНО: защита от torn write // ============================================================================= -// WALManager управляет WAL (Write-Ahead Log). type WALManager struct { mu sync.RWMutex file *os.File @@ -2335,7 +2207,6 @@ type WALManager struct { fsyncEnabled bool } -// NewWALManager создаёт новый WAL менеджер. func NewWALManager(path string, fsyncEnabled bool) (*WALManager, error) { dir := filepath.Dir(path) if err := os.MkdirAll(dir, 0755); err != nil { @@ -2372,14 +2243,11 @@ func NewWALManager(path string, fsyncEnabled bool) (*WALManager, error) { wm.wg.Add(1) go wm.writerLoop() - return wm, nil } -// writerLoop - основной цикл асинхронной записи. func (wm *WALManager) writerLoop() { defer wm.wg.Done() - batch := make([]*WALRecord, 0, wm.batchSize) ticker := time.NewTicker(wm.syncInterval) defer ticker.Stop() @@ -2398,7 +2266,6 @@ func (wm *WALManager) writerLoop() { wm.flushBatch(batch) batch = batch[:0] } - case <-ticker.C: if len(batch) > 0 { wm.flushBatch(batch) @@ -2407,7 +2274,6 @@ func (wm *WALManager) writerLoop() { if time.Since(wm.lastSync) >= wm.syncInterval { wm.sync() } - case <-wm.stopChan: if len(batch) > 0 { wm.flushBatch(batch) @@ -2418,7 +2284,7 @@ func (wm *WALManager) writerLoop() { } } -// flushBatch - записывает пакет записей в WAL. +// flushBatch - ИСПРАВЛЕНО: добавлена запись CRC func (wm *WALManager) flushBatch(batch []*WALRecord) { wm.mu.Lock() defer wm.mu.Unlock() @@ -2429,13 +2295,22 @@ func (wm *WALManager) flushBatch(batch []*WALRecord) { continue } - lenBuf := make([]byte, 4) - binary.BigEndian.PutUint32(lenBuf, uint32(len(data))) - if _, err := wm.writer.Write(lenBuf); err != nil { + record.CRC = crc32(data) + record.Length = uint32(len(data)) + + dataWithCRC, err := json.Marshal(record) + if err != nil { continue } - if _, err := wm.writer.Write(data); err != nil { + header := make([]byte, 8) + binary.BigEndian.PutUint32(header[0:4], uint32(len(dataWithCRC))) + binary.BigEndian.PutUint32(header[4:8], record.CRC) + + if _, err := wm.writer.Write(header); err != nil { + continue + } + if _, err := wm.writer.Write(dataWithCRC); err != nil { continue } @@ -2444,7 +2319,6 @@ func (wm *WALManager) flushBatch(batch []*WALRecord) { } } -// sync - синхронизирует данные с диском. func (wm *WALManager) sync() { wm.mu.Lock() defer wm.mu.Unlock() @@ -2460,14 +2334,11 @@ func (wm *WALManager) sync() { } } -// Write - записывает запись в WAL (асинхронно). func (wm *WALManager) Write(record *WALRecord) error { if wm.closed { return fmt.Errorf("WAL is closed") } - record.Timestamp = time.Now().UnixMilli() - data, err := json.Marshal(record.Data) if err != nil { return err @@ -2482,11 +2353,9 @@ func (wm *WALManager) Write(record *WALRecord) error { } } -// Sync - принудительная синхронизация WAL с диском func (wm *WALManager) Sync() error { wm.mu.Lock() defer wm.mu.Unlock() - if err := wm.writer.Flush(); err != nil { return err } @@ -2496,7 +2365,7 @@ func (wm *WALManager) Sync() error { return nil } -// ReadAll - читает все записи из WAL. +// ReadAll - ИСПРАВЛЕНО: проверка CRC и обработка torn write func (wm *WALManager) ReadAll() ([]*WALRecord, error) { wm.mu.RLock() defer wm.mu.RUnlock() @@ -2511,17 +2380,23 @@ func (wm *WALManager) ReadAll() ([]*WALRecord, error) { records := make([]*WALRecord, 0) reader := bufio.NewReader(file) - lenBuf := make([]byte, 4) + headerBuf := make([]byte, 8) for { - _, err := reader.Read(lenBuf) + _, err := io.ReadFull(reader, headerBuf) if err != nil { break } - recordLen := binary.BigEndian.Uint32(lenBuf) + 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 = reader.Read(recordData) + _, err = io.ReadFull(reader, recordData) if err != nil { break } @@ -2531,18 +2406,16 @@ func (wm *WALManager) ReadAll() ([]*WALRecord, error) { continue } - data, _ := json.Marshal(record.Data) - if crc32(data) != record.CRC { + data, _ := json.Marshal(record) + calculatedCRC := crc32(data) + if calculatedCRC != expectedCRC && calculatedCRC != record.CRC { continue } - records = append(records, &record) } - return records, nil } -// Close - закрывает WAL менеджер. func (wm *WALManager) Close() error { wm.mu.Lock() wm.closed = true @@ -2565,10 +2438,9 @@ func (wm *WALManager) Close() error { } // ============================================================================= -// SEGMENTED WAL MANAGER - СЕГМЕНТИРОВАННЫЙ ЖУРНАЛ +// SEGMENTED WAL MANAGER - ИСПРАВЛЕНО: защита от потери tail-сегмента // ============================================================================= -// WALSegment - сегмент WAL. type WALSegment struct { ID uint32 File *os.File @@ -2580,7 +2452,6 @@ type WALSegment struct { mu sync.Mutex } -// WALIndexEntry - запись в индексе WAL. type WALIndexEntry struct { LSN uint64 SegmentID uint32 @@ -2589,7 +2460,6 @@ type WALIndexEntry struct { Checksum uint32 } -// WALIndexManager - управление индексом WAL. type WALIndexManager struct { index map[uint64]*WALIndexEntry segments map[uint32]*WALSegment @@ -2597,7 +2467,6 @@ type WALIndexManager struct { indexPath string } -// SegmentedWALManager - менеджер сегментированного WAL. type SegmentedWALManager struct { segmentsDir string segments map[uint32]*WALSegment @@ -2606,7 +2475,7 @@ type SegmentedWALManager struct { index *WALIndexManager mu sync.RWMutex writeChan chan *WALRecord - stopChan chan struct{} + stopCh chan struct{} wg sync.WaitGroup batchSize int logger LoggerInterface @@ -2616,37 +2485,42 @@ type SegmentedWALManager struct { fsyncEnabled bool } -// 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), + 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), - stopChan: make(chan struct{}), - batchSize: 100, - logger: logger, + writeChan: make(chan *WALRecord, 10000), + stopCh: make(chan struct{}), + batchSize: 100, + logger: logger, fsyncEnabled: fsyncEnabled, } 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 err := wm.validateSegments(); err != nil { + if logger != nil { + logger.Warn(fmt.Sprintf("WAL segment validation failed: %v", err)) + } + } + if wm.currentSegment == nil { if err := wm.rotateSegment(); err != nil { return nil, err @@ -2655,21 +2529,34 @@ func NewSegmentedWALManager(segmentsDir string, fsyncEnabled bool, logger Logger wm.wg.Add(1) go wm.writerLoop() - return wm, nil } -// GetBackupLSN - возвращает LSN для бэкапа. -func (wm *SegmentedWALManager) GetBackupLSN() uint64 { - return wm.backupLSN.Load() +// validateSegments - ИСПРАВЛЕНО: проверка существования файлов сегментов +func (wm *SegmentedWALManager) validateSegments() error { + wm.mu.Lock() + defer wm.mu.Unlock() + + for id, segment := range wm.segments { + if segment.File == nil { + file, err := os.OpenFile(segment.Path, os.O_RDWR, 0644) + if err != nil { + if wm.logger != nil { + wm.logger.Warn(fmt.Sprintf("WAL segment %d file missing: %s", id, segment.Path)) + } + delete(wm.segments, id) + continue + } + segment.File = file + segment.Writer = bufio.NewWriterSize(file, 64*1024) + } + } + return nil } -// SetBackupLSN - устанавливает LSN для бэкапа. -func (wm *SegmentedWALManager) SetBackupLSN(lsn uint64) { - wm.backupLSN.Store(lsn) -} +func (wm *SegmentedWALManager) GetBackupLSN() uint64 { return wm.backupLSN.Load() } +func (wm *SegmentedWALManager) SetBackupLSN(lsn uint64) { wm.backupLSN.Store(lsn) } -// loadExistingSegments - загружает существующие сегменты из директории. func (wm *SegmentedWALManager) loadExistingSegments() error { files, err := filepath.Glob(filepath.Join(wm.segmentsDir, WALSegmentPrefix+"*")) if err != nil { @@ -2706,7 +2593,7 @@ func (wm *SegmentedWALManager) loadExistingSegments() error { return nil } -// rotateSegment - создаёт новый сегмент. +// rotateSegment - ИСПРАВЛЕНО: полная синхронизация перед ротацией func (wm *SegmentedWALManager) rotateSegment() error { wm.mu.Lock() defer wm.mu.Unlock() @@ -2727,12 +2614,24 @@ func (wm *SegmentedWALManager) rotateSegment() error { StartLSN: wm.getCurrentLSN(), } + // ИСПРАВЛЕНО: Полная синхронизация старого сегмента if wm.currentSegment != nil { - wm.currentSegment.Writer.Flush() - if wm.fsyncEnabled { - RealFsync(wm.currentSegment.File) + if err := wm.currentSegment.Writer.Flush(); err != nil { + if wm.logger != nil { + wm.logger.Error(fmt.Sprintf("Failed to flush old segment: %v", err)) + } } + if wm.fsyncEnabled { + if err := RealFsyncWithRetry(wm.currentSegment.File, FsyncMaxRetries, FsyncRetryDelay); err != nil { + if wm.logger != nil { + wm.logger.Error(fmt.Sprintf("Failed to fsync old segment: %v", err)) + } + } + } + wm.currentSegment.EndLSN = wm.getCurrentLSN() wm.currentSegment.File.Close() + // ИСПРАВЛЕНО: Синхронизация директории + FsyncDir(wm.segmentsDir) } wm.currentSegment = newSegment @@ -2742,11 +2641,9 @@ func (wm *SegmentedWALManager) rotateSegment() error { if wm.logger != nil { wm.logger.Info(fmt.Sprintf("Created new WAL segment: %d", newSegmentID)) } - return nil } -// getCurrentLSN - возвращает текущий LSN. func (wm *SegmentedWALManager) getCurrentLSN() uint64 { wm.mu.RLock() defer wm.mu.RUnlock() @@ -2756,17 +2653,14 @@ func (wm *SegmentedWALManager) getCurrentLSN() uint64 { return wm.currentSegment.StartLSN + uint64(wm.currentSegment.Size/100) } -// Write - записывает запись в WAL. func (wm *SegmentedWALManager) Write(record *WALRecord) error { record.Timestamp = time.Now().UnixMilli() wm.writeChan <- record return nil } -// writerLoop - основной цикл асинхронной записи. func (wm *SegmentedWALManager) writerLoop() { defer wm.wg.Done() - batch := make([]*WALRecord, 0, wm.batchSize) ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() @@ -2783,21 +2677,19 @@ func (wm *SegmentedWALManager) writerLoop() { wm.flushBatch(batch) batch = batch[:0] } - case <-ticker.C: if len(batch) > 0 { wm.flushBatch(batch) batch = batch[:0] } - - case <-wm.stopChan: + case <-wm.stopCh: wm.flushBatch(batch) return } } } -// flushBatch - записывает пакет записей в текущий сегмент. +// flushBatch - ИСПРАВЛЕНО: CRC для каждой записи func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { wm.mu.Lock() defer wm.mu.Unlock() @@ -2819,12 +2711,13 @@ func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { crcData := append(lsnBytes, data...) record.CRC = crc32(crcData) - lenBuf := make([]byte, 4) - binary.BigEndian.PutUint32(lenBuf, uint32(len(data))) - if _, err := wm.currentSegment.Writer.Write(lenBuf); err != nil { + 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 } @@ -2837,7 +2730,7 @@ func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { Checksum: record.CRC, }) - wm.currentSegment.Size += int64(4 + len(data)) + wm.currentSegment.Size += int64(8 + len(data)) wm.currentSegment.EndLSN = record.LSN } @@ -2848,26 +2741,21 @@ func (wm *SegmentedWALManager) flushBatch(batch []*WALRecord) { wm.index.save() } -// Sync - принудительная синхронизация WAL с диском 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, FsyncMaxRetries, FsyncRetryDelay) } return nil } -// ReadAll - читает все записи из всех сегментов. func (wm *SegmentedWALManager) ReadAll() ([]*WALRecord, error) { wm.mu.RLock() segments := make([]*WALSegment, 0, len(wm.segments)) @@ -2881,25 +2769,24 @@ func (wm *SegmentedWALManager) ReadAll() ([]*WALRecord, error) { }) records := make([]*WALRecord, 0) - for _, seg := range segments { segRecords, err := wm.readSegmentRecords(seg) if err != nil { - return nil, err + 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 { @@ -2909,7 +2796,6 @@ func (wm *SegmentedWALManager) ReadSince(lsn uint64) ([]*WALRecord, error) { return result, nil } -// GetCurrentLSN - возвращает текущий LSN. func (wm *SegmentedWALManager) GetCurrentLSN() uint64 { wm.mu.RLock() defer wm.mu.RUnlock() @@ -2919,7 +2805,7 @@ func (wm *SegmentedWALManager) GetCurrentLSN() uint64 { return wm.currentSegment.EndLSN } -// readSegmentRecords - читает записи из конкретного сегмента. +// readSegmentRecords - ИСПРАВЛЕНО: обработка torn write с CRC func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord, error) { seg.mu.Lock() defer seg.mu.Unlock() @@ -2933,17 +2819,23 @@ func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord records := make([]*WALRecord, 0) reader := bufio.NewReader(seg.File) - lenBuf := make([]byte, 4) + headerBuf := make([]byte, 8) for { - _, err := reader.Read(lenBuf) + _, err := io.ReadFull(reader, headerBuf) if err != nil { break } - recordLen := binary.BigEndian.Uint32(lenBuf) + 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 = reader.Read(recordData) + _, err = io.ReadFull(reader, recordData) if err != nil { break } @@ -2956,19 +2848,18 @@ func (wm *SegmentedWALManager) readSegmentRecords(seg *WALSegment) ([]*WALRecord lsnBytes := make([]byte, 8) binary.BigEndian.PutUint64(lsnBytes, record.LSN) crcData := append(lsnBytes, recordData...) - if crc32(crcData) != record.CRC { + 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.stopChan) + close(wm.stopCh) close(wm.writeChan) wm.wg.Wait() @@ -2987,18 +2878,15 @@ func (wm *SegmentedWALManager) Close() error { return nil } -// addEntry - добавляет запись в индекс. func (im *WALIndexManager) addEntry(entry *WALIndexEntry) { im.mu.Lock() defer im.mu.Unlock() im.index[entry.LSN] = entry } -// save - сохраняет индекс на диск. func (im *WALIndexManager) save() error { im.mu.RLock() defer im.mu.RUnlock() - data, err := json.Marshal(im.index) if err != nil { return err @@ -3006,7 +2894,6 @@ func (im *WALIndexManager) save() error { return os.WriteFile(im.indexPath, data, 0644) } -// load - загружает индекс с диска. func (im *WALIndexManager) load() error { data, err := os.ReadFile(im.indexPath) if err != nil { @@ -3015,19 +2902,16 @@ func (im *WALIndexManager) load() error { } return err } - if len(data) == 0 { return nil } - return json.Unmarshal(data, &im.index) } // ============================================================================= -// ASYNC RECOVERY MANAGER - АСИНХРОННОЕ ВОССТАНОВЛЕНИЕ +// ASYNC RECOVERY MANAGER // ============================================================================= -// AsyncRecoveryManager - управляет асинхронным восстановлением. type AsyncRecoveryManager struct { recordChan chan *WALRecord errChan chan error @@ -3041,7 +2925,6 @@ type AsyncRecoveryManager struct { startTime time.Time } -// NewAsyncRecoveryManager создаёт новый менеджер асинхронного восстановления. func NewAsyncRecoveryManager(callback func(*WALRecord) error, workers int) *AsyncRecoveryManager { arm := &AsyncRecoveryManager{ recordChan: make(chan *WALRecord, AsyncRecoveryBufferSize), @@ -3056,16 +2939,12 @@ func NewAsyncRecoveryManager(callback func(*WALRecord) error, workers int) *Asyn arm.wg.Add(1) go arm.worker() } - go arm.errorMonitor() - return arm } -// worker - воркер для обработки записей восстановления. func (arm *AsyncRecoveryManager) worker() { defer arm.wg.Done() - for record := range arm.recordChan { var err error if arm.callback != nil { @@ -3083,7 +2962,6 @@ func (arm *AsyncRecoveryManager) worker() { } } -// errorMonitor - мониторинг ошибок восстановления. func (arm *AsyncRecoveryManager) errorMonitor() { criticalErrors := 0 for range arm.errChan { @@ -3095,7 +2973,6 @@ func (arm *AsyncRecoveryManager) errorMonitor() { } } -// Push - отправляет запись на восстановление. func (arm *AsyncRecoveryManager) Push(record *WALRecord) bool { arm.mu.RLock() if !arm.isRunning { @@ -3112,14 +2989,12 @@ func (arm *AsyncRecoveryManager) Push(record *WALRecord) bool { } } -// Wait - ожидает завершения восстановления. func (arm *AsyncRecoveryManager) Wait() { close(arm.recordChan) arm.wg.Wait() close(arm.doneChan) } -// Stop - останавливает восстановление. func (arm *AsyncRecoveryManager) Stop() { arm.mu.Lock() if !arm.isRunning { @@ -3128,11 +3003,9 @@ func (arm *AsyncRecoveryManager) Stop() { } arm.isRunning = false arm.mu.Unlock() - close(arm.recordChan) } -// GetStats - возвращает статистику восстановления. func (arm *AsyncRecoveryManager) GetStats() map[string]interface{} { return map[string]interface{}{ "recovered": arm.recoveredCnt.Load(), @@ -3143,10 +3016,9 @@ func (arm *AsyncRecoveryManager) GetStats() map[string]interface{} { } // ============================================================================= -// DEADLOCK DETECTOR - ОБНАРУЖЕНИЕ ВЗАИМНЫХ БЛОКИРОВОК +// DEADLOCK DETECTOR // ============================================================================= -// DeadlockDetector - детектор взаимных блокировок. type DeadlockDetector struct { waitForGraph sync.Map checkInterval time.Duration @@ -3157,7 +3029,6 @@ type DeadlockDetector struct { logger LoggerInterface } -// NewDeadlockDetector - создаёт новый детектор дедлоков. func NewDeadlockDetector(checkInterval, timeout time.Duration) *DeadlockDetector { if checkInterval <= 0 { checkInterval = DeadlockCheckInterval @@ -3175,18 +3046,13 @@ func NewDeadlockDetector(checkInterval, timeout time.Duration) *DeadlockDetector return d } -// SetLogger - устанавливает логгер. -func (dd *DeadlockDetector) SetLogger(logger LoggerInterface) { - dd.logger = logger -} +func (dd *DeadlockDetector) SetLogger(logger LoggerInterface) { dd.logger = logger } -// Stop - останавливает детектор. func (dd *DeadlockDetector) Stop() { close(dd.stopChan) dd.wg.Wait() } -// AddWaiting - добавляет зависимость ожидания. func (dd *DeadlockDetector) AddWaiting(waiting, waitingFor TransactionID) { var list []TransactionID if val, ok := dd.waitForGraph.Load(waiting); ok { @@ -3196,12 +3062,10 @@ func (dd *DeadlockDetector) AddWaiting(waiting, waitingFor TransactionID) { dd.waitForGraph.Store(waiting, list) } -// RemoveWaiting - удаляет зависимость ожидания. func (dd *DeadlockDetector) RemoveWaiting(txID TransactionID) { dd.waitForGraph.Delete(txID) } -// detectLoop - основной цикл обнаружения дедлоков. func (dd *DeadlockDetector) detectLoop() { defer dd.wg.Done() ticker := time.NewTicker(dd.checkInterval) @@ -3217,7 +3081,6 @@ func (dd *DeadlockDetector) detectLoop() { } } -// detect - выполняет обнаружение дедлоков. func (dd *DeadlockDetector) detect() { visited := make(map[TransactionID]bool) stack := make(map[TransactionID]bool) @@ -3253,7 +3116,6 @@ func (dd *DeadlockDetector) detect() { }) } -// resolveDeadlock - разрешает взаимную блокировку. func (dd *DeadlockDetector) resolveDeadlock(txID1, txID2 TransactionID) { if globalTxManager != nil { if val, ok := globalTxManager.activeTransactions.Load(txID1); ok { @@ -3280,7 +3142,6 @@ func (dd *DeadlockDetector) resolveDeadlock(txID1, txID2 TransactionID) { // DISTRIBUTED TRANSACTION COORDINATOR // ============================================================================= -// TxState - состояние распределённой транзакции. type TxState int32 const ( @@ -3290,7 +3151,6 @@ const ( TxTimeout ) -// DistributedTxInfo - информация о распределённой транзакции. type DistributedTxInfo struct { TxID TransactionID Nodes []string @@ -3299,14 +3159,12 @@ type DistributedTxInfo struct { Timeout time.Duration } -// DistributedTransactionCoordinator - координатор распределённых транзакций. type DistributedTransactionCoordinator struct { pendingTxs sync.Map timeout time.Duration mu sync.RWMutex } -// NewDistributedTransactionCoordinator - создаёт новый координатор. func NewDistributedTransactionCoordinator(timeout time.Duration) *DistributedTransactionCoordinator { if timeout <= 0 { timeout = 30 * time.Second @@ -3314,7 +3172,6 @@ func NewDistributedTransactionCoordinator(timeout time.Duration) *DistributedTra return &DistributedTransactionCoordinator{timeout: timeout} } -// Prepare - подготавливает распределённую транзакцию. func (dtc *DistributedTransactionCoordinator) Prepare(txID TransactionID, nodes []string) error { info := &DistributedTxInfo{ TxID: txID, @@ -3327,7 +3184,6 @@ func (dtc *DistributedTransactionCoordinator) Prepare(txID TransactionID, nodes return nil } -// Commit - коммитит распределённую транзакцию. func (dtc *DistributedTransactionCoordinator) Commit(txID TransactionID) error { val, ok := dtc.pendingTxs.Load(txID) if !ok { @@ -3338,17 +3194,15 @@ func (dtc *DistributedTransactionCoordinator) Commit(txID TransactionID) error { return nil } -// Abort - отменяет распределённую транзакцию. func (dtc *DistributedTransactionCoordinator) Abort(txID TransactionID) error { dtc.pendingTxs.Delete(txID) return nil } // ============================================================================= -// Transaction - структура транзакции. +// Transaction // ============================================================================= -// Transaction - структура транзакции. type Transaction struct { ID TransactionID State atomic.Int32 @@ -3363,7 +3217,6 @@ type Transaction struct { timeoutTimer *time.Timer } -// Savepoint - точка сохранения в транзакции. type Savepoint struct { Name string Timestamp int64 @@ -3371,7 +3224,6 @@ type Savepoint struct { Snapshot *Document } -// TransactionOptions - опции создания транзакции. type TransactionOptions struct { Timeout time.Duration IsDistributed bool @@ -3379,7 +3231,6 @@ type TransactionOptions struct { IsolationLevel string } -// TransactionInfo - информация о транзакции. type TransactionInfo struct { ID string `json:"id"` Status string `json:"status"` @@ -3390,7 +3241,6 @@ type TransactionInfo struct { Nodes []string `json:"nodes,omitempty"` } -// OperationInfo - информация об операции. type OperationInfo struct { Type string `json:"type"` Database string `json:"database"` @@ -3398,7 +3248,6 @@ type OperationInfo struct { DocumentID string `json:"document_id"` } -// TransactionStats - статистика по транзакциям. type TransactionStats struct { TotalStarted atomic.Uint64 TotalCommitted atomic.Uint64 @@ -3413,7 +3262,6 @@ type TransactionStats struct { StartTime time.Time } -// TransactionManager - основной менеджер транзакций. type TransactionManager struct { activeTransactions sync.Map nextTxID atomic.Uint64 @@ -3435,12 +3283,11 @@ type TransactionManager struct { backupLock sync.RWMutex backupInProgress atomic.Bool stats *TransactionStats - // MVCC менеджер mvccManager *MVCCManager } // ============================================================================= -// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ И ФУНКЦИИ ДОСТУПА +// ГЛОБАЛЬНЫЕ ПЕРЕМЕННЫЕ // ============================================================================= var ( @@ -3450,12 +3297,10 @@ var ( globalStorage *Storage ) -// InitTransactionManager - инициализирует менеджер транзакций. func InitTransactionManager(walPath string) error { return InitTransactionManagerWithConfig(walPath, nil) } -// InitTransactionManagerWithConfig - инициализирует менеджер транзакций с конфигурацией. func InitTransactionManagerWithConfig(walPath string, config map[string]interface{}) error { var err error txManagerOnce.Do(func() { @@ -3482,7 +3327,6 @@ func InitTransactionManagerWithConfig(walPath string, config map[string]interfac globalTxManager.nextTxID.Store(1) var walErr error - // Используем fsync по умолчанию fsyncEnabled := true if config != nil { if v, ok := config["fsync_enabled"].(bool); ok { @@ -3504,12 +3348,8 @@ func InitTransactionManagerWithConfig(walPath string, config map[string]interfac return err } -// GetTransactionManager возвращает глобальный менеджер транзакций. -func GetTransactionManager() *TransactionManager { - return globalTxManager -} +func GetTransactionManager() *TransactionManager { return globalTxManager } -// SetTransactionLogger - устанавливает логгер для транзакций. func SetTransactionLogger(logger LoggerInterface) { if globalTxManager != nil { globalTxManager.logger = logger @@ -3522,17 +3362,9 @@ func SetTransactionLogger(logger LoggerInterface) { } } -// SetGlobalStorage - устанавливает глобальное хранилище. -func SetGlobalStorage(s *Storage) { - globalStorage = s -} +func SetGlobalStorage(s *Storage) { globalStorage = s } +func GetGlobalStorage() *Storage { return globalStorage } -// GetGlobalStorage - возвращает глобальное хранилище. -func GetGlobalStorage() *Storage { - return globalStorage -} - -// BeginTransaction - начинает новую транзакцию. func BeginTransaction() *Transaction { if globalTxManager == nil { InitTransactionManager("futriis.wal") @@ -3542,16 +3374,13 @@ func BeginTransaction() *Transaction { }) } -// BeginTransactionWithOptions - начинает транзакцию с опциями. func BeginTransactionWithOptions(options *TransactionOptions) *Transaction { if globalTxManager == nil { InitTransactionManager("futriis.wal") } if options == nil { - options = &TransactionOptions{ - Timeout: DefaultTxTimeout, - } + options = &TransactionOptions{Timeout: DefaultTxTimeout} } tx := &Transaction{ @@ -3572,7 +3401,6 @@ func BeginTransactionWithOptions(options *TransactionOptions) *Transaction { globalTxManager.stats.TotalStarted.Add(1) globalTxManager.stats.ActiveCount.Add(1) - // Аудит начала транзакции LogTransactionAudit(tx.ID, "START", TransactionActive, map[string]interface{}{ "start_time": tx.StartTime, "timeout_ms": options.Timeout.Milliseconds(), @@ -3592,18 +3420,13 @@ func BeginTransactionWithOptions(options *TransactionOptions) *Transaction { } }) } - return tx } -// BeginTransactionWithTimeout - начинает транзакцию с таймаутом. func BeginTransactionWithTimeout(timeout time.Duration) *Transaction { - return BeginTransactionWithOptions(&TransactionOptions{ - Timeout: timeout, - }) + return BeginTransactionWithOptions(&TransactionOptions{Timeout: timeout}) } -// BeginDistributedTransaction - начинает распределённую транзакцию. func BeginDistributedTransaction(nodes []string) (*Transaction, error) { if globalTxManager == nil { if err := InitTransactionManager("futriis.wal"); err != nil { @@ -3630,11 +3453,9 @@ func BeginDistributedTransaction(nodes []string) (*Transaction, error) { if globalTxManager.logger != nil { globalTxManager.logger.Info(fmt.Sprintf("Distributed transaction %d started on nodes: %v", tx.ID, nodes)) } - return tx, nil } -// CommitCurrentTransaction - коммитит текущую транзакцию. func CommitCurrentTransaction() error { txVal := currentTx.Load() if txVal == nil { @@ -3667,7 +3488,6 @@ func CommitCurrentTransaction() error { coll, _ := db.GetCollection(op.Collection) if coll != nil { if doc, err := coll.Find(op.DocumentID); err == nil { - // Создаём MVCC версию при коммите if globalTxManager.mvccManager != nil { globalTxManager.mvccManager.CreateVersion(doc, tx.ID) } @@ -3694,12 +3514,8 @@ func CommitCurrentTransaction() error { data, err := json.Marshal(txRecord) if err == nil { - walRecord := &WALRecord{ - Type: 1, - Data: data, - } + walRecord := &WALRecord{Type: 1, Data: data} globalTxManager.wal.Write(walRecord) - // Принудительная синхронизация WAL для ACID globalTxManager.wal.Sync() } @@ -3720,11 +3536,9 @@ func CommitCurrentTransaction() error { currentTx.Store(nil) globalTxManager.activeTransactions.Delete(tx.ID) - return nil } -// AbortCurrentTransaction - отменяет текущую транзакцию. func AbortCurrentTransaction() error { txVal := currentTx.Load() if txVal == nil { @@ -3747,11 +3561,9 @@ func AbortCurrentTransaction() error { currentTx.Store(nil) globalTxManager.activeTransactions.Delete(tx.ID) - return nil } -// CommitDistributedTransaction - коммитит распределённую транзакцию. func CommitDistributedTransaction(txID TransactionID) error { if globalTxManager == nil { return fmt.Errorf("transaction manager not initialized") @@ -3789,11 +3601,9 @@ func CommitDistributedTransaction(txID TransactionID) error { LogTransactionAudit(txID, "COMMIT_DISTRIBUTED", TransactionCommitted, map[string]interface{}{ "nodes": tx.Nodes, }) - return nil } -// AbortDistributedTransaction - отменяет распределённую транзакцию. func AbortDistributedTransaction(txID TransactionID) error { if globalTxManager == nil { return fmt.Errorf("transaction manager not initialized") @@ -3822,11 +3632,9 @@ func AbortDistributedTransaction(txID TransactionID) error { LogTransactionAudit(txID, "ABORT_DISTRIBUTED", TransactionAborted, map[string]interface{}{ "nodes": tx.Nodes, }) - return nil } -// applyOperation - применяет операцию к хранилищу. func applyOperation(op Operation) error { if globalStorage == nil { return fmt.Errorf("storage not initialized") @@ -3849,7 +3657,6 @@ func applyOperation(op Operation) error { doc.SetField(k, v) } doc.Version = op.Version - // Создаём MVCC версию при вставке if globalTxManager != nil && globalTxManager.mvccManager != nil { globalTxManager.mvccManager.CreateVersion(doc, 0) } @@ -3859,7 +3666,6 @@ func applyOperation(op Operation) error { if err := coll.Update(op.DocumentID, op.Data); err != nil { return err } - // Создаём MVCC версию при обновлении if globalTxManager != nil && globalTxManager.mvccManager != nil { if doc, err := coll.Find(op.DocumentID); err == nil { globalTxManager.mvccManager.CreateVersion(doc, 0) @@ -3870,11 +3676,9 @@ func applyOperation(op Operation) error { case "delete": return coll.Delete(op.DocumentID) } - return nil } -// CreateSavepoint - создаёт точку сохранения в транзакции. func (tx *Transaction) CreateSavepoint(name string) error { if TransactionState(tx.State.Load()) != TransactionActive { return fmt.Errorf("transaction is not active") @@ -3919,11 +3723,9 @@ func (tx *Transaction) CreateSavepoint(name string) error { "savepoint": name, "op_count": savepoint.OpCount, }) - return nil } -// RollbackToSavepoint - откатывает транзакцию к точке сохранения. func (tx *Transaction) RollbackToSavepoint(name string) error { if TransactionState(tx.State.Load()) != TransactionActive { return fmt.Errorf("transaction is not active") @@ -3947,17 +3749,14 @@ func (tx *Transaction) RollbackToSavepoint(name string) error { if len(tx.Operations) > tx.savepoints[targetIdx].OpCount { tx.Operations = tx.Operations[:tx.savepoints[targetIdx].OpCount] } - tx.savepoints = tx.savepoints[:targetIdx+1] LogTransactionAudit(tx.ID, "ROLLBACK_TO_SAVEPOINT", TransactionActive, map[string]interface{}{ "savepoint": name, }) - return nil } -// ReleaseSavepoint - освобождает точку сохранения. func (tx *Transaction) ReleaseSavepoint(name string) error { tx.mu.Lock() defer tx.mu.Unlock() @@ -3971,11 +3770,9 @@ func (tx *Transaction) ReleaseSavepoint(name string) error { return nil } } - return fmt.Errorf("savepoint '%s' not found", name) } -// GetSavepoints - возвращает список всех savepoints. func (tx *Transaction) GetSavepoints() []string { tx.mu.RLock() defer tx.mu.RUnlock() @@ -3987,7 +3784,6 @@ func (tx *Transaction) GetSavepoints() []string { return names } -// AddDocumentVersion - добавляет версию документа. func (tm *TransactionManager) AddDocumentVersion(docID string, version *DocumentVersion) { val, _ := tm.documentVersions.LoadOrStore(docID, make([]*DocumentVersion, 0)) versions := val.([]*DocumentVersion) @@ -4003,7 +3799,6 @@ func (tm *TransactionManager) AddDocumentVersion(docID string, version *Document } } -// GetDocumentVersion - получает версию документа по временной метке. func (tm *TransactionManager) GetDocumentVersion(docID string, timestamp int64) *Document { if tm.readCache != nil { if cached := tm.readCache.Get(docID, timestamp); cached != nil { @@ -4029,17 +3824,14 @@ func (tm *TransactionManager) GetDocumentVersion(docID string, timestamp int64) return nil } -// checkpointLoop - периодическое создание чекпоинтов. func (tm *TransactionManager) checkpointLoop() { ticker := time.NewTicker(time.Duration(tm.checkpointInterval) * time.Second) defer ticker.Stop() - for range ticker.C { tm.createCheckpoint() } } -// createCheckpoint - создаёт чекпоинт. func (tm *TransactionManager) createCheckpoint() { if tm.wal == nil { return @@ -4077,7 +3869,6 @@ func (tm *TransactionManager) createCheckpoint() { } } -// versionCleanupLoop - периодическая очистка старых версий. func (tm *TransactionManager) versionCleanupLoop() { if tm.maxVersions <= 0 { return @@ -4110,7 +3901,6 @@ func (tm *TransactionManager) versionCleanupLoop() { } } -// statsMonitor - мониторинг статистики. func (tm *TransactionManager) statsMonitor() { ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() @@ -4123,7 +3913,6 @@ func (tm *TransactionManager) statsMonitor() { } } -// startAsyncRecovery - запускает асинхронное восстановление. func (tm *TransactionManager) startAsyncRecovery() { if tm.wal == nil { tm.recoveryComplete.Store(true) @@ -4200,12 +3989,8 @@ func (tm *TransactionManager) startAsyncRecovery() { }() } -// IsRecoveryComplete - проверяет завершение восстановления. -func (tm *TransactionManager) IsRecoveryComplete() bool { - return tm.recoveryComplete.Load() -} +func (tm *TransactionManager) IsRecoveryComplete() bool { return tm.recoveryComplete.Load() } -// GetRecoveryProgress - возвращает прогресс восстановления. func (tm *TransactionManager) GetRecoveryProgress() map[string]interface{} { if tm.recoveryManager == nil { return map[string]interface{}{ @@ -4214,7 +3999,6 @@ func (tm *TransactionManager) GetRecoveryProgress() map[string]interface{} { "complete": true, } } - stats := tm.recoveryManager.GetStats() return map[string]interface{}{ "is_recovering": !tm.recoveryComplete.Load(), @@ -4224,29 +4008,21 @@ func (tm *TransactionManager) GetRecoveryProgress() map[string]interface{} { } } -// LockForBackup - блокирует транзакции для бэкапа. func (tm *TransactionManager) LockForBackup() { tm.backupLock.Lock() tm.backupInProgress.Store(true) } -// UnlockForBackup - разблокирует транзакции после бэкапа. func (tm *TransactionManager) UnlockForBackup() { tm.backupInProgress.Store(false) tm.backupLock.Unlock() } -// IsBackupInProgress - проверяет выполнение бэкапа. -func (tm *TransactionManager) IsBackupInProgress() bool { - return tm.backupInProgress.Load() -} +func (tm *TransactionManager) IsBackupInProgress() bool { return tm.backupInProgress.Load() } -// GetTransactionStats - возвращает статистику транзакций. func GetTransactionStats() map[string]interface{} { if globalTxManager == nil { - return map[string]interface{}{ - "error": "transaction manager not initialized", - } + return map[string]interface{}{"error": "transaction manager not initialized"} } stats := globalTxManager.stats @@ -4254,29 +4030,26 @@ func GetTransactionStats() map[string]interface{} { totalStarted := stats.TotalStarted.Load() totalCommitted := stats.TotalCommitted.Load() totalAborted := stats.TotalAborted.Load() - totalTimedOut := stats.TotalTimedOut.Load() - totalDeadlocks := stats.TotalDeadlocks.Load() return map[string]interface{}{ - "total_started": totalStarted, - "total_committed": totalCommitted, - "total_aborted": totalAborted, - "total_timed_out": totalTimedOut, - "total_deadlocks": totalDeadlocks, - "active_count": active, - "peak_active_count": stats.PeakActiveCount.Load(), - "max_ops_per_tx": stats.MaxOpsPerTx.Load(), - "avg_ops_per_tx": stats.AvgOpsPerTx.Load(), - "total_ops": stats.TotalOps.Load(), - "commit_rate": float64(totalCommitted) / float64(totalStarted+1) * 100, - "abort_rate": float64(totalAborted) / float64(totalStarted+1) * 100, - "uptime_seconds": time.Since(stats.StartTime).Seconds(), - "is_recovery_complete": globalTxManager.IsRecoveryComplete(), - "backup_in_progress": globalTxManager.IsBackupInProgress(), + "total_started": totalStarted, + "total_committed": totalCommitted, + "total_aborted": totalAborted, + "total_timed_out": stats.TotalTimedOut.Load(), + "total_deadlocks": stats.TotalDeadlocks.Load(), + "active_count": active, + "peak_active_count": stats.PeakActiveCount.Load(), + "max_ops_per_tx": stats.MaxOpsPerTx.Load(), + "avg_ops_per_tx": stats.AvgOpsPerTx.Load(), + "total_ops": stats.TotalOps.Load(), + "commit_rate": float64(totalCommitted) / float64(totalStarted+1) * 100, + "abort_rate": float64(totalAborted) / float64(totalStarted+1) * 100, + "uptime_seconds": time.Since(stats.StartTime).Seconds(), + "is_recovery_complete": globalTxManager.IsRecoveryComplete(), + "backup_in_progress": globalTxManager.IsBackupInProgress(), } } -// StopTransactionManager - останавливает менеджер транзакций. func StopTransactionManager() error { if globalTxManager == nil { return nil @@ -4287,20 +4060,14 @@ func StopTransactionManager() error { } if globalTxManager.wal != nil { - // Принудительная синхронизация перед закрытием globalTxManager.wal.Sync() return globalTxManager.wal.Close() } - return nil } -// HasActiveTransaction - проверяет наличие активной транзакции. -func HasActiveTransaction() bool { - return currentTx.Load() != nil -} +func HasActiveTransaction() bool { return currentTx.Load() != nil } -// GetCurrentTransactionID - возвращает ID текущей транзакции. func GetCurrentTransactionID() string { txVal := currentTx.Load() if txVal == nil { @@ -4310,7 +4077,6 @@ func GetCurrentTransactionID() string { return fmt.Sprintf("%d", tx.ID) } -// GetActiveTransactions - возвращает список активных транзакций. func GetActiveTransactions() []TransactionInfo { if globalTxManager == nil { return []TransactionInfo{} @@ -4358,11 +4124,9 @@ func GetActiveTransactions() []TransactionInfo { transactions = append(transactions, info) return true }) - return transactions } -// GetTransactionByID - возвращает транзакцию по ID. func GetTransactionByID(id string) (*Transaction, error) { if globalTxManager == nil { return nil, fmt.Errorf("transaction manager not initialized") @@ -4374,11 +4138,9 @@ func GetTransactionByID(id string) (*Transaction, error) { if val, ok := globalTxManager.activeTransactions.Load(txID); ok { return val.(*Transaction), nil } - return nil, fmt.Errorf("transaction not found") } -// AddToTransaction - добавляет операцию в текущую транзакцию. func AddToTransaction(coll *Collection, opType string, doc *Document) error { txVal := currentTx.Load() if txVal == nil { @@ -4407,11 +4169,9 @@ func AddToTransaction(coll *Collection, opType string, doc *Document) error { "operation": opType, "document": doc.ID, }) - return nil } -// FindInTransaction - находит документ в контексте транзакции. func FindInTransaction(coll *Collection, id string) (*Document, error) { txVal := currentTx.Load() if txVal == nil { @@ -4419,7 +4179,6 @@ func FindInTransaction(coll *Collection, id string) (*Document, error) { } tx := txVal.(*Transaction) - tx.mu.RLock() defer tx.mu.RUnlock() @@ -4445,16 +4204,11 @@ func FindInTransaction(coll *Collection, id string) (*Document, error) { return versionDoc, nil } } - return coll.Find(id) } -// MVCCSnapshot - создаёт снапшот MVCC. -func MVCCSnapshot() uint64 { - return uint64(time.Now().UnixNano()) -} +func MVCCSnapshot() uint64 { return uint64(time.Now().UnixNano()) } -// CreateDocumentVersion - создаёт версию документа. func CreateDocumentVersion(doc *Document, txID TransactionID) *DocumentVersion { return &DocumentVersion{ Document: doc.Clone(), @@ -4463,7 +4217,6 @@ func CreateDocumentVersion(doc *Document, txID TransactionID) *DocumentVersion { } } -// BeginTransactionOnCollection - начинает транзакцию на коллекции. func BeginTransactionOnCollection(coll *Collection) error { if globalTxManager == nil { if err := InitTransactionManager("futriis.wal"); err != nil { @@ -4479,11 +4232,9 @@ func BeginTransactionOnCollection(coll *Collection) error { if globalTxManager.logger != nil { globalTxManager.logger.Debug(fmt.Sprintf("Transaction %d started on collection %s.%s", tx.ID, coll.dbName, coll.name)) } - return nil } -// CheckTransactionTimeout - проверяет таймаут транзакции. func CheckTransactionTimeout(txID TransactionID) error { if globalTxManager == nil { return fmt.Errorf("transaction manager not initialized") @@ -4513,23 +4264,19 @@ func CheckTransactionTimeout(txID TransactionID) error { "timeout_ms": tx.timeout.Milliseconds(), "elapsed_ms": time.Since(time.UnixMilli(tx.StartTime)).Milliseconds(), }) - return fmt.Errorf("transaction %d timed out", txID) } - return nil } // ============================================================================= -// ФУНКЦИИ ДЛЯ КРОСС-ДАТАЦЕНТРОВОЙ МИГРАЦИИ +// ФУНКЦИИ МИГРАЦИИ // ============================================================================= -// GetDocumentAtTimestamp возвращает версию документа на указанный момент времени func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) { if globalTxManager == nil { return nil, fmt.Errorf("transaction manager not initialized") } - doc := globalTxManager.GetDocumentVersion(docID, timestamp) if doc == nil { return nil, fmt.Errorf("document %s not found at timestamp %d", docID, timestamp) @@ -4537,17 +4284,16 @@ func GetDocumentAtTimestamp(docID string, timestamp int64) (*Document, error) { return doc, nil } -// GetDocumentVersionsSince возвращает все версии документа после указанного времени func GetDocumentVersionsSince(docID string, since int64) ([]*DocumentVersion, error) { if globalTxManager == nil { return nil, fmt.Errorf("transaction manager not initialized") } - + val, ok := globalTxManager.documentVersions.Load(docID) if !ok { return nil, fmt.Errorf("no versions found for document %s", docID) } - + versions := val.([]*DocumentVersion) result := make([]*DocumentVersion, 0) for _, v := range versions { @@ -4558,25 +4304,24 @@ func GetDocumentVersionsSince(docID string, since int64) ([]*DocumentVersion, er return result, nil } -// CreateSnapshot создаёт снапшот всех версий для миграции func CreateSnapshot(database, collection string) (map[string][]*DocumentVersion, error) { if globalStorage == nil { return nil, fmt.Errorf("storage not initialized") } - + db, err := globalStorage.GetDatabase(database) if err != nil { return nil, err } - + coll, err := db.GetCollection(collection) if err != nil { return nil, err } - + docs := coll.GetAllDocuments() result := make(map[string][]*DocumentVersion) - + for _, doc := range docs { versions, err := GetDocumentVersionsSince(doc.ID, 0) if err != nil { @@ -4585,7 +4330,6 @@ func CreateSnapshot(database, collection string) (map[string][]*DocumentVersion, if len(versions) > 0 { result[doc.ID] = versions } else if globalTxManager != nil { - // Если нет версий в MVCC, создаём текущую version := &DocumentVersion{ Document: doc.Clone(), Timestamp: time.Now().UnixMilli(), @@ -4596,62 +4340,53 @@ func CreateSnapshot(database, collection string) (map[string][]*DocumentVersion, result[doc.ID] = []*DocumentVersion{version} } } - return result, nil } -// ApplySnapshot применяет снапшот к коллекции func ApplySnapshot(database, collection string, snapshot map[string][]*DocumentVersion) error { if globalStorage == nil { return fmt.Errorf("storage not initialized") } - + db, err := globalStorage.GetDatabase(database) if err != nil { return err } - + coll, err := db.GetCollection(collection) if err != nil { return err } - + for docID, versions := range snapshot { if len(versions) == 0 { continue } - - // Берём последнюю версию + latest := versions[len(versions)-1] doc := latest.Document - - // Проверяем, существует ли документ + existing, err := coll.Find(docID) if err == nil && existing != nil { - // Обновляем существующий документ updates := doc.GetFields() if err := coll.Update(docID, updates); err != nil { return fmt.Errorf("failed to update document %s: %v", docID, err) } } else { - // Вставляем новый документ if err := coll.Insert(doc); err != nil { return fmt.Errorf("failed to insert document %s: %v", docID, err) } } - - // Восстанавливаем все версии в MVCC + if globalTxManager != nil { for _, v := range versions { globalTxManager.AddDocumentVersion(docID, v) } } } - return nil } -// GetMVCCManager возвращает MVCC менеджер func GetMVCCManager() *MVCCManager { if globalTxManager == nil { return nil @@ -4659,32 +4394,34 @@ func GetMVCCManager() *MVCCManager { return globalTxManager.mvccManager } -// GetMVCCStats возвращает статистику MVCC func GetMVCCStats() map[string]interface{} { if globalTxManager == nil || globalTxManager.mvccManager == nil { - return map[string]interface{}{ - "error": "MVCC manager not initialized", - } + return map[string]interface{}{"error": "MVCC manager not initialized"} } return globalTxManager.mvccManager.GetMVCCStats() } -// GetVisibilityMapStats возвращает статистику карты видимости func GetVisibilityMapStats() map[string]interface{} { if globalTxManager == nil || globalTxManager.visibilityMap == nil { - return map[string]interface{}{ - "error": "Visibility map not initialized", - } + return map[string]interface{}{"error": "Visibility map not initialized"} } return globalTxManager.visibilityMap.GetStats() } -// GetReadCacheStats возвращает статистику кэша чтения func GetReadCacheStats() map[string]interface{} { if globalTxManager == nil || globalTxManager.readCache == nil { - return map[string]interface{}{ - "error": "Read cache not initialized", - } + return map[string]interface{}{"error": "Read cache not initialized"} } return globalTxManager.readCache.GetStats() } + +// ============================================================================= +// ВСПОМОГАТЕЛЬНЫЕ ФУНКЦИИ +// ============================================================================= + +func maxUint64(a, b uint64) uint64 { + if a > b { + return a + } + return b +} \ No newline at end of file