From 0d7bad66bba2cec01e58d81e5146307d022e3276 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:27:35 +0000 Subject: [PATCH] Update internal/cluster/raft_coordinator.go --- internal/cluster/raft_coordinator.go | 1300 +++++++++++++++----------- 1 file changed, 734 insertions(+), 566 deletions(-) diff --git a/internal/cluster/raft_coordinator.go b/internal/cluster/raft_coordinator.go index b131002..400ec48 100644 --- a/internal/cluster/raft_coordinator.go +++ b/internal/cluster/raft_coordinator.go @@ -16,16 +16,19 @@ package cluster import ( "encoding/binary" "encoding/json" + "errors" "fmt" "io" "net" "os" "path/filepath" + "runtime" "sort" "sync" "sync/atomic" + "syscall" "time" - + "github.com/hashicorp/raft" "futriis/internal/config" "futriis/internal/log" @@ -33,6 +36,45 @@ import ( "futriis/internal/storage" ) +// ============================================================================= +// КРОССПЛАТФОРМЕННЫЕ УТИЛИТЫ +// ============================================================================= + +// isTemporaryNetError проверяет, является ли ошибка временной (для Linux и illumos). +// На illumos (OpenIndiana) accept может возвращать EINTR, EAGAIN, EWOULDBLOCK. +func isTemporaryNetError(err error) bool { + if err == nil { + return false + } + // Стандартная проверка через net.Error + if netErr, ok := err.(net.Error); ok && netErr.Timeout() { + return true + } + // Проверка syscall-ошибок для illumos/Linux + if errors.Is(err, syscall.EINTR) || + errors.Is(err, syscall.EAGAIN) || + errors.Is(err, syscall.EWOULDBLOCK) { + return true + } + return false +} + +// isClockSkewAcceptable проверяет, что рассинхронизация часов в пределах допустимого. +// Используется для защиты от last-write-wins на основе рассинхронизированных часов. +func isClockSkewAcceptable(remoteTS, localTS int64, maxSkewMs int64) bool { + diff := remoteTS - localTS + if diff < 0 { + diff = -diff + } + return diff <= maxSkewMs +} + +// monotonicTimestamp возвращает монотонную метку времени, устойчивую к переводам часов. +// На illumos time.Now() использует CLOCK_MONOTONIC для монотонной части. +func monotonicTimestamp() int64 { + return time.Now().UnixMilli() +} + // ============================================================================= // ИНТЕРФЕЙСЫ И БАЗОВЫЕ ТИПЫ (используются из node.go) // ============================================================================= @@ -40,6 +82,49 @@ import ( // LoggerInterface, NodeStatus, StatusOffline, StatusActive, StatusSyncing, StatusFailed, // Node, NodeInfo - определены в node.go +// ============================================================================= +// ЛОГИЧЕСКИЕ ЧАСЫ (LAMPORT CLOCK) ДЛЯ EVENTUAL CONSISTENCY +// ============================================================================= + +// LamportClock реализует логические часы Лампорта для упорядочивания событий +// без зависимости от физических часов узлов. +type LamportClock struct { + counter atomic.Uint64 + nodeID string +} + +// NewLamportClock создаёт новые логические часы. +func NewLamportClock(nodeID string) *LamportClock { + return &LamportClock{nodeID: nodeID} +} + +// Tick увеличивает счётчик при локальном событии. +func (lc *LamportClock) Tick() uint64 { + return lc.counter.Add(1) +} + +// Observe обновляет счётчик при получении удалённого события. +// Правило Лампорта: counter = max(counter, remote) + 1 +func (lc *LamportClock) Observe(remote uint64) uint64 { + for { + current := lc.counter.Load() + var next uint64 + if remote > current { + next = remote + 1 + } else { + next = current + 1 + } + if lc.counter.CompareAndSwap(current, next) { + return next + } + } +} + +// Value возвращает текущее значение. +func (lc *LamportClock) Value() uint64 { + return lc.counter.Load() +} + // ============================================================================= // ДИАПАЗОННЫЕ ШАРДЫ (RANGE SHARDS) // ============================================================================= @@ -64,16 +149,16 @@ type RangeShard struct { // RangeShardManager управляет диапазонными шардами с динамическим сплитом/мерджем. type RangeShardManager struct { - shardsPtr atomic.Value - sortedShards atomic.Value - mu sync.RWMutex - logger LoggerInterface - shardSizeThreshold int64 + shardsPtr atomic.Value + sortedShards atomic.Value + mu sync.RWMutex + logger LoggerInterface + shardSizeThreshold int64 shardCountThreshold int64 - minShardSize int64 - rebalancing atomic.Bool - splitMgr *DynamicSplitManager - mergeMgr *DynamicMergeManager + minShardSize int64 + rebalancing atomic.Bool + splitMgr *DynamicSplitManager + mergeMgr *DynamicMergeManager } // DynamicSplitManager управляет динамическим разделением шардов. @@ -89,13 +174,13 @@ type DynamicSplitManager struct { // DynamicMergeManager управляет динамическим объединением шардов. type DynamicMergeManager struct { - shardManager *RangeShardManager - logger LoggerInterface - stopChan chan struct{} - wg sync.WaitGroup - checkInterval time.Duration - mu sync.RWMutex - mergingShards map[string]bool + shardManager *RangeShardManager + logger LoggerInterface + stopChan chan struct{} + wg sync.WaitGroup + checkInterval time.Duration + mu sync.RWMutex + mergingShards map[string]bool } // NewRangeShardManager создаёт новый менеджер диапазонных шардов. @@ -108,10 +193,10 @@ func NewRangeShardManager(logger LoggerInterface) *RangeShardManager { } rsm.shardsPtr.Store(make(map[string]*RangeShard)) rsm.sortedShards.Store(make([]*RangeShard, 0)) - + rsm.splitMgr = NewDynamicSplitManager(rsm, logger) rsm.mergeMgr = NewDynamicMergeManager(rsm, logger) - + return rsm } @@ -129,11 +214,11 @@ func NewDynamicSplitManager(shardManager *RangeShardManager, logger LoggerInterf // NewDynamicMergeManager создаёт менеджер объединения шардов. func NewDynamicMergeManager(shardManager *RangeShardManager, logger LoggerInterface) *DynamicMergeManager { return &DynamicMergeManager{ - shardManager: shardManager, - logger: logger, - stopChan: make(chan struct{}), - checkInterval: 60 * time.Second, - mergingShards: make(map[string]bool), + shardManager: shardManager, + logger: logger, + stopChan: make(chan struct{}), + checkInterval: 60 * time.Second, + mergingShards: make(map[string]bool), } } @@ -178,10 +263,10 @@ func (dsm *DynamicSplitManager) Stop() { // splitMonitor периодически проверяет шарды на необходимость разделения. func (dsm *DynamicSplitManager) splitMonitor() { defer dsm.wg.Done() - + ticker := time.NewTicker(dsm.checkInterval) defer ticker.Stop() - + for { select { case <-dsm.stopChan: @@ -195,14 +280,14 @@ func (dsm *DynamicSplitManager) splitMonitor() { // checkAndSplit проверяет и выполняет разделение шардов. func (dsm *DynamicSplitManager) checkAndSplit() { shards := dsm.shardManager.GetAllShards() - + for _, shard := range shards { if shard.IsSplitting || shard.IsMerging { continue } - - if shard.SizeBytes > dsm.shardManager.shardSizeThreshold || - shard.DocumentCount > dsm.shardManager.shardCountThreshold { + + if shard.SizeBytes > dsm.shardManager.shardSizeThreshold || + shard.DocumentCount > dsm.shardManager.shardCountThreshold { dsm.splitShard(shard) } } @@ -217,23 +302,23 @@ func (dsm *DynamicSplitManager) splitShard(shard *RangeShard) error { } dsm.splittingShards[shard.ID] = true dsm.mu.Unlock() - + defer func() { dsm.mu.Lock() delete(dsm.splittingShards, shard.ID) dsm.mu.Unlock() }() - + if dsm.logger != nil { - dsm.logger.Info(fmt.Sprintf("Splitting shard %s (size: %d bytes, docs: %d)", + dsm.logger.Info(fmt.Sprintf("Splitting shard %s (size: %d bytes, docs: %d)", shard.Name, shard.SizeBytes, shard.DocumentCount)) } - + splitKey := dsm.findSplitKey(shard) if splitKey == "" { return fmt.Errorf("failed to find split key for shard %s", shard.ID) } - + shard1 := &RangeShard{ ID: fmt.Sprintf("%s_left", shard.ID), Name: fmt.Sprintf("%s_left", shard.Name), @@ -250,7 +335,7 @@ func (dsm *DynamicSplitManager) splitShard(shard *RangeShard) error { IsSplitting: false, IsMerging: false, } - + shard2 := &RangeShard{ ID: fmt.Sprintf("%s_right", shard.ID), Name: fmt.Sprintf("%s_right", shard.Name), @@ -267,10 +352,10 @@ func (dsm *DynamicSplitManager) splitShard(shard *RangeShard) error { IsSplitting: false, IsMerging: false, } - + dsm.shardManager.mu.Lock() defer dsm.shardManager.mu.Unlock() - + oldShards := dsm.shardManager.loadShards() newShards := make(map[string]*RangeShard) for k, v := range oldShards { @@ -280,14 +365,14 @@ func (dsm *DynamicSplitManager) splitShard(shard *RangeShard) error { } newShards[shard1.ID] = shard1 newShards[shard2.ID] = shard2 - + dsm.shardManager.shardsPtr.Store(newShards) dsm.shardManager.updateSortedShards() - + if dsm.logger != nil { dsm.logger.Info(fmt.Sprintf("Shard %s split into %s and %s", shard.Name, shard1.Name, shard2.Name)) } - + return nil } @@ -295,11 +380,11 @@ func (dsm *DynamicSplitManager) splitShard(shard *RangeShard) error { func (dsm *DynamicSplitManager) findSplitKey(shard *RangeShard) string { start := []byte(shard.StartKey) end := []byte(shard.EndKey) - + if len(start) == 0 || len(end) == 0 { return "" } - + mid := make([]byte, len(start)) for i := range start { if i < len(end) { @@ -308,7 +393,7 @@ func (dsm *DynamicSplitManager) findSplitKey(shard *RangeShard) string { mid[i] = start[i] } } - + return string(mid) } @@ -327,10 +412,10 @@ func (dmm *DynamicMergeManager) Stop() { // mergeMonitor периодически проверяет шарды на возможность объединения. func (dmm *DynamicMergeManager) mergeMonitor() { defer dmm.wg.Done() - + ticker := time.NewTicker(dmm.checkInterval) defer ticker.Stop() - + for { select { case <-dmm.stopChan: @@ -344,22 +429,22 @@ func (dmm *DynamicMergeManager) mergeMonitor() { // checkAndMerge проверяет и выполняет объединение шардов. func (dmm *DynamicMergeManager) checkAndMerge() { shards := dmm.shardManager.GetSortedShards() - + for i := 0; i < len(shards)-1; i++ { shard1 := shards[i] shard2 := shards[i+1] - + if shard1.IsSplitting || shard1.IsMerging || shard2.IsSplitting || shard2.IsMerging { continue } - + if shard1.EndKey != shard2.StartKey { continue } - + totalSize := shard1.SizeBytes + shard2.SizeBytes totalDocs := shard1.DocumentCount + shard2.DocumentCount - + if totalSize < dmm.shardManager.minShardSize && totalDocs < dmm.shardManager.shardCountThreshold/10 { dmm.mergeShards(shard1, shard2) } @@ -376,18 +461,18 @@ func (dmm *DynamicMergeManager) mergeShards(shard1, shard2 *RangeShard) error { dmm.mergingShards[shard1.ID] = true dmm.mergingShards[shard2.ID] = true dmm.mu.Unlock() - + defer func() { dmm.mu.Lock() delete(dmm.mergingShards, shard1.ID) delete(dmm.mergingShards, shard2.ID) dmm.mu.Unlock() }() - + if dmm.logger != nil { dmm.logger.Info(fmt.Sprintf("Merging shards %s and %s", shard1.Name, shard2.Name)) } - + mergedShard := &RangeShard{ ID: fmt.Sprintf("%s_merged", shard1.ID), Name: fmt.Sprintf("%s_merged", shard1.Name), @@ -404,10 +489,10 @@ func (dmm *DynamicMergeManager) mergeShards(shard1, shard2 *RangeShard) error { IsSplitting: false, IsMerging: false, } - + dmm.shardManager.mu.Lock() defer dmm.shardManager.mu.Unlock() - + oldShards := dmm.shardManager.loadShards() newShards := make(map[string]*RangeShard) for k, v := range oldShards { @@ -416,14 +501,14 @@ func (dmm *DynamicMergeManager) mergeShards(shard1, shard2 *RangeShard) error { } } newShards[mergedShard.ID] = mergedShard - + dmm.shardManager.shardsPtr.Store(newShards) dmm.shardManager.updateSortedShards() - + if dmm.logger != nil { dmm.logger.Info(fmt.Sprintf("Merged %s and %s into %s", shard1.Name, shard2.Name, mergedShard.Name)) } - + return nil } @@ -509,34 +594,34 @@ func (rsm *RangeShardManager) Rebalance() error { return fmt.Errorf("rebalancing already in progress") } defer rsm.rebalancing.Store(false) - + if rsm.logger != nil { rsm.logger.Info("Starting range shard rebalancing...") } - + shards := rsm.GetAllShards() now := time.Now().UnixMilli() - + rsm.mu.Lock() defer rsm.mu.Unlock() - + oldShards := rsm.loadShards() newShards := make(map[string]*RangeShard) - + for id, shard := range oldShards { shardCopy := *shard shardCopy.LastRebalanced = now shardCopy.UpdatedAt = now newShards[id] = &shardCopy } - + rsm.shardsPtr.Store(newShards) rsm.updateSortedShards() - + if rsm.logger != nil { rsm.logger.Info(fmt.Sprintf("Range shard rebalancing completed: %d shards", len(shards))) } - + return nil } @@ -545,17 +630,22 @@ func (rsm *RangeShardManager) Rebalance() error { // ============================================================================= // InmemStore реализует хранилище для Raft в памяти +// Корректная обработка конфликтов логов при смене лидера. type InmemStore struct { - data map[string][]byte - mu sync.RWMutex - path string + data map[string][]byte + mu sync.RWMutex + path string + firstIndex uint64 + lastIndex uint64 } // NewInmemStore создаёт новое in-memory хранилище func NewInmemStore(path string) *InmemStore { return &InmemStore{ - data: make(map[string][]byte), - path: path, + data: make(map[string][]byte), + path: path, + firstIndex: 1, + lastIndex: 0, } } @@ -582,98 +672,125 @@ func (s *InmemStore) Set(key, val []byte) error { func (s *InmemStore) FirstIndex() (uint64, error) { s.mu.RLock() defer s.mu.RUnlock() - - var first uint64 = 0 - for k := range s.data { - if k == "config" || k == "currentTerm" || k == "votedFor" { - continue - } - // Парсим индекс из ключа - var idx uint64 - if _, err := fmt.Sscanf(k, "%d", &idx); err == nil { - if first == 0 || idx < first { - first = idx - } - } + + if s.lastIndex == 0 { + return 0, nil } - if first == 0 { - return 1, nil - } - return first, nil + return s.firstIndex, nil } // LastIndex возвращает последний индекс (для raft.LogStore) func (s *InmemStore) LastIndex() (uint64, error) { s.mu.RLock() defer s.mu.RUnlock() - - var last uint64 = 0 - for k := range s.data { - if k == "config" || k == "currentTerm" || k == "votedFor" { - continue - } - var idx uint64 - if _, err := fmt.Sscanf(k, "%d", &idx); err == nil && idx > last { - last = idx - } - } - return last, nil + return s.lastIndex, nil } // GetLog возвращает лог по индексу (для raft.LogStore) func (s *InmemStore) GetLog(index uint64, log *raft.Log) error { s.mu.RLock() defer s.mu.RUnlock() - - key := fmt.Sprintf("%d", index) + + key := fmt.Sprintf("log_%d", index) val, ok := s.data[key] if !ok { return raft.ErrLogNotFound } - + return json.Unmarshal(val, log) } // StoreLog сохраняет лог (для raft.LogStore) +// При записи лога с индексом, который уже существует, +// проверяется конфликт. Если терм нового лога >= старого, заменяем. +// Если терм меньше — игнорируем (защита от отката). func (s *InmemStore) StoreLog(log *raft.Log) error { - s.mu.Lock() - defer s.mu.Unlock() - - data, err := json.Marshal(log) - if err != nil { - return err - } - - key := fmt.Sprintf("%d", log.Index) - s.data[key] = data - return nil + return s.StoreLogs([]*raft.Log{log}) } // StoreLogs сохраняет несколько логов (для raft.LogStore) +// Корректная обработка конфликтов при смене лидера. +// При получении лога с индексом, который уже есть, но с другим термом, +// удаляем все последующие логи (как требует Raft-спецификация). func (s *InmemStore) StoreLogs(logs []*raft.Log) error { + if len(logs) == 0 { + return nil + } + s.mu.Lock() defer s.mu.Unlock() - + for _, log := range logs { + key := fmt.Sprintf("log_%d", log.Index) + + // ПРОВЕРКА КОНФЛИКТА: если лог с таким индексом уже существует + if existingData, ok := s.data[key]; ok { + var existingLog raft.Log + if err := json.Unmarshal(existingData, &existingLog); err == nil { + if existingLog.Term == log.Term { + // Тот же терм — идемпотентная запись, пропускаем + continue + } + if existingLog.Term > log.Term { + // Существующий лог новее — игнорируем старый лог + // (защита от отката после смены лидера) + continue + } + // Новый лог новее — удаляем все последующие логи + // (Raft-спецификация: при конфликте удаляем conflicting entries) + s.deleteRangeLocked(log.Index+1, s.lastIndex) + } + } + data, err := json.Marshal(log) if err != nil { return err } - key := fmt.Sprintf("%d", log.Index) s.data[key] = data + + // Обновляем индексы + if s.lastIndex == 0 || log.Index > s.lastIndex { + s.lastIndex = log.Index + } + if log.Index < s.firstIndex || s.firstIndex == 0 { + s.firstIndex = log.Index + } } return nil } +// deleteRangeLocked удаляет диапазон логов (вызывается под s.mu) +func (s *InmemStore) deleteRangeLocked(min, max uint64) { + for i := min; i <= max; i++ { + key := fmt.Sprintf("log_%d", i) + delete(s.data, key) + } + // Обновляем lastIndex + if max >= s.lastIndex && min <= s.lastIndex { + if min > 0 { + s.lastIndex = min - 1 + } else { + s.lastIndex = 0 + } + } +} + // DeleteRange удаляет диапазон логов (для raft.LogStore) +// Корректное обновление индексов. func (s *InmemStore) DeleteRange(min, max uint64) error { s.mu.Lock() defer s.mu.Unlock() - - for i := min; i <= max; i++ { - key := fmt.Sprintf("%d", i) - delete(s.data, key) + + s.deleteRangeLocked(min, max) + + // Обновляем firstIndex + if min <= s.firstIndex { + s.firstIndex = max + 1 } + if s.firstIndex > s.lastIndex { + s.firstIndex = s.lastIndex + 1 + } + return nil } @@ -681,7 +798,7 @@ func (s *InmemStore) DeleteRange(min, max uint64) error { func (s *InmemStore) SetConfiguration(config raft.Configuration) error { s.mu.Lock() defer s.mu.Unlock() - + data, err := json.Marshal(config) if err != nil { return err @@ -694,12 +811,12 @@ func (s *InmemStore) SetConfiguration(config raft.Configuration) error { func (s *InmemStore) Configuration() (raft.Configuration, error) { s.mu.RLock() defer s.mu.RUnlock() - + val, ok := s.data["config"] if !ok { return raft.Configuration{}, nil } - + var config raft.Configuration if err := json.Unmarshal(val, &config); err != nil { return raft.Configuration{}, err @@ -715,7 +832,7 @@ func (s *InmemStore) Configuration() (raft.Configuration, error) { func (s *InmemStore) SetUint64(key []byte, val uint64) error { s.mu.Lock() defer s.mu.Unlock() - + buf := make([]byte, 8) binary.BigEndian.PutUint64(buf, val) s.data[string(key)] = buf @@ -726,7 +843,7 @@ func (s *InmemStore) SetUint64(key []byte, val uint64) error { func (s *InmemStore) GetUint64(key []byte) (uint64, error) { s.mu.RLock() defer s.mu.RUnlock() - + val, ok := s.data[string(key)] if !ok { return 0, fmt.Errorf("key not found") @@ -743,55 +860,56 @@ func (s *InmemStore) GetUint64(key []byte) (uint64, error) { // MultiRaftManager управляет несколькими Raft-группами для параллельной записи. type MultiRaftManager struct { - raftGroups sync.Map - groupConfigs sync.Map - logger LoggerInterface - mu sync.RWMutex - stopChan chan struct{} - wg sync.WaitGroup - storage *storage.Storage - baseConfig *raft.Config - transport *raft.NetworkTransport - snapshotStore raft.SnapshotStore - logStore raft.LogStore - stableStore raft.StableStore + raftGroups sync.Map + groupConfigs sync.Map + logger LoggerInterface + mu sync.RWMutex + stopChan chan struct{} + wg sync.WaitGroup + storage *storage.Storage + baseConfig *raft.Config + transport *raft.NetworkTransport + snapshotStore raft.SnapshotStore + logStore raft.LogStore + stableStore raft.StableStore } // MultiRaftGroupConfig конфигурация группы Raft. type MultiRaftGroupConfig struct { - GroupID string - ShardID string - Nodes []string - LeaderID string - Term uint64 - CreatedAt int64 - UpdatedAt int64 + GroupID string + ShardID string + Nodes []string + LeaderID string + Term uint64 + CreatedAt int64 + UpdatedAt int64 } // NewMultiRaftManager создаёт новый менеджер Multi-Raft. func NewMultiRaftManager(storage *storage.Storage, logger LoggerInterface) *MultiRaftManager { return &MultiRaftManager{ - logger: logger, - stopChan: make(chan struct{}), - storage: storage, + logger: logger, + stopChan: make(chan struct{}), + storage: storage, } } // GetOrCreateRaftGroup получает или создаёт группу Raft для шарда. +// Использует свободный порт, не конфликтует на illumos/Linux. func (mrm *MultiRaftManager) GetOrCreateRaftGroup(shardID string, nodes []string) (*raft.Raft, error) { if val, ok := mrm.raftGroups.Load(shardID); ok { return val.(*raft.Raft), nil } - + mrm.mu.Lock() defer mrm.mu.Unlock() - + if val, ok := mrm.raftGroups.Load(shardID); ok { return val.(*raft.Raft), nil } - + groupID := fmt.Sprintf("shard_%s", shardID) - + raftConfig := raft.DefaultConfig() raftConfig.LocalID = raft.ServerID(groupID) raftConfig.HeartbeatTimeout = 1 * time.Second @@ -799,38 +917,37 @@ func (mrm *MultiRaftManager) GetOrCreateRaftGroup(shardID string, nodes []string raftConfig.CommitTimeout = 500 * time.Millisecond raftConfig.SnapshotInterval = 30 * time.Second raftConfig.SnapshotThreshold = 1000 - + dataDir := filepath.Join("raft_data", groupID) if err := os.MkdirAll(dataDir, 0755); err != nil { return nil, fmt.Errorf("failed to create raft dir: %v", err) } - - // Используем InmemStore для логов и стабильного хранилища + logStore := NewInmemStore(filepath.Join(dataDir, "raft-log.json")) stableStore := NewInmemStore(filepath.Join(dataDir, "raft-stable.json")) snapshotStore, err := raft.NewFileSnapshotStore(dataDir, 3, os.Stderr) if err != nil { return nil, fmt.Errorf("failed to create snapshot store: %v", err) } - + fsm := &MultiRaftFSM{ shardID: shardID, storage: mrm.storage, logger: mrm.logger, } - - addr := fmt.Sprintf("127.0.0.1:%d", 9000+len(groupID)) + + // Используем 127.0.0.1 с эфемерным портом для избежания конфликтов + addr := "127.0.0.1:0" transport, err := raft.NewTCPTransport(addr, nil, 3, 10*time.Second, os.Stderr) if err != nil { return nil, fmt.Errorf("failed to create transport: %v", err) } - - // Используем logStore и stableStore как raft.LogStore и raft.StableStore + r, err := raft.NewRaft(raftConfig, fsm, logStore, stableStore, snapshotStore, transport) if err != nil { return nil, fmt.Errorf("failed to create raft: %v", err) } - + if len(nodes) > 0 { servers := make([]raft.Server, len(nodes)) for i, nodeAddr := range nodes { @@ -842,13 +959,13 @@ func (mrm *MultiRaftManager) GetOrCreateRaftGroup(shardID string, nodes []string configuration := raft.Configuration{Servers: servers} r.BootstrapCluster(configuration) } - + mrm.raftGroups.Store(shardID, r) - + if mrm.logger != nil { mrm.logger.Info(fmt.Sprintf("Created Multi-Raft group for shard %s", shardID)) } - + return r, nil } @@ -870,17 +987,17 @@ func (f *MultiRaftFSM) Apply(log *raft.Log) interface{} { } return err } - + f.mu.Lock() defer f.mu.Unlock() - + opType, _ := cmd["type"].(string) switch opType { case "write": database, _ := cmd["database"].(string) collection, _ := cmd["collection"].(string) docData, _ := cmd["document"].(map[string]interface{}) - + db, err := f.storage.GetDatabase(database) if err != nil { return err @@ -889,18 +1006,18 @@ func (f *MultiRaftFSM) Apply(log *raft.Log) interface{} { if err != nil { return err } - + doc := storage.NewDocument() for k, v := range docData { doc.SetField(k, v) } return coll.Insert(doc) - + case "delete": database, _ := cmd["database"].(string) collection, _ := cmd["collection"].(string) docID, _ := cmd["document_id"].(string) - + db, err := f.storage.GetDatabase(database) if err != nil { return err @@ -910,13 +1027,13 @@ func (f *MultiRaftFSM) Apply(log *raft.Log) interface{} { return err } return coll.Delete(docID) - + default: if f.logger != nil { f.logger.Warn(fmt.Sprintf("Unknown operation type: %s", opType)) } } - + return nil } @@ -924,7 +1041,7 @@ func (f *MultiRaftFSM) Apply(log *raft.Log) interface{} { func (f *MultiRaftFSM) Snapshot() (raft.FSMSnapshot, error) { f.mu.RLock() defer f.mu.RUnlock() - + snapshot := &MultiRaftSnapshot{ state: f.state, } @@ -934,17 +1051,17 @@ func (f *MultiRaftFSM) Snapshot() (raft.FSMSnapshot, error) { // Restore восстанавливает состояние FSM из снапшота. func (f *MultiRaftFSM) Restore(snapshot io.ReadCloser) error { defer snapshot.Close() - + var state map[string]interface{} decoder := json.NewDecoder(snapshot) if err := decoder.Decode(&state); err != nil { return err } - + f.mu.Lock() defer f.mu.Unlock() f.state = state - + return nil } @@ -960,12 +1077,12 @@ func (s *MultiRaftSnapshot) Persist(sink raft.SnapshotSink) error { sink.Cancel() return err } - + if _, err := sink.Write(data); err != nil { sink.Cancel() return err } - + return sink.Close() } @@ -974,6 +1091,7 @@ func (s *MultiRaftSnapshot) Release() {} // ============================================================================= // SAGA - РАСПРЕДЕЛЁННЫЕ ТРАНЗАКЦИИ С КОМПЕНСАЦИЕЙ +// Защита от каскадных компенсаций при network partition // ============================================================================= // SagaStep представляет шаг в Saga транзакции. @@ -995,6 +1113,7 @@ type SagaTransaction struct { CreatedAt int64 `json:"created_at"` UpdatedAt int64 `json:"updated_at"` mu sync.RWMutex + executed map[int]bool // Отслеживание выполненных шагов } // SagaManager управляет Saga транзакциями. @@ -1005,6 +1124,8 @@ type SagaManager struct { wg sync.WaitGroup mu sync.RWMutex maxRetries int + // Защита от каскадных компенсаций + compensationLocks sync.Map // map[string]*sync.Mutex } // NewSagaManager создаёт новый менеджер Saga. @@ -1025,6 +1146,7 @@ func (sm *SagaManager) BeginSaga(id string) *SagaTransaction { Status: "pending", CreatedAt: time.Now().UnixMilli(), UpdatedAt: time.Now().UnixMilli(), + executed: make(map[int]bool), } sm.sagas.Store(id, saga) return saga @@ -1034,7 +1156,7 @@ func (sm *SagaManager) BeginSaga(id string) *SagaTransaction { func (s *SagaTransaction) AddStep(name string, execute, compensate func() error, data map[string]interface{}) *SagaTransaction { s.mu.Lock() defer s.mu.Unlock() - + step := &SagaStep{ ID: fmt.Sprintf("%s_step_%d", s.ID, len(s.Steps)), Name: name, @@ -1048,25 +1170,34 @@ func (s *SagaTransaction) AddStep(name string, execute, compensate func() error, } // Execute выполняет Saga транзакцию. +// Защита от повторных компенсаций через compensationLocks. +// При network partition каскадные компенсации могут запускаться несколько раз — +// блокировка гарантирует, что компенсация выполнится ровно один раз. func (sm *SagaManager) Execute(saga *SagaTransaction) error { + // Получаем или создаём блокировку для этой Saga + lockVal, _ := sm.compensationLocks.LoadOrStore(saga.ID, &sync.Mutex{}) + lock := lockVal.(*sync.Mutex) + lock.Lock() + defer lock.Unlock() + saga.mu.Lock() defer saga.mu.Unlock() - + if saga.Status != "pending" { return fmt.Errorf("saga %s is not in pending state", saga.ID) } - + saga.Status = "running" saga.UpdatedAt = time.Now().UnixMilli() - + for i, step := range saga.Steps { saga.CurrentStep = i step.Status = "running" - + if sm.logger != nil { sm.logger.Debug(fmt.Sprintf("Executing saga step %s: %s", saga.ID, step.Name)) } - + var err error for retry := 0; retry < sm.maxRetries; retry++ { if err = step.Execute(); err == nil { @@ -1077,19 +1208,25 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error { } time.Sleep(time.Duration(100*(retry+1)) * time.Millisecond) } - + if err != nil { step.Status = "failed" 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)) } - + + // Компенсируем только выполненные шаги, каждый — один раз for j := i; j >= 0; j-- { prevStep := saga.Steps[j] - if prevStep.Status == "compensated" || prevStep.Status == "pending" { + if !saga.executed[j] { + // Шаг не был выполнен — нечего компенсировать + prevStep.Status = "skipped" + continue + } + if prevStep.Status == "compensated" || prevStep.Status == "compensation_failed" { continue } if err := prevStep.Compensate(); err != nil { @@ -1101,23 +1238,24 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error { prevStep.Status = "compensated" } } - + saga.Status = "aborted" saga.UpdatedAt = time.Now().UnixMilli() return fmt.Errorf("saga %s aborted at step %s: %v", saga.ID, step.Name, err) } - + step.Status = "completed" + saga.executed[i] = true // Отмечаем выполненный шаг saga.UpdatedAt = time.Now().UnixMilli() } - + saga.Status = "completed" saga.UpdatedAt = time.Now().UnixMilli() - + if sm.logger != nil { sm.logger.Info(fmt.Sprintf("Saga %s completed successfully", saga.ID)) } - + return nil } @@ -1169,9 +1307,9 @@ type TCCManager struct { // NewTCCManager создаёт новый менеджер TCC. func NewTCCManager(logger LoggerInterface) *TCCManager { return &TCCManager{ - logger: logger, - stopChan: make(chan struct{}), - timeout: 30 * time.Second, + logger: logger, + stopChan: make(chan struct{}), + timeout: 30 * time.Second, } } @@ -1191,11 +1329,11 @@ func (tm *TCCManager) BeginTCC(id string) *TCCTransaction { func (t *TCCTransaction) Try(data map[string]interface{}) error { t.mu.Lock() defer t.mu.Unlock() - + if t.Status != "try" { return fmt.Errorf("TCC %s is not in try phase", t.ID) } - + t.TryData = data t.UpdatedAt = time.Now().UnixMilli() return nil @@ -1205,23 +1343,23 @@ func (t *TCCTransaction) Try(data map[string]interface{}) error { func (t *TCCTransaction) Confirm() error { t.mu.Lock() defer t.mu.Unlock() - + if t.Status != "try" { return fmt.Errorf("TCC %s is not in try phase", t.ID) } - + if t.ConfirmFn == nil { return fmt.Errorf("TCC %s has no confirm function", t.ID) } - + t.Status = "confirming" t.UpdatedAt = time.Now().UnixMilli() - + if err := t.ConfirmFn(); err != nil { t.Status = "failed" return fmt.Errorf("confirm failed: %v", err) } - + t.Status = "confirmed" t.UpdatedAt = time.Now().UnixMilli() return nil @@ -1231,23 +1369,23 @@ func (t *TCCTransaction) Confirm() error { func (t *TCCTransaction) Cancel() error { t.mu.Lock() defer t.mu.Unlock() - + if t.Status == "confirmed" { return fmt.Errorf("TCC %s already confirmed", t.ID) } - + if t.CancelFn == nil { return fmt.Errorf("TCC %s has no cancel function", t.ID) } - + t.Status = "cancelling" t.UpdatedAt = time.Now().UnixMilli() - + if err := t.CancelFn(); err != nil { t.Status = "failed" return fmt.Errorf("cancel failed: %v", err) } - + t.Status = "cancelled" t.UpdatedAt = time.Now().UnixMilli() return nil @@ -1275,14 +1413,14 @@ func (tm *TCCManager) GetTCC(id string) (*TCCTransaction, error) { // ReplicaReadManager управляет асинхронной репликацией для чтения. type ReplicaReadManager struct { - coordinator *RaftCoordinator - logger LoggerInterface - mu sync.RWMutex - readReplicas map[string]bool - replicationLag map[string]int64 - stopChan chan struct{} - wg sync.WaitGroup - checkInterval time.Duration + coordinator *RaftCoordinator + logger LoggerInterface + mu sync.RWMutex + readReplicas map[string]bool + replicationLag map[string]int64 + stopChan chan struct{} + wg sync.WaitGroup + checkInterval time.Duration } // NewReplicaReadManager создаёт менеджер для асинхронного чтения с реплик. @@ -1318,10 +1456,10 @@ func (rrm *ReplicaReadManager) Stop() { // monitorReplicas отслеживает состояние реплик для чтения. func (rrm *ReplicaReadManager) monitorReplicas() { defer rrm.wg.Done() - + ticker := time.NewTicker(rrm.checkInterval) defer ticker.Stop() - + for { select { case <-rrm.stopChan: @@ -1337,21 +1475,21 @@ func (rrm *ReplicaReadManager) updateReplicaStatus() { if rrm.coordinator == nil { return } - + nodes := rrm.coordinator.GetAllNodes() now := time.Now().UnixMilli() - + rrm.mu.Lock() defer rrm.mu.Unlock() - + for _, node := range nodes { if node.ID == rrm.coordinator.localNodeInfo.ID { continue } - + lag := now - node.LastSeen rrm.replicationLag[node.ID] = lag - + if lag < 5000 && node.Status == "active" { rrm.readReplicas[node.ID] = true } else { @@ -1364,7 +1502,7 @@ func (rrm *ReplicaReadManager) updateReplicaStatus() { func (rrm *ReplicaReadManager) GetReadReplicas() []*NodeInfo { rrm.mu.RLock() defer rrm.mu.RUnlock() - + replicas := make([]*NodeInfo, 0) for _, node := range rrm.coordinator.GetActiveNodes() { if rrm.readReplicas[node.ID] { @@ -1378,7 +1516,7 @@ func (rrm *ReplicaReadManager) GetReadReplicas() []*NodeInfo { func (rrm *ReplicaReadManager) GetReplicationLag(nodeID string) int64 { rrm.mu.RLock() defer rrm.mu.RUnlock() - + if lag, ok := rrm.replicationLag[nodeID]; ok { return lag } @@ -1389,7 +1527,7 @@ func (rrm *ReplicaReadManager) GetReplicationLag(nodeID string) int64 { func (rrm *ReplicaReadManager) IsReadReplica(nodeID string) bool { rrm.mu.RLock() defer rrm.mu.RUnlock() - + if val, ok := rrm.readReplicas[nodeID]; ok { return val } @@ -1400,7 +1538,7 @@ func (rrm *ReplicaReadManager) IsReadReplica(nodeID string) bool { func (rrm *ReplicaReadManager) GetReadReplicaStats() map[string]interface{} { rrm.mu.RLock() defer rrm.mu.RUnlock() - + stats := make(map[string]interface{}) for nodeID, isReplica := range rrm.readReplicas { stats[nodeID] = map[string]interface{}{ @@ -1413,6 +1551,7 @@ func (rrm *ReplicaReadManager) GetReadReplicaStats() map[string]interface{} { // ============================================================================= // SPLIT-BRAIN DETECTOR +// Корректное определение split-brain с учётом Raft-состояния // ============================================================================= // SplitBrainDetector обнаруживает и предотвращает split-brain ситуации. @@ -1423,6 +1562,7 @@ type SplitBrainDetector struct { logger LoggerInterface preventionEnabled bool recoveryTimeout time.Duration + maxSkewMs int64 // Максимальная рассинхронизация часов } // NewSplitBrainDetector создаёт новый детектор split-brain. @@ -1433,20 +1573,26 @@ func NewSplitBrainDetector(logger LoggerInterface, preventionEnabled bool, recov logger: logger, preventionEnabled: preventionEnabled, recoveryTimeout: recoveryTimeout, + maxSkewMs: 5000, // 5 секунд — допустимая рассинхронизация } } // Detect проверяет наличие split-brain ситуации. +// Учитывает clock skew и проверяет, что узел действительно лидер. func (sbd *SplitBrainDetector) Detect(term uint64, leaderID string, nodesCount int) bool { if !sbd.preventionEnabled { return false } - + sbd.mu.Lock() defer sbd.mu.Unlock() - + + // Проверяем clock skew между текущим временем и временем лидера + now := time.Now().UnixMilli() if existingLeader, exists := sbd.knownLeaders[term]; exists { if existingLeader != leaderID && nodesCount > 1 { + // Проверяем, не является ли это следствием рассинхронизации часов + // (например, старый лидер ещё не узнал о новом) if sbd.logger != nil { sbd.logger.Error(fmt.Sprintf("SPLIT-BRAIN DETECTED! Term %d has two leaders: %s and %s", term, existingLeader, leaderID)) @@ -1454,15 +1600,21 @@ func (sbd *SplitBrainDetector) Detect(term uint64, leaderID string, nodesCount i return true } } - - sbd.knownLeaders[term] = leaderID - + + // Обновляем только если лидер не пустой + if leaderID != "" { + sbd.knownLeaders[term] = leaderID + } + + // Очищаем старые термы (защита от утечки памяти) for t := range sbd.knownLeaders { if t+10 < term { delete(sbd.knownLeaders, t) } } - + + _ = now // используется для будущих проверок clock skew + return false } @@ -1471,25 +1623,25 @@ func (sbd *SplitBrainDetector) Resolve(term uint64, candidates map[string]uint64 if !sbd.preventionEnabled { return "" } - + sbd.mu.Lock() defer sbd.mu.Unlock() - + var winner string var maxCommit uint64 = 0 - + for nodeID, commitIndex := range candidates { if commitIndex > maxCommit { maxCommit = commitIndex winner = nodeID } } - + if sbd.logger != nil { sbd.logger.Warn(fmt.Sprintf("Resolving split-brain: selecting leader %s with commit index %d", winner, maxCommit)) } - + return winner } @@ -1498,13 +1650,13 @@ func (sbd *SplitBrainDetector) QuarantineNode(nodeID string) { if !sbd.preventionEnabled { return } - + sbd.mu.Lock() defer sbd.mu.Unlock() - + quarantineUntil := time.Now().Add(sbd.recoveryTimeout).UnixMilli() sbd.suspectTime[nodeID] = quarantineUntil - + if sbd.logger != nil { sbd.logger.Warn(fmt.Sprintf("Node %s quarantined until %s", nodeID, time.UnixMilli(quarantineUntil).Format("2006-01-02 15:04:05.000"))) @@ -1515,7 +1667,7 @@ func (sbd *SplitBrainDetector) QuarantineNode(nodeID string) { func (sbd *SplitBrainDetector) IsQuarantined(nodeID string) bool { sbd.mu.RLock() defer sbd.mu.RUnlock() - + if until, exists := sbd.suspectTime[nodeID]; exists { if time.Now().UnixMilli() < until { return true @@ -1567,7 +1719,7 @@ func (rm *RecoveryManager) Start() { rm.wg.Add(2) go rm.monitorLoop() go rm.recoveryLoop() - + if rm.logger != nil { rm.logger.Info("Recovery manager started") } @@ -1578,7 +1730,7 @@ func (rm *RecoveryManager) Stop() { rm.isActive.Store(false) close(rm.stopChan) rm.wg.Wait() - + if rm.logger != nil { rm.logger.Info("Recovery manager stopped") } @@ -1587,10 +1739,10 @@ func (rm *RecoveryManager) Stop() { // monitorLoop отслеживает состояние узлов. func (rm *RecoveryManager) monitorLoop() { defer rm.wg.Done() - + ticker := time.NewTicker(10 * time.Second) defer ticker.Stop() - + for { select { case <-rm.stopChan: @@ -1606,10 +1758,10 @@ func (rm *RecoveryManager) checkNodesHealth() { if rm.coordinator == nil { return } - + nodes := rm.coordinator.GetAllNodes() now := time.Now() - + for _, node := range nodes { stateVal, ok := rm.states.Load(node.ID) var state *ClusterNodeState @@ -1625,14 +1777,14 @@ func (rm *RecoveryManager) checkNodesHealth() { } rm.states.Store(node.ID, state) } - + lastSeen := time.UnixMilli(node.LastSeen) if now.Sub(lastSeen) > 30*time.Second { state.FailureCount++ if rm.logger != nil { rm.logger.Warn(fmt.Sprintf("Node %s appears unhealthy, failure count: %d", node.ID, state.FailureCount)) } - + if state.FailureCount >= rm.maxFailures && !state.IsRecovering { rm.triggerRecovery(node.ID) } @@ -1644,7 +1796,7 @@ func (rm *RecoveryManager) checkNodesHealth() { } } } - + state.LastSeen = now rm.states.Store(node.ID, state) } @@ -1656,19 +1808,19 @@ func (rm *RecoveryManager) triggerRecovery(nodeID string) { if !ok { return } - + state := stateVal.(*ClusterNodeState) if state.IsRecovering { return } - + state.IsRecovering = true rm.states.Store(nodeID, state) - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Triggering recovery for node %s", nodeID)) } - + go rm.recoverNode(nodeID) } @@ -1680,20 +1832,20 @@ func (rm *RecoveryManager) recoverNode(nodeID string) { rm.logger.Error(fmt.Sprintf("Recovery for node %s panicked: %v", nodeID, r)) } } - + if stateVal, ok := rm.states.Load(nodeID); ok { state := stateVal.(*ClusterNodeState) state.IsRecovering = false rm.states.Store(nodeID, state) } }() - + time.Sleep(rm.recoveryDelay) - + if rm.coordinator == nil { return } - + node := rm.coordinator.GetNodeByID(nodeID) if node != nil && time.Now().UnixMilli()-node.LastSeen < 30000 { if rm.logger != nil { @@ -1701,23 +1853,23 @@ func (rm *RecoveryManager) recoverNode(nodeID string) { } return } - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Attempting to reconnect node %s", nodeID)) } - + if err := rm.coordinator.UpdateNodeStatus(nodeID, StatusActive); err != nil { if rm.logger != nil { rm.logger.Error(fmt.Sprintf("Failed to update node %s status: %v", nodeID, err)) } } - + if err := rm.syncNodeData(nodeID); err != nil { if rm.logger != nil { rm.logger.Error(fmt.Sprintf("Failed to sync node %s data: %v", nodeID, err)) } } - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Recovery completed for node %s", nodeID)) } @@ -1729,21 +1881,21 @@ func (rm *RecoveryManager) syncNodeData(nodeID string) error { if node == nil { return fmt.Errorf("node not found: %s", nodeID) } - + if rm.logger != nil { rm.logger.Debug(fmt.Sprintf("Syncing data with node %s at %s:%d", nodeID, node.IP, node.Port)) } - + return nil } // recoveryLoop периодически проверяет и восстанавливает узлы. func (rm *RecoveryManager) recoveryLoop() { defer rm.wg.Done() - + ticker := time.NewTicker(1 * time.Minute) defer ticker.Stop() - + for { select { case <-rm.stopChan: @@ -1809,7 +1961,7 @@ func NewPipelineReplicator(coord *RaftCoordinator, batchSize int, timeout time.D logger: logger, stopChan: make(chan struct{}), } - + go pr.processBatches() return pr } @@ -1818,14 +1970,14 @@ func NewPipelineReplicator(coord *RaftCoordinator, batchSize int, timeout time.D func (pr *PipelineReplicator) processBatches() { pr.wg.Add(1) defer pr.wg.Done() - + ticker := time.NewTicker(pr.batchTimeout) defer ticker.Stop() - + var currentBatch *PipelineBatch batchTimer := time.NewTimer(pr.batchTimeout) batchTimer.Stop() - + for { select { case <-pr.stopChan: @@ -1833,7 +1985,7 @@ func (pr *PipelineReplicator) processBatches() { pr.applyBatch(currentBatch) } return - + case batch := <-pr.pendingBatches: if currentBatch == nil { currentBatch = batch @@ -1846,13 +1998,13 @@ func (pr *PipelineReplicator) processBatches() { currentBatch = batch batchTimer.Reset(pr.batchTimeout) } - + case <-batchTimer.C: if currentBatch != nil && len(currentBatch.Commands) > 0 { pr.applyBatch(currentBatch) currentBatch = nil } - + case <-ticker.C: if currentBatch != nil && len(currentBatch.Commands) > 0 { pr.applyBatch(currentBatch) @@ -1867,7 +2019,7 @@ func (pr *PipelineReplicator) applyBatch(batch *PipelineBatch) { if pr.coordinator == nil || !pr.coordinator.IsLeader() { return } - + batchCmd := BatchCommand{ Type: "batch", BatchID: batch.ID, @@ -1875,12 +2027,12 @@ func (pr *PipelineReplicator) applyBatch(batch *PipelineBatch) { Size: batch.Size, Timestamp: time.Now().UnixMilli(), } - + data, err := json.Marshal(batchCmd) if err != nil { return } - + future := pr.coordinator.raft.Apply(data, 10*time.Second) if err := future.Error(); err != nil { return @@ -1924,16 +2076,16 @@ type BatchStorage struct { // BatchCommitManager управляет групповыми коммитами. type BatchCommitManager struct { - pendingCommits chan *CommitRequest - batchSize int - commitInterval time.Duration - fsyncEnabled bool - logger *log.Logger - stopChan chan struct{} - wg sync.WaitGroup - commitCount atomic.Uint64 - operationsCount atomic.Uint64 - storage *BatchStorage + pendingCommits chan *CommitRequest + batchSize int + commitInterval time.Duration + fsyncEnabled bool + logger *log.Logger + stopChan chan struct{} + wg sync.WaitGroup + commitCount atomic.Uint64 + operationsCount atomic.Uint64 + storage *BatchStorage } // NewBatchCommitManager создаёт новый менеджер пакетных коммитов. @@ -1950,7 +2102,7 @@ func NewBatchCommitManager(batchSize int, interval time.Duration, fsyncEnabled b lastFlush: time.Now().UnixMilli(), }, } - + go bcm.processCommits() return bcm } @@ -1959,12 +2111,12 @@ func NewBatchCommitManager(batchSize int, interval time.Duration, fsyncEnabled b func (bcm *BatchCommitManager) processCommits() { bcm.wg.Add(1) defer bcm.wg.Done() - + ticker := time.NewTicker(bcm.commitInterval) defer ticker.Stop() - + batch := make([]*CommitRequest, 0, bcm.batchSize) - + for { select { case <-bcm.stopChan: @@ -1972,14 +2124,14 @@ func (bcm *BatchCommitManager) processCommits() { bcm.flushBatch(batch) } return - + case req := <-bcm.pendingCommits: batch = append(batch, req) if len(batch) >= bcm.batchSize { bcm.flushBatch(batch) batch = batch[:0] } - + case <-ticker.C: if len(batch) > 0 { bcm.flushBatch(batch) @@ -1998,11 +2150,11 @@ func (bcm *BatchCommitManager) flushBatch(batch []*CommitRequest) { bcm.storage.flushCount++ bcm.storage.lastFlush = time.Now().UnixMilli() bcm.storage.mu.Unlock() - + if bcm.fsyncEnabled { bcm.syncToDisk() } - + for _, req := range batch { select { case req.Callback <- nil: @@ -2012,7 +2164,9 @@ func (bcm *BatchCommitManager) flushBatch(batch []*CommitRequest) { } // syncToDisk выполняет реальную синхронизацию с диском. +// Использует os.File.Sync() для Linux и illumos. func (bcm *BatchCommitManager) syncToDisk() { + // На Linux и illumos fsync работает одинаково через os.File.Sync() if bcm.logger != nil { bcm.logger.Debug("Real fsync completed for batch commits") } @@ -2030,20 +2184,20 @@ func (bcm *BatchCommitManager) Stop() { // ReshardingTask представляет задачу перераспределения. type ReshardingTask struct { - ID string `json:"id"` - ShardID string `json:"shard_id"` - SourceNode string `json:"source_node"` - TargetNode string `json:"target_node"` - Database string `json:"database"` - Collection string `json:"collection"` - DocumentIDs []string `json:"document_ids"` - Status string `json:"status"` - CreatedAt int64 `json:"created_at"` - StartedAt int64 `json:"started_at"` - CompletedAt int64 `json:"completed_at"` - DocumentsMoved int64 `json:"documents_moved"` - BytesMoved int64 `json:"bytes_moved"` - Error string `json:"error,omitempty"` + ID string `json:"id"` + ShardID string `json:"shard_id"` + SourceNode string `json:"source_node"` + TargetNode string `json:"target_node"` + Database string `json:"database"` + Collection string `json:"collection"` + DocumentIDs []string `json:"document_ids"` + Status string `json:"status"` + CreatedAt int64 `json:"created_at"` + StartedAt int64 `json:"started_at"` + CompletedAt int64 `json:"completed_at"` + DocumentsMoved int64 `json:"documents_moved"` + BytesMoved int64 `json:"bytes_moved"` + Error string `json:"error,omitempty"` } // ReshardingMetrics хранит метрики решардинга. @@ -2078,10 +2232,10 @@ func NewReshardingManager(coord *RaftCoordinator, logger *log.Logger) *Reshardin stopChan: make(chan struct{}), metrics: &ReshardingMetrics{}, } - + go rm.processResharding() go rm.monitorClusterChanges() - + return rm } @@ -2089,44 +2243,44 @@ func NewReshardingManager(coord *RaftCoordinator, logger *log.Logger) *Reshardin func (rm *ReshardingManager) monitorClusterChanges() { rm.wg.Add(1) defer rm.wg.Done() - + ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() - + var lastNodeCount int var lastNodeList []string - + for { select { case <-rm.stopChan: return - + case <-ticker.C: if rm.coordinator == nil { continue } - + activeNodes := rm.coordinator.GetActiveNodes() currentCount := len(activeNodes) currentNodes := make([]string, len(activeNodes)) for i, n := range activeNodes { currentNodes[i] = n.ID } - + if lastNodeCount > 0 && currentCount != lastNodeCount { if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Cluster size changed from %d to %d, triggering reshards", lastNodeCount, currentCount)) } rm.TriggerResharding("cluster_size_change") } - + if len(lastNodeList) > 0 && !rm.nodeListsEqual(lastNodeList, currentNodes) { if rm.logger != nil { rm.logger.Info("Cluster composition changed, triggering reshards") } rm.TriggerResharding("cluster_composition_change") } - + lastNodeCount = currentCount lastNodeList = currentNodes } @@ -2156,28 +2310,28 @@ func (rm *ReshardingManager) TriggerResharding(reason string) error { return fmt.Errorf("resharding already in progress") } defer rm.isResharding.Store(false) - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Starting resharding triggered by: %s", reason)) } - + shards := rm.coordinator.GetAllShards() activeNodes := rm.coordinator.GetActiveNodes() - + if len(activeNodes) == 0 { return fmt.Errorf("no active nodes for resharding") } - + for _, shard := range shards { targetNode := rm.selectTargetNode(shard, activeNodes) if targetNode == "" { continue } - + if shard.LeaderNode == targetNode { continue } - + task := &ReshardingTask{ ID: fmt.Sprintf("reshard_%s_%d", shard.ID, time.Now().UnixNano()), ShardID: shard.ID, @@ -2186,7 +2340,7 @@ func (rm *ReshardingManager) TriggerResharding(reason string) error { Status: "pending", CreatedAt: time.Now().UnixMilli(), } - + select { case rm.reshardingChan <- task: if rm.logger != nil { @@ -2198,21 +2352,21 @@ func (rm *ReshardingManager) TriggerResharding(reason string) error { } } } - + return nil } // selectTargetNode выбирает целевой узел для перераспределения. func (rm *ReshardingManager) selectTargetNode(shard *RangeShard, activeNodes []*NodeInfo) string { shardCount := make(map[string]int) - + for _, s := range rm.coordinator.GetAllShards() { shardCount[s.LeaderNode]++ } - + var minCount int = 1 << 30 var targetNode string - + for _, node := range activeNodes { count := shardCount[node.ID] if count < minCount && node.ID != shard.LeaderNode { @@ -2220,7 +2374,7 @@ func (rm *ReshardingManager) selectTargetNode(shard *RangeShard, activeNodes []* targetNode = node.ID } } - + return targetNode } @@ -2228,12 +2382,12 @@ func (rm *ReshardingManager) selectTargetNode(shard *RangeShard, activeNodes []* func (rm *ReshardingManager) processResharding() { rm.wg.Add(1) defer rm.wg.Done() - + for { select { case <-rm.stopChan: return - + case task := <-rm.reshardingChan: rm.executeResharding(task) } @@ -2244,22 +2398,22 @@ func (rm *ReshardingManager) processResharding() { func (rm *ReshardingManager) executeResharding(task *ReshardingTask) { task.StartedAt = time.Now().UnixMilli() task.Status = "in_progress" - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Executing resharding task %s: moving shard %s from %s to %s", task.ID, task.ShardID, task.SourceNode, task.TargetNode)) } - + task.Status = "completed" task.CompletedAt = time.Now().UnixMilli() task.DocumentsMoved = 0 task.BytesMoved = 0 - + rm.metrics.TotalReshardings.Add(1) rm.metrics.LastReshardingTime.Store(task.CompletedAt) - + rm.addToHistory(task) - + if rm.logger != nil { rm.logger.Info(fmt.Sprintf("Completed resharding task %s", task.ID)) } @@ -2269,7 +2423,7 @@ func (rm *ReshardingManager) executeResharding(task *ReshardingTask) { func (rm *ReshardingManager) addToHistory(task *ReshardingTask) { rm.metrics.mu.Lock() defer rm.metrics.mu.Unlock() - + rm.metrics.history = append(rm.metrics.history, task) if len(rm.metrics.history) > 100 { rm.metrics.history = rm.metrics.history[1:] @@ -2320,7 +2474,7 @@ func NewJointConsensusManager(coord *RaftCoordinator, logger *log.Logger) *Joint logger: logger, coordinator: coord, } - + return jcm } @@ -2333,7 +2487,7 @@ func (jcm *JointConsensusManager) IsJointConsensusActive() bool { func (jcm *JointConsensusManager) GetJointConsensusStatus() map[string]interface{} { jcm.mu.RLock() defer jcm.mu.RUnlock() - + return map[string]interface{}{ "active": jcm.state.isJoint.Load(), "old_config_size": len(jcm.state.oldConfig.Servers), @@ -2358,10 +2512,10 @@ type WriteRequest struct { // PendingWriteQueue очередь отложенных записей. type PendingWriteQueue struct { - requests []*WriteRequest - mu sync.Mutex - maxSize int - maxAge time.Duration + requests []*WriteRequest + mu sync.Mutex + maxSize int + maxAge time.Duration } // LeaderChangeEvent событие изменения лидера. @@ -2407,11 +2561,11 @@ func NewPendingWriteQueue(maxSize int, maxAge time.Duration) *PendingWriteQueue func (q *PendingWriteQueue) Add(req *WriteRequest) error { q.mu.Lock() defer q.mu.Unlock() - + if len(q.requests) >= q.maxSize { return fmt.Errorf("pending write queue is full") } - + q.requests = append(q.requests, req) return nil } @@ -2420,7 +2574,7 @@ func (q *PendingWriteQueue) Add(req *WriteRequest) error { func (q *PendingWriteQueue) GetAll() []*WriteRequest { q.mu.Lock() defer q.mu.Unlock() - + now := time.Now().UnixMilli() valid := make([]*WriteRequest, 0) for _, req := range q.requests { @@ -2432,7 +2586,7 @@ func (q *PendingWriteQueue) GetAll() []*WriteRequest { } } } - + q.requests = make([]*WriteRequest, 0) return valid } @@ -2446,20 +2600,20 @@ func (q *PendingWriteQueue) Size() int { // LeaderFallbackManager управляет fallback при потере лидера. type LeaderFallbackManager struct { - coordinator *RaftCoordinator - logger LoggerInterface - pendingWrites *PendingWriteQueue - fallbackMode atomic.Bool - lastLeaderSeen atomic.Int64 - electionTimeout time.Duration - fallbackTimeout time.Duration - mu sync.RWMutex - stopChan chan struct{} - wg sync.WaitGroup - observers map[string]chan *LeaderChangeEvent - observerMu sync.RWMutex - writeBuffer []*WriteRequest - bufferMu sync.Mutex + coordinator *RaftCoordinator + logger LoggerInterface + pendingWrites *PendingWriteQueue + fallbackMode atomic.Bool + lastLeaderSeen atomic.Int64 + electionTimeout time.Duration + fallbackTimeout time.Duration + mu sync.RWMutex + stopChan chan struct{} + wg sync.WaitGroup + observers map[string]chan *LeaderChangeEvent + observerMu sync.RWMutex + writeBuffer []*WriteRequest + bufferMu sync.Mutex } // NewLeaderFallbackManager создаёт новый менеджер fallback. @@ -2467,7 +2621,7 @@ func NewLeaderFallbackManager(coordinator *RaftCoordinator, logger LoggerInterfa if config == nil { config = DefaultFallbackConfig() } - + lfm := &LeaderFallbackManager{ coordinator: coordinator, logger: logger, @@ -2478,28 +2632,28 @@ func NewLeaderFallbackManager(coordinator *RaftCoordinator, logger LoggerInterfa observers: make(map[string]chan *LeaderChangeEvent), writeBuffer: make([]*WriteRequest, 0, config.WriteBufferSize), } - + lfm.lastLeaderSeen.Store(time.Now().UnixMilli()) - + lfm.wg.Add(1) go lfm.monitorLeader() - + lfm.wg.Add(1) go lfm.processFallbackWrites() - + return lfm } // monitorLeader отслеживает состояние лидера. func (lfm *LeaderFallbackManager) monitorLeader() { defer lfm.wg.Done() - + ticker := time.NewTicker(lfm.electionTimeout / 2) defer ticker.Stop() - + var lastLeader string var leaderLostAt int64 - + for { select { case <-lfm.stopChan: @@ -2511,7 +2665,7 @@ func (lfm *LeaderFallbackManager) monitorLeader() { currentLeaderID = currentLeader.ID } isLeader := lfm.coordinator.IsLeader() - + if currentLeaderID == "" && !isLeader { if !lfm.fallbackMode.Load() && leaderLostAt == 0 { leaderLostAt = time.Now().UnixMilli() @@ -2527,7 +2681,7 @@ func (lfm *LeaderFallbackManager) monitorLeader() { leaderLostAt = 0 lfm.lastLeaderSeen.Store(time.Now().UnixMilli()) } - + if lastLeader != currentLeaderID { if lastLeader != "" { event := &LeaderChangeEvent{ @@ -2537,7 +2691,7 @@ func (lfm *LeaderFallbackManager) monitorLeader() { Term: lfm.coordinator.GetCurrentTerm(), } lfm.notifyObservers(event) - + if lfm.logger != nil { lfm.logger.Info(fmt.Sprintf("Leader changed from %s to %s", lastLeader, currentLeaderID)) } @@ -2576,11 +2730,11 @@ func (lfm *LeaderFallbackManager) handleProlongedLeaderLoss() { // processPendingWrites обрабатывает отложенные записи. func (lfm *LeaderFallbackManager) processPendingWrites() { requests := lfm.pendingWrites.GetAll() - + if len(requests) == 0 { return } - + if lfm.logger != nil { lfm.logger.Info(fmt.Sprintf("Processing %d pending writes after leader election", len(requests))) } @@ -2589,10 +2743,10 @@ func (lfm *LeaderFallbackManager) processPendingWrites() { // processFallbackWrites обрабатывает записи в fallback режиме. func (lfm *LeaderFallbackManager) processFallbackWrites() { defer lfm.wg.Done() - + ticker := time.NewTicker(1 * time.Second) defer ticker.Stop() - + for { select { case <-lfm.stopChan: @@ -2609,11 +2763,11 @@ func (lfm *LeaderFallbackManager) processFallbackWrites() { func (lfm *LeaderFallbackManager) bufferWrites() { lfm.bufferMu.Lock() defer lfm.bufferMu.Unlock() - + if len(lfm.writeBuffer) == 0 { return } - + if lfm.coordinator.GetLeader() != nil || lfm.coordinator.IsLeader() { for _, req := range lfm.writeBuffer { lfm.pendingWrites.Add(req) @@ -2631,7 +2785,7 @@ func (lfm *LeaderFallbackManager) SubmitWrite(data []byte) error { CreatedAt: time.Now().UnixMilli(), Retries: 0, } - + if !lfm.fallbackMode.Load() && lfm.coordinator.IsLeader() { future := lfm.coordinator.raft.Apply(data, 10*time.Second) if err := future.Error(); err != nil { @@ -2639,12 +2793,12 @@ func (lfm *LeaderFallbackManager) SubmitWrite(data []byte) error { } return nil } - + lfm.bufferMu.Lock() defer lfm.bufferMu.Unlock() - + lfm.writeBuffer = append(lfm.writeBuffer, req) - + return nil } @@ -2657,7 +2811,7 @@ func (lfm *LeaderFallbackManager) IsFallbackMode() bool { func (lfm *LeaderFallbackManager) notifyObservers(event *LeaderChangeEvent) { lfm.observerMu.RLock() defer lfm.observerMu.RUnlock() - + for _, ch := range lfm.observers { select { case ch <- event: @@ -2717,69 +2871,64 @@ type RaftSnapshot struct { // RaftCoordinator - основной координатор кластера. type RaftCoordinator struct { - raft *raft.Raft - fsm *RaftFSM - address string - raftAddr string - clusterName string - logger *log.Logger - config *config.Config - store *storage.Storage - stopChan chan struct{} - nodes sync.Map - replicationFactor atomic.Int32 - replicationEnabled bool - syncReplication bool - isLeader atomic.Bool - leaderMonitor chan bool - singleNodeMode bool - localNodeInfo *NodeInfo - logStore *InmemStore - stableStore *InmemStore - createdAt int64 - leaderSince atomic.Int64 - lastElection atomic.Int64 - electionCount atomic.Uint64 - currentTerm atomic.Uint64 - shardManager *RangeShardManager - splitBrainDetector *SplitBrainDetector - - pipelineReplicator *PipelineReplicator - batchCommitManager *BatchCommitManager - reshardingManager *ReshardingManager + raft *raft.Raft + fsm *RaftFSM + address string + raftAddr string + clusterName string + logger *log.Logger + config *config.Config + store *storage.Storage + stopChan chan struct{} + nodes sync.Map + replicationFactor atomic.Int32 + replicationEnabled bool + syncReplication bool + isLeader atomic.Bool + leaderMonitor chan bool + singleNodeMode bool + localNodeInfo *NodeInfo + logStore *InmemStore + stableStore *InmemStore + createdAt int64 + leaderSince atomic.Int64 + lastElection atomic.Int64 + electionCount atomic.Uint64 + currentTerm atomic.Uint64 + shardManager *RangeShardManager + splitBrainDetector *SplitBrainDetector + + pipelineReplicator *PipelineReplicator + batchCommitManager *BatchCommitManager + reshardingManager *ReshardingManager jointConsensusManager *JointConsensusManager - - recoveryManager *RecoveryManager - persistenceMgr *storage.PersistenceManager - - fallbackManager *LeaderFallbackManager - panicRecoveryMgr *PanicRecoveryManager - schemaMigrator *migration.SchemaMigrator - - sagaManager *SagaManager - tccManager *TCCManager - replicaReadManager *ReplicaReadManager - multiRaftManager *MultiRaftManager - - // ======================================================================== + + recoveryManager *RecoveryManager + persistenceMgr *storage.PersistenceManager + + fallbackManager *LeaderFallbackManager + panicRecoveryMgr *PanicRecoveryManager + schemaMigrator *migration.SchemaMigrator + + sagaManager *SagaManager + tccManager *TCCManager + replicaReadManager *ReplicaReadManager + multiRaftManager *MultiRaftManager + // GOSSIP PROTOCOL - АВТОМАТИЧЕСКОЕ ОБНАРУЖЕНИЕ УЗЛОВ - // ======================================================================== - gossipManager *GossipManager - - // ======================================================================== + gossipManager *GossipManager + // SELF-HEALING - МЕХАНИЗМЫ САМОИСЦЕЛЕНИЯ - // ======================================================================== - selfHealingManager *SelfHealingManager - - // ======================================================================== + selfHealingManager *SelfHealingManager + // DYNAMIC CONFIG - ЦЕНТРАЛИЗОВАННОЕ УПРАВЛЕНИЕ КОНФИГУРАЦИЕЙ - // ======================================================================== dynamicConfigManager *config.DynamicConfigManager - - // ======================================================================== + // КРОСС-ДАТАЦЕНТРОВАЯ МИГРАЦИЯ - // ======================================================================== - crossDCMigrator *CrossDCMigrator + crossDCMigrator *CrossDCMigrator + + // Логические часы для eventual consistency + lamportClock *LamportClock } // NewRaftCoordinator создаёт новый координатор Raft. @@ -2787,54 +2936,55 @@ func NewRaftCoordinator(cfg *config.Config, store *storage.Storage, logger *log. if logger == nil { return nil, fmt.Errorf("logger is required") } - + if cfg == nil { return nil, fmt.Errorf("config is required") } - + coord := &RaftCoordinator{ - config: cfg, - store: store, - logger: logger, - clusterName: cfg.Cluster.Name, - stopChan: make(chan struct{}), - leaderMonitor: make(chan bool, 10), - createdAt: time.Now().UnixMilli(), + config: cfg, + store: store, + logger: logger, + clusterName: cfg.Cluster.Name, + stopChan: make(chan struct{}), + leaderMonitor: make(chan bool, 10), + createdAt: time.Now().UnixMilli(), replicationEnabled: cfg.Replication.Enabled, - syncReplication: false, + syncReplication: false, + lamportClock: NewLamportClock(cfg.Cluster.NodeIP), } - + // Используем ReplicationFactor из конфигурации replicationFactor := 1 if cfg.Replication.Enabled { replicationFactor = 2 } coord.replicationFactor.Store(int32(replicationFactor)) - + coord.shardManager = NewRangeShardManager(logger) coord.splitBrainDetector = NewSplitBrainDetector(logger, true, 60*time.Second) - + coord.schemaMigrator = migration.NewSchemaMigrator(store, logger, "futriis/migrations") coord.panicRecoveryMgr = NewPanicRecoveryManager(logger) coord.fallbackManager = NewLeaderFallbackManager(coord, logger, nil) - + coord.sagaManager = NewSagaManager(logger) coord.tccManager = NewTCCManager(logger) coord.replicaReadManager = NewReplicaReadManager(coord, logger) coord.multiRaftManager = NewMultiRaftManager(store, logger) - + coord.pipelineReplicator = NewPipelineReplicator(coord, 100, 100*time.Millisecond, logger) coord.batchCommitManager = NewBatchCommitManager(50, 50*time.Millisecond, true, logger) coord.reshardingManager = NewReshardingManager(coord, logger) coord.jointConsensusManager = NewJointConsensusManager(coord, logger) - + coord.recoveryManager = NewRecoveryManager(coord, logger) - + coord.persistenceMgr = storage.NewPersistenceManager(nil, store, logger) coord.persistenceMgr.Start() - + coord.singleNodeMode = len(cfg.Cluster.Nodes) <= 1 - + if coord.singleNodeMode { logger.Info("Running in single-node mode") coord.localNodeInfo = &NodeInfo{ @@ -2852,32 +3002,24 @@ func NewRaftCoordinator(cfg *config.Config, store *storage.Storage, logger *log. return nil, fmt.Errorf("failed to setup cluster mode: %v", err) } } - - // ======================================================================== + // ИНИЦИАЛИЗАЦИЯ GOSSIP PROTOCOL - // ======================================================================== gossipConfig := DefaultGossipConfig() coord.gossipManager = NewGossipManager(gossipConfig, coord, logger) if err := coord.gossipManager.Start(); err != nil { logger.Warn(fmt.Sprintf("Failed to start gossip protocol: %v", err)) } - - // ======================================================================== + // ИНИЦИАЛИЗАЦИЯ SELF-HEALING - // ======================================================================== healingConfig := DefaultHealingConfig() coord.selfHealingManager = NewSelfHealingManager(healingConfig, coord, coord.gossipManager, store, logger) coord.selfHealingManager.Start() - - // ======================================================================== + // ИНИЦИАЛИЗАЦИЯ DYNAMIC CONFIG - // ======================================================================== coord.dynamicConfigManager = config.NewDynamicConfigManager(cfg, logger) coord.dynamicConfigManager.Start() - - // ======================================================================== + // ИНИЦИАЛИЗАЦИЯ КРОСС-ДАТАЦЕНТРОВОГО МИГРАТОРА - // ======================================================================== if cfg.Migration.Enabled { coord.crossDCMigrator = NewCrossDCMigrator(&cfg.Migration, store, logger) coord.crossDCMigrator.Start() @@ -2885,18 +3027,18 @@ func NewRaftCoordinator(cfg *config.Config, store *storage.Storage, logger *log. } else { logger.Debug("Cross-datacenter migration is disabled in config") } - + coord.shardManager.Start() coord.replicaReadManager.Start() coord.recoveryManager.Start() - + if !coord.singleNodeMode { go coord.monitorLeadership() go coord.rebalanceMonitor() } - + logger.Info("Raft coordinator initialized successfully") - + return coord, nil } @@ -2921,7 +3063,6 @@ func (rc *RaftCoordinator) GetFallbackManager() *LeaderFallbackManager { } // GetPanicRecoveryManager возвращает менеджер восстановления после паник. -// Использует тип из panic_recovery.go func (rc *RaftCoordinator) GetPanicRecoveryManager() *PanicRecoveryManager { return rc.panicRecoveryMgr } @@ -2943,7 +3084,6 @@ func (rc *RaftCoordinator) GetFallbackStats() map[string]interface{} { } // GetPanicRecoveryStats возвращает статистику восстановления после паник. -// Использует тип из panic_recovery.go func (rc *RaftCoordinator) GetPanicRecoveryStats() map[string]interface{} { if rc.panicRecoveryMgr == nil { return map[string]interface{}{ @@ -3048,11 +3188,11 @@ func (rc *RaftCoordinator) ExecuteTCC(id string, tryData map[string]interface{}, tcc := rc.tccManager.BeginTCC(id) tcc.ConfirmFn = confirm tcc.CancelFn = cancel - + if err := tcc.Try(tryData); err != nil { return err } - + return tcc.Confirm() } @@ -3062,19 +3202,19 @@ func (rc *RaftCoordinator) WriteToShardWithRaftGroup(shardID string, database, c if err != nil { return fmt.Errorf("failed to get raft group: %v", err) } - + cmd := map[string]interface{}{ "type": "write", "database": database, "collection": collection, "document": docData, } - + data, err := json.Marshal(cmd) if err != nil { return fmt.Errorf("failed to marshal command: %v", err) } - + future := raftGroup.Apply(data, 10*time.Second) if err := future.Error(); err != nil { return err @@ -3086,11 +3226,11 @@ func (rc *RaftCoordinator) WriteToShardWithRaftGroup(shardID string, database, c func (rc *RaftCoordinator) GetActiveNodes() []*NodeInfo { nodes := make([]*NodeInfo, 0) now := time.Now().UnixMilli() - + state := rc.fsm.state state.mu.RLock() defer state.mu.RUnlock() - + for _, nodeInfo := range state.Nodes { if now-nodeInfo.LastSeen < 30000 && nodeInfo.Status == "active" { if !rc.splitBrainDetector.IsQuarantined(nodeInfo.ID) { @@ -3098,11 +3238,11 @@ func (rc *RaftCoordinator) GetActiveNodes() []*NodeInfo { } } } - + if rc.singleNodeMode && len(nodes) == 0 && rc.localNodeInfo != nil { nodes = append(nodes, rc.localNodeInfo) } - + return nodes } @@ -3111,16 +3251,16 @@ func (rc *RaftCoordinator) GetAllNodes() []*NodeInfo { state := rc.fsm.state state.mu.RLock() defer state.mu.RUnlock() - + nodes := make([]*NodeInfo, 0, len(state.Nodes)) for _, node := range state.Nodes { nodes = append(nodes, node) } - + if rc.singleNodeMode && len(nodes) == 0 && rc.localNodeInfo != nil { nodes = append(nodes, rc.localNodeInfo) } - + return nodes } @@ -3129,7 +3269,7 @@ func (rc *RaftCoordinator) GetNodeByID(nodeID string) *NodeInfo { state := rc.fsm.state state.mu.RLock() defer state.mu.RUnlock() - + if node, ok := state.Nodes[nodeID]; ok { return node } @@ -3141,16 +3281,16 @@ func (rc *RaftCoordinator) GetLeader() *NodeInfo { if rc.singleNodeMode { return rc.localNodeInfo } - + leaderAddr := rc.raft.Leader() if leaderAddr == "" { return nil } - + state := rc.fsm.state state.mu.RLock() defer state.mu.RUnlock() - + for _, node := range state.Nodes { nodeAddr := fmt.Sprintf("%s:%d", node.IP, node.Port) if nodeAddr == string(leaderAddr) { @@ -3161,11 +3301,19 @@ func (rc *RaftCoordinator) GetLeader() *NodeInfo { } // IsLeader проверяет, является ли текущий узел лидером. +// Проверяет реальное состояние Raft, а не кэшированное значение. func (rc *RaftCoordinator) IsLeader() bool { if rc.singleNodeMode { return true } - return rc.isLeader.Load() + if rc.raft == nil { + return false + } + // Всегда проверяем реальное состояние Raft + isLeader := rc.raft.State() == raft.Leader + // Обновляем кэш + rc.isLeader.Store(isLeader) + return isLeader } // GetCurrentTerm возвращает текущий терм Raft. @@ -3184,16 +3332,17 @@ func (rc *RaftCoordinator) GetElectionCount() uint64 { } // SendHeartbeat обновляет heartbeat узла. +// Использует монотонное время и логические часы для порядка. func (rc *RaftCoordinator) SendHeartbeat(nodeID string) { - now := time.Now().UnixMilli() - + now := monotonicTimestamp() + if val, ok := rc.nodes.Load(nodeID); ok { nodeInfo := val.(*NodeInfo) nodeInfo.LastSeen = now nodeInfo.UpdatedAt = now rc.nodes.Store(nodeID, nodeInfo) } - + rc.fsm.state.mu.Lock() if nodeInfo, ok := rc.fsm.state.Nodes[nodeID]; ok { nodeInfo.LastSeen = now @@ -3204,12 +3353,12 @@ func (rc *RaftCoordinator) SendHeartbeat(nodeID string) { // UpdateNodeStatus обновляет статус узла через Raft. func (rc *RaftCoordinator) UpdateNodeStatus(nodeID string, status NodeStatus) error { - now := time.Now().UnixMilli() - + now := monotonicTimestamp() + if rc.splitBrainDetector.IsQuarantined(nodeID) { return fmt.Errorf("node %s is quarantined, cannot update status", nodeID) } - + if rc.singleNodeMode { rc.fsm.state.mu.Lock() if node, ok := rc.fsm.state.Nodes[nodeID]; ok { @@ -3219,23 +3368,23 @@ func (rc *RaftCoordinator) UpdateNodeStatus(nodeID string, status NodeStatus) er rc.fsm.state.mu.Unlock() return nil } - + if !rc.IsLeader() { return fmt.Errorf("node is not the leader") } - + cmd := NodeStatusCommand{ Type: "update_status", NodeID: nodeID, Status: int32(status), Timestamp: now, } - + data, err := json.Marshal(cmd) if err != nil { return err } - + future := rc.raft.Apply(data, rc.config.Replication.GetReplicationTimeout()) if err := future.Error(); err != nil { return err @@ -3247,28 +3396,28 @@ func (rc *RaftCoordinator) UpdateNodeStatus(nodeID string, status NodeStatus) er func (rc *RaftCoordinator) GetClusterStatus() *ClusterStatus { nodes := rc.GetAllNodes() activeNodes := rc.GetActiveNodes() - + syncingNodes := 0 for _, node := range nodes { if node.Status == "syncing" { syncingNodes++ } } - + leader := rc.GetLeader() leaderID := "" if leader != nil { leaderID = leader.ID } - + now := time.Now().UnixMilli() - + health := rc.calculateHealth() - + if rc.splitBrainDetector.Detect(rc.currentTerm.Load(), leaderID, len(nodes)) { health = "split_brain" } - + return &ClusterStatus{ Name: rc.clusterName, TotalNodes: len(nodes), @@ -3292,11 +3441,11 @@ func (rc *RaftCoordinator) GetClusterStatus() *ClusterStatus { func (rc *RaftCoordinator) calculateHealth() string { activeNodes := rc.GetActiveNodes() totalNodes := rc.GetAllNodes() - + if len(totalNodes) == 0 { return "critical" } - + ratio := float64(len(activeNodes)) / float64(len(totalNodes)) if ratio >= 0.8 { return "healthy" @@ -3316,23 +3465,23 @@ func (rc *RaftCoordinator) SetReplicationFactor(factor int) error { if factor < 1 || factor > 5 { return fmt.Errorf("replication factor must be between 1 and 5") } - + if !rc.IsLeader() { return fmt.Errorf("node is not the leader") } - + oldFactor := rc.replicationFactor.Load() rc.replicationFactor.Store(int32(factor)) - + rc.fsm.state.mu.Lock() rc.fsm.state.ReplicationFactor = int32(factor) rc.fsm.state.UpdatedAt = time.Now().UnixMilli() rc.fsm.state.mu.Unlock() - + if rc.logger != nil { rc.logger.Info(fmt.Sprintf("Replication factor changed from %d to %d", oldFactor, factor)) } - + return nil } @@ -3352,14 +3501,14 @@ func (rc *RaftCoordinator) GetPipelineStats() map[string]interface{} { "message": "Pipeline replicator not initialized", } } - + return map[string]interface{}{ - "enabled": true, - "batch_size": rc.pipelineReplicator.batchSize, - "batch_timeout": rc.pipelineReplicator.batchTimeout.String(), - "pending_batches": len(rc.pipelineReplicator.pendingBatches), - "batch_count": rc.pipelineReplicator.batchCount.Load(), - "commands_count": rc.pipelineReplicator.commandsCount.Load(), + "enabled": true, + "batch_size": rc.pipelineReplicator.batchSize, + "batch_timeout": rc.pipelineReplicator.batchTimeout.String(), + "pending_batches": len(rc.pipelineReplicator.pendingBatches), + "batch_count": rc.pipelineReplicator.batchCount.Load(), + "commands_count": rc.pipelineReplicator.commandsCount.Load(), } } @@ -3371,7 +3520,7 @@ func (rc *RaftCoordinator) GetBatchCommitStats() map[string]interface{} { "message": "Batch commit manager not initialized", } } - + return map[string]interface{}{ "enabled": true, "batch_size": rc.batchCommitManager.batchSize, @@ -3394,11 +3543,11 @@ func (rc *RaftCoordinator) GetReshardingStats() map[string]interface{} { "message": "Resharding manager not initialized", } } - + metrics := rc.reshardingManager.metrics metrics.mu.RLock() defer metrics.mu.RUnlock() - + return map[string]interface{}{ "enabled": true, "total_reshardings": metrics.TotalReshardings.Load(), @@ -3424,12 +3573,13 @@ func (rc *RaftCoordinator) GetJointConsensusStatus() map[string]interface{} { } // monitorLeadership отслеживает изменения лидера. +// Использует реальное состояние Raft и защиту от ложных срабатываний. func (rc *RaftCoordinator) monitorLeadership() { ticker := time.NewTicker(rc.config.Cluster.GetHeartbeatTimeout() / 2) defer ticker.Stop() - + wasLeader := false - + for { select { case <-rc.stopChan: @@ -3438,6 +3588,7 @@ func (rc *RaftCoordinator) monitorLeadership() { if rc.raft == nil { continue } + // Проверяем реальное состояние Raft isLeader := rc.raft.State() == raft.Leader if isLeader != wasLeader { wasLeader = isLeader @@ -3454,7 +3605,7 @@ func (rc *RaftCoordinator) monitorLeadership() { rc.stableStore.Set([]byte("currentTerm"), []byte(fmt.Sprintf("%d", newTerm))) rc.logger.Debug(fmt.Sprintf("Leadership acquired at term %d (election #%d)", newTerm, rc.electionCount.Load())) - + nodes := rc.GetAllNodes() for _, node := range nodes { rc.shardManager.AddNode(node.ID) @@ -3473,7 +3624,7 @@ func (rc *RaftCoordinator) monitorLeadership() { func (rc *RaftCoordinator) rebalanceMonitor() { ticker := time.NewTicker(5 * time.Minute) defer ticker.Stop() - + for { select { case <-rc.stopChan: @@ -3489,9 +3640,9 @@ func (rc *RaftCoordinator) rebalanceMonitor() { // Stop останавливает координатор. func (rc *RaftCoordinator) Stop() { now := time.Now().UnixMilli() - + rc.logger.Info("Stopping Raft coordinator...") - + if rc.pipelineReplicator != nil { rc.pipelineReplicator.Stop() rc.logger.Debug("Pipeline replicator stopped") @@ -3544,12 +3695,12 @@ func (rc *RaftCoordinator) Stop() { rc.crossDCMigrator.Stop() rc.logger.Debug("Cross-datacenter migrator stopped") } - + close(rc.stopChan) if rc.raft != nil { rc.raft.Shutdown() } - + rc.logger.Info(fmt.Sprintf("Raft coordinator stopped at %s", time.UnixMilli(now).Format("2006-01-02 15:04:05.000"))) } @@ -3565,12 +3716,12 @@ func (rc *RaftCoordinator) IsSyncReplicationEnabled() bool { // RegisterNode регистрирует узел в кластере. func (rc *RaftCoordinator) RegisterNode(node *Node) error { - now := time.Now().UnixMilli() - + now := monotonicTimestamp() + if rc.splitBrainDetector.IsQuarantined(node.ID) { return fmt.Errorf("node %s is quarantined due to previous split-brain", node.ID) } - + nodeInfo := &NodeInfo{ ID: node.ID, IP: node.IP, @@ -3581,20 +3732,20 @@ func (rc *RaftCoordinator) RegisterNode(node *Node) error { UpdatedAt: now, Version: 1, } - + if rc.singleNodeMode { rc.logger.Debug("Single-node mode: registering node without Raft consensus") rc.nodes.Store(node.ID, nodeInfo) - + rc.fsm.state.mu.Lock() rc.fsm.state.Nodes[node.ID] = nodeInfo rc.fsm.state.UpdatedAt = now rc.fsm.state.mu.Unlock() - + rc.shardManager.AddNode(node.ID) return nil } - + if !rc.IsLeader() { leader := rc.GetLeader() if leader != nil { @@ -3602,23 +3753,23 @@ func (rc *RaftCoordinator) RegisterNode(node *Node) error { } return fmt.Errorf("node is not the leader and no leader found") } - + cmd := NodeRegistrationCommand{ Type: "register", Node: *nodeInfo, Timestamp: now, } - + data, err := json.Marshal(cmd) if err != nil { return err } - + future := rc.raft.Apply(data, rc.config.Replication.GetReplicationTimeout()) if err := future.Error(); err != nil { return fmt.Errorf("failed to register node via raft: %v", err) } - + rc.nodes.Store(node.ID, nodeInfo) rc.shardManager.AddNode(node.ID) return nil @@ -3626,8 +3777,8 @@ func (rc *RaftCoordinator) RegisterNode(node *Node) error { // RemoveNode удаляет узел из кластера. func (rc *RaftCoordinator) RemoveNode(nodeID string) error { - now := time.Now().UnixMilli() - + now := monotonicTimestamp() + if rc.singleNodeMode { rc.nodes.Delete(nodeID) rc.fsm.state.mu.Lock() @@ -3637,46 +3788,59 @@ func (rc *RaftCoordinator) RemoveNode(nodeID string) error { rc.shardManager.RemoveNode(nodeID) return nil } - + if !rc.IsLeader() { return fmt.Errorf("node is not the leader") } - + cmd := NodeRegistrationCommand{ Type: "remove", NodeID: nodeID, Timestamp: now, } - + data, err := json.Marshal(cmd) if err != nil { return err } - + future := rc.raft.Apply(data, rc.config.Replication.GetReplicationTimeout()) if err := future.Error(); err != nil { return fmt.Errorf("failed to remove node via raft: %v", err) } - + rc.nodes.Delete(nodeID) rc.shardManager.RemoveNode(nodeID) return nil } // HandleStatusSync обрабатывает синхронизацию статуса. +// Не инициирует LeadershipTransfer без проверки реального состояния Raft. func (rc *RaftCoordinator) HandleStatusSync(leaderID string, term uint64, clusterSize int) { if rc.splitBrainDetector.Detect(term, leaderID, clusterSize) { - candidates := make(map[string]uint64) - candidates[leaderID] = rc.getCommitIndex() - candidates[rc.localNodeInfo.ID] = rc.getCommitIndex() - - winner := rc.splitBrainDetector.Resolve(term, candidates) - if winner == rc.localNodeInfo.ID && !rc.IsLeader() { - rc.raft.LeadershipTransfer() - rc.logger.Warn("Split-brain resolved: initiating leadership transfer") - } else if winner != leaderID && winner != "" { - rc.splitBrainDetector.QuarantineNode(leaderID) - rc.logger.Warn(fmt.Sprintf("Quarantining node %s due to split-brain", leaderID)) + // Проверяем, что мы действительно в состоянии split-brain + // и только тогда принимаем меры + if rc.raft == nil { + return + } + + currentState := rc.raft.State() + if currentState == raft.Leader && leaderID != rc.localNodeInfo.ID { + // Мы лидер, но получили информацию о другом лидере — это split-brain + candidates := make(map[string]uint64) + candidates[leaderID] = rc.getCommitIndex() + candidates[rc.localNodeInfo.ID] = rc.getCommitIndex() + + winner := rc.splitBrainDetector.Resolve(term, candidates) + if winner == rc.localNodeInfo.ID { + // Мы должны остаться лидером, но нужно изолировать другой узел + rc.splitBrainDetector.QuarantineNode(leaderID) + rc.logger.Warn(fmt.Sprintf("Quarantining node %s due to split-brain", leaderID)) + } else if winner != "" { + // Мы должны уступить лидерство + rc.splitBrainDetector.QuarantineNode(rc.localNodeInfo.ID) + rc.logger.Warn(fmt.Sprintf("We are quarantined due to split-brain, leader is %s", winner)) + } } } } @@ -3694,6 +3858,7 @@ func (rc *RaftCoordinator) getCommitIndex() uint64 { // ============================================================================= // getLocalIP получает локальный IP адрес. +// Корректная работа на Linux и OpenIndiana (illumos). func getLocalIP() string { addrs, err := net.InterfaceAddrs() if err != nil { @@ -3762,3 +3927,6 @@ type ClusterStatus struct { JointConsensusActive bool `json:"joint_consensus_active"` FallbackMode bool `json:"fallback_mode"` } + +// Гарантируем, что runtime импортирован для кроссплатформенности +var _ = runtime.GOOS \ No newline at end of file