Update internal/cluster/raft_coordinator.go
This commit is contained in:
1 parent
384553da25
commit
0d7bad66bb
1 file changed
+243
-75
@@ -16,14 +16,17 @@ 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"
|
||||
@@ -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)
|
||||
// =============================================================================
|
||||
@@ -545,10 +630,13 @@ func (rsm *RangeShardManager) Rebalance() error {
|
||||
// =============================================================================
|
||||
|
||||
// InmemStore реализует хранилище для Raft в памяти
|
||||
// Корректная обработка конфликтов логов при смене лидера.
|
||||
type InmemStore struct {
|
||||
data map[string][]byte
|
||||
mu sync.RWMutex
|
||||
path string
|
||||
firstIndex uint64
|
||||
lastIndex uint64
|
||||
}
|
||||
|
||||
// NewInmemStore создаёт новое in-memory хранилище
|
||||
@@ -556,6 +644,8 @@ func NewInmemStore(path string) *InmemStore {
|
||||
return &InmemStore{
|
||||
data: make(map[string][]byte),
|
||||
path: path,
|
||||
firstIndex: 1,
|
||||
lastIndex: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -583,41 +673,17 @@ 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
|
||||
if s.lastIndex == 0 {
|
||||
return 0, nil
|
||||
}
|
||||
// Парсим индекс из ключа
|
||||
var idx uint64
|
||||
if _, err := fmt.Sscanf(k, "%d", &idx); err == nil {
|
||||
if first == 0 || idx < first {
|
||||
first = idx
|
||||
}
|
||||
}
|
||||
}
|
||||
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)
|
||||
@@ -625,7 +691,7 @@ 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
|
||||
@@ -635,45 +701,96 @@ func (s *InmemStore) GetLog(index uint64, log *raft.Log) error {
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
@@ -778,6 +895,7 @@ func NewMultiRaftManager(storage *storage.Storage, logger LoggerInterface) *Mult
|
||||
}
|
||||
|
||||
// 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
|
||||
@@ -805,7 +923,6 @@ func (mrm *MultiRaftManager) GetOrCreateRaftGroup(shardID string, nodes []string
|
||||
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)
|
||||
@@ -819,13 +936,13 @@ func (mrm *MultiRaftManager) GetOrCreateRaftGroup(shardID string, nodes []string
|
||||
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)
|
||||
@@ -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
|
||||
@@ -1048,7 +1170,16 @@ 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()
|
||||
|
||||
@@ -1087,9 +1218,15 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error {
|
||||
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 {
|
||||
@@ -1108,6 +1245,7 @@ func (sm *SagaManager) Execute(saga *SagaTransaction) error {
|
||||
}
|
||||
|
||||
step.Status = "completed"
|
||||
saga.executed[i] = true // Отмечаем выполненный шаг
|
||||
saga.UpdatedAt = time.Now().UnixMilli()
|
||||
}
|
||||
|
||||
@@ -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,10 +1573,12 @@ 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
|
||||
@@ -1445,8 +1587,12 @@ func (sbd *SplitBrainDetector) Detect(term uint64, leaderID string, nodesCount i
|
||||
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))
|
||||
@@ -1455,14 +1601,20 @@ func (sbd *SplitBrainDetector) Detect(term uint64, leaderID string, nodesCount i
|
||||
}
|
||||
}
|
||||
|
||||
// Обновляем только если лидер не пустой
|
||||
if leaderID != "" {
|
||||
sbd.knownLeaders[term] = leaderID
|
||||
}
|
||||
|
||||
// Очищаем старые термы (защита от утечки памяти)
|
||||
for t := range sbd.knownLeaders {
|
||||
if t+10 < term {
|
||||
delete(sbd.knownLeaders, t)
|
||||
}
|
||||
}
|
||||
|
||||
_ = now // используется для будущих проверок clock skew
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
@@ -2761,25 +2915,20 @@ type RaftCoordinator struct {
|
||||
replicaReadManager *ReplicaReadManager
|
||||
multiRaftManager *MultiRaftManager
|
||||
|
||||
// ========================================================================
|
||||
// GOSSIP PROTOCOL - АВТОМАТИЧЕСКОЕ ОБНАРУЖЕНИЕ УЗЛОВ
|
||||
// ========================================================================
|
||||
gossipManager *GossipManager
|
||||
|
||||
// ========================================================================
|
||||
// SELF-HEALING - МЕХАНИЗМЫ САМОИСЦЕЛЕНИЯ
|
||||
// ========================================================================
|
||||
selfHealingManager *SelfHealingManager
|
||||
|
||||
// ========================================================================
|
||||
// DYNAMIC CONFIG - ЦЕНТРАЛИЗОВАННОЕ УПРАВЛЕНИЕ КОНФИГУРАЦИЕЙ
|
||||
// ========================================================================
|
||||
dynamicConfigManager *config.DynamicConfigManager
|
||||
|
||||
// ========================================================================
|
||||
// КРОСС-ДАТАЦЕНТРОВАЯ МИГРАЦИЯ
|
||||
// ========================================================================
|
||||
crossDCMigrator *CrossDCMigrator
|
||||
|
||||
// Логические часы для eventual consistency
|
||||
lamportClock *LamportClock
|
||||
}
|
||||
|
||||
// NewRaftCoordinator создаёт новый координатор Raft.
|
||||
@@ -2802,6 +2951,7 @@ func NewRaftCoordinator(cfg *config.Config, store *storage.Storage, logger *log.
|
||||
createdAt: time.Now().UnixMilli(),
|
||||
replicationEnabled: cfg.Replication.Enabled,
|
||||
syncReplication: false,
|
||||
lamportClock: NewLamportClock(cfg.Cluster.NodeIP),
|
||||
}
|
||||
|
||||
// Используем ReplicationFactor из конфигурации
|
||||
@@ -2853,31 +3003,23 @@ func NewRaftCoordinator(cfg *config.Config, store *storage.Storage, logger *log.
|
||||
}
|
||||
}
|
||||
|
||||
// ========================================================================
|
||||
// ИНИЦИАЛИЗАЦИЯ 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()
|
||||
@@ -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{}{
|
||||
@@ -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,8 +3332,9 @@ 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)
|
||||
@@ -3204,7 +3353,7 @@ 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)
|
||||
@@ -3424,6 +3573,7 @@ func (rc *RaftCoordinator) GetJointConsensusStatus() map[string]interface{} {
|
||||
}
|
||||
|
||||
// monitorLeadership отслеживает изменения лидера.
|
||||
// Использует реальное состояние Raft и защиту от ложных срабатываний.
|
||||
func (rc *RaftCoordinator) monitorLeadership() {
|
||||
ticker := time.NewTicker(rc.config.Cluster.GetHeartbeatTimeout() / 2)
|
||||
defer ticker.Stop()
|
||||
@@ -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
|
||||
@@ -3565,7 +3716,7 @@ 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)
|
||||
@@ -3626,7 +3777,7 @@ 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)
|
||||
@@ -3664,19 +3815,32 @@ func (rc *RaftCoordinator) RemoveNode(nodeID string) error {
|
||||
}
|
||||
|
||||
// HandleStatusSync обрабатывает синхронизацию статуса.
|
||||
// Не инициирует LeadershipTransfer без проверки реального состояния Raft.
|
||||
func (rc *RaftCoordinator) HandleStatusSync(leaderID string, term uint64, clusterSize int) {
|
||||
if rc.splitBrainDetector.Detect(term, leaderID, clusterSize) {
|
||||
// Проверяем, что мы действительно в состоянии 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.IsLeader() {
|
||||
rc.raft.LeadershipTransfer()
|
||||
rc.logger.Warn("Split-brain resolved: initiating leadership transfer")
|
||||
} else if winner != leaderID && winner != "" {
|
||||
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
|
||||
Reference in new issue
Block a user