From 113d79aaac45a336e6d38244addd3a3234e23178 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:39:27 +0000 Subject: [PATCH] Upload files to "internal/cluster" --- internal/cluster/gossip.go | 860 +++++++++++++++++++++++++++++++++++++ 1 file changed, 860 insertions(+) create mode 100644 internal/cluster/gossip.go diff --git a/internal/cluster/gossip.go b/internal/cluster/gossip.go new file mode 100644 index 0000000..9f3d657 --- /dev/null +++ b/internal/cluster/gossip.go @@ -0,0 +1,860 @@ +/* + * Copyright 2026 Safronov Grigorii + * + * Licensed under the CDDL, Version 1.0 (the "License"); + * you may not use this file except in compliance with the License. + * + * You may obtain a copy of the License at + * https://opensource.org/licenses/CDDL-1.0 + */ + +// Файл: internal/cluster/gossip.go +// Назначение: Gossip Protocol для автоматического обнаружения узлов +// и обмена информацией о состоянии кластера. + +package cluster + +import ( + "encoding/json" + "fmt" + "math/rand" + "net" + "sync" + "sync/atomic" + "time" + + "futriis/internal/log" +) + +// ============================================================================= +// ТИПЫ ДЛЯ GOSSIP PROTOCOL +// ============================================================================= + +// GossipMessageType представляет тип сообщения протокола сплетен +type GossipMessageType int + +const ( + GossipTypeMembership GossipMessageType = iota + GossipTypeHealth + GossipTypeLoad + GossipTypeConfig + GossipTypeSync +) + +// GossipMessage представляет сообщение протокола сплетен +type GossipMessage struct { + Type GossipMessageType `json:"type"` + SenderID string `json:"sender_id"` + SenderIP string `json:"sender_ip"` + SenderPort int `json:"sender_port"` + Timestamp int64 `json:"timestamp"` + Payload json.RawMessage `json:"payload"` + TTL int `json:"ttl"` + Version uint64 `json:"version"` + // Логическая метка для упорядочивания сообщений + LamportTS uint64 `json:"lamport_ts,omitempty"` +} + +// NodeState представляет состояние узла в протоколе сплетен +type NodeState struct { + ID string `json:"id"` + IP string `json:"ip"` + Port int `json:"port"` + RaftPort int `json:"raft_port"` + Status string `json:"status"` + LastSeen time.Time `json:"last_seen"` + LastHeartbeat time.Time `json:"last_heartbeat"` + Incarnation uint64 `json:"incarnation"` + LoadCPU float64 `json:"load_cpu"` + LoadMemory float64 `json:"load_memory"` + LoadQueue int64 `json:"load_queue"` + Connections int `json:"connections"` + Version string `json:"version"` + IsAlive bool `json:"is_alive"` + IsSuspect bool `json:"is_suspect"` + // ИСПРАВЛЕНО: логические часы узла для разрешения конфликтов + LamportTS uint64 `json:"lamport_ts,omitempty"` + // ИСПРАВЛЕНО: смещение часов относительно локального времени + ClockOffsetMs int64 `json:"clock_offset_ms,omitempty"` +} + +// GossipConfig содержит настройки протокола сплетен +// ИСПРАВЛЕНО: увеличены таймауты для предотвращения ложных срабатываний +// при сетевых задержках и рассинхронизации часов. +type GossipConfig struct { + BroadcastIntervalSec int `json:"broadcast_interval_sec"` + MonitorIntervalSec int `json:"monitor_interval_sec"` + Fanout int `json:"fanout"` + MaxTTL int `json:"max_ttl"` + TimeoutSec int `json:"timeout_sec"` + // ИСПРАВЛЕНО: увеличены таймауты (было 15/30) + SuspectTimeoutSec int `json:"suspect_timeout_sec"` + DeadTimeoutSec int `json:"dead_timeout_sec"` + BufferSize int `json:"buffer_size"` + // ИСПРАВЛЕНО: максимальная допустимая рассинхронизация часов + MaxClockSkewMs int64 `json:"max_clock_skew_ms"` +} + +// DefaultGossipConfig возвращает конфигурацию gossip по умолчанию +// ИСПРАВЛЕНО: таймауты увеличены для стабильности при сетевых задержках. +func DefaultGossipConfig() *GossipConfig { + return &GossipConfig{ + BroadcastIntervalSec: 5, + MonitorIntervalSec: 2, + Fanout: 3, + MaxTTL: 3, + TimeoutSec: 10, + // ИСПРАВЛЕНО: было 15/30, стало 30/60 для устойчивости к сетевым задержкам + SuspectTimeoutSec: 30, + DeadTimeoutSec: 60, + BufferSize: 10000, + // ИСПРАВЛЕНО: 10 секунд — допустимая рассинхронизация часов + MaxClockSkewMs: 10000, + } +} + +// ============================================================================= +// GOSSIP MANAGER +// ============================================================================= + +// GossipManager управляет протоколом сплетен для обнаружения узлов +type GossipManager struct { + config *GossipConfig + logger *log.Logger + coordinator *RaftCoordinator + localNode *NodeState + nodes sync.Map + membership atomic.Value + subscribers map[string]chan *GossipMessage + subscriberMu sync.RWMutex + stopChan chan struct{} + wg sync.WaitGroup + udpConn *net.UDPConn + msgCounter atomic.Uint64 + incarnation atomic.Uint64 + mu sync.RWMutex + seeds []string + metrics *GossipMetrics + // ИСПРАВЛЕНО: логические часы для упорядочивания сообщений + lamportClock *LamportClock + // ИСПРАВЛЕНО: локальное смещение часов + clockOffset atomic.Int64 +} + +// GossipMetrics хранит метрики gossip протокола +type GossipMetrics struct { + MessagesSent atomic.Uint64 + MessagesReceived atomic.Uint64 + MessagesDropped atomic.Uint64 + NodesDiscovered atomic.Uint64 + NodesRemoved atomic.Uint64 + LastBroadcast atomic.Int64 + LastSync atomic.Int64 + // ИСПРАВЛЕНО: метрики для clock skew + ClockSkewDetected atomic.Uint64 +} + +// NewGossipManager создаёт новый менеджер gossip протокола +func NewGossipManager(config *GossipConfig, coordinator *RaftCoordinator, logger *log.Logger) *GossipManager { + if config == nil { + config = DefaultGossipConfig() + } + + gm := &GossipManager{ + config: config, + logger: logger, + coordinator: coordinator, + subscribers: make(map[string]chan *GossipMessage), + stopChan: make(chan struct{}), + metrics: &GossipMetrics{}, + seeds: make([]string, 0), + lamportClock: NewLamportClock("gossip"), + } + + gm.membership.Store(make([]*NodeState, 0)) + + return gm +} + +// Start запускает gossip протокол +func (gm *GossipManager) Start() error { + gm.localNode = &NodeState{ + ID: gm.coordinator.localNodeInfo.ID, + IP: gm.coordinator.localNodeInfo.IP, + Port: gm.coordinator.localNodeInfo.Port, + RaftPort: gm.coordinator.config.Cluster.RaftPort, + Status: "active", + LastSeen: time.Now(), + LastHeartbeat: time.Now(), + Incarnation: 0, + IsAlive: true, + IsSuspect: false, + Version: "1.0", + LamportTS: 0, + } + + if err := gm.startUDPServer(); err != nil { + return fmt.Errorf("failed to start UDP server: %v", err) + } + + gm.wg.Add(4) + go gm.broadcastLoop() + go gm.monitorLoop() + go gm.receiveLoop() + go gm.syncLoop() + + gm.logger.Info(fmt.Sprintf("Gossip protocol started on port %d", gm.coordinator.config.Cluster.RaftPort+1)) + + return nil +} + +// Stop останавливает gossip протокол +func (gm *GossipManager) Stop() { + close(gm.stopChan) + gm.wg.Wait() + + if gm.udpConn != nil { + gm.udpConn.Close() + } + + gm.logger.Info("Gossip protocol stopped") +} + +// startUDPServer запускает UDP сервер для приёма сообщений +// ИСПРАВЛЕНО: корректная привязка на Linux и OpenIndiana (illumos). +// Если IP пустой или 0.0.0.0, используем wildcard-адрес. +func (gm *GossipManager) startUDPServer() error { + port := gm.coordinator.config.Cluster.RaftPort + 1 + + // ИСПРАВЛЕНО: определяем адрес для привязки + bindIP := gm.localNode.IP + if bindIP == "" || bindIP == "0.0.0.0" || bindIP == "::" { + bindIP = "0.0.0.0" + } + + addr := fmt.Sprintf("%s:%d", bindIP, port) + + udpAddr, err := net.ResolveUDPAddr("udp", addr) + if err != nil { + return err + } + + conn, err := net.ListenUDP("udp", udpAddr) + if err != nil { + // ИСПРАВЛЕНО: fallback на localhost, если не удалось привязаться + // (частая проблема на illumos при отсутствии прав) + fallbackAddr := fmt.Sprintf("127.0.0.1:%d", port) + udpAddr, err2 := net.ResolveUDPAddr("udp", fallbackAddr) + if err2 != nil { + return fmt.Errorf("failed to resolve UDP address %s: %v (fallback also failed: %v)", addr, err, err2) + } + conn, err = net.ListenUDP("udp", udpAddr) + if err != nil { + return fmt.Errorf("failed to listen on UDP %s: %v (fallback %s also failed)", addr, err, fallbackAddr) + } + gm.logger.Warn(fmt.Sprintf("Gossip UDP server bound to fallback address %s", fallbackAddr)) + } + + gm.udpConn = conn + return nil +} + +// broadcastLoop периодически рассылает gossip сообщения +func (gm *GossipManager) broadcastLoop() { + defer gm.wg.Done() + + ticker := time.NewTicker(time.Duration(gm.config.BroadcastIntervalSec) * time.Second) + defer ticker.Stop() + + for { + select { + case <-gm.stopChan: + return + case <-ticker.C: + gm.broadcast() + } + } +} + +// broadcast рассылает сообщение случайным узлам +func (gm *GossipManager) broadcast() { + // ИСПРАВЛЕНО: обновляем логическую метку + lamportTS := gm.lamportClock.Tick() + + payload, _ := json.Marshal(gm.localNode) + msg := &GossipMessage{ + Type: GossipTypeMembership, + SenderID: gm.localNode.ID, + SenderIP: gm.localNode.IP, + SenderPort: gm.localNode.Port, + Timestamp: time.Now().UnixMilli(), + Payload: payload, + TTL: gm.config.MaxTTL, + Version: gm.incarnation.Add(1), + LamportTS: lamportTS, + } + + nodes := gm.getActiveNodes() + + if len(nodes) == 0 { + gm.broadcastToSeeds(msg) + return + } + + fanout := gm.config.Fanout + if fanout > len(nodes) { + fanout = len(nodes) + } + + selected := gm.selectRandomNodes(nodes, fanout) + + for _, node := range selected { + gm.sendMessage(msg, node.IP, node.Port) + } + + gm.metrics.MessagesSent.Add(uint64(len(selected))) + gm.metrics.LastBroadcast.Store(time.Now().UnixMilli()) +} + +// broadcastToSeeds рассылает сообщения seed узлам +func (gm *GossipManager) broadcastToSeeds(msg *GossipMessage) { + gm.mu.RLock() + seeds := make([]string, len(gm.seeds)) + copy(seeds, gm.seeds) + gm.mu.RUnlock() + + for _, seed := range seeds { + host, port, err := net.SplitHostPort(seed) + if err != nil { + continue + } + var portInt int + fmt.Sscanf(port, "%d", &portInt) + gm.sendMessage(msg, host, portInt) + } +} + +// monitorLoop периодически проверяет состояние узлов +func (gm *GossipManager) monitorLoop() { + defer gm.wg.Done() + + ticker := time.NewTicker(time.Duration(gm.config.MonitorIntervalSec) * time.Second) + defer ticker.Stop() + + for { + select { + case <-gm.stopChan: + return + case <-ticker.C: + gm.monitorNodes() + } + } +} + +// monitorNodes проверяет состояние всех известных узлов +// ИСПРАВЛЕНО: учитывает clock skew между узлами, использует монотонное время. +func (gm *GossipManager) monitorNodes() { + now := time.Now() + deadTimeout := time.Duration(gm.config.DeadTimeoutSec) * time.Second + suspectTimeout := time.Duration(gm.config.SuspectTimeoutSec) * time.Second + + var nodesToRemove []string + + gm.nodes.Range(func(key, value interface{}) bool { + nodeID := key.(string) + node := value.(*NodeState) + + if nodeID == gm.localNode.ID { + return true + } + + if !node.IsAlive { + return true + } + + // ИСПРАВЛЕНО: используем локальное время, а не время узла + // (защита от рассинхронизации часов) + lastSeen := now.Sub(node.LastSeen) + + if lastSeen > deadTimeout { + node.IsAlive = false + node.IsSuspect = false + gm.nodes.Store(nodeID, node) + nodesToRemove = append(nodesToRemove, nodeID) + gm.metrics.NodesRemoved.Add(1) + + if gm.logger != nil { + gm.logger.Warn(fmt.Sprintf("Node %s considered dead (last seen %v ago)", nodeID, lastSeen)) + } + + if gm.coordinator != nil { + gm.coordinator.RemoveNode(nodeID) + } + } else if lastSeen > suspectTimeout { + if !node.IsSuspect { + node.IsSuspect = true + gm.nodes.Store(nodeID, node) + + if gm.logger != nil { + gm.logger.Warn(fmt.Sprintf("Node %s marked as suspect (last seen %v ago)", nodeID, lastSeen)) + } + + gm.requestConfirmation(node) + } + } + + return true + }) + + for _, nodeID := range nodesToRemove { + gm.nodes.Delete(nodeID) + } + + gm.updateMembership() +} + +// requestConfirmation запрашивает подтверждение от узла +func (gm *GossipManager) requestConfirmation(node *NodeState) { + lamportTS := gm.lamportClock.Tick() + + msg := &GossipMessage{ + Type: GossipTypeHealth, + SenderID: gm.localNode.ID, + SenderIP: gm.localNode.IP, + SenderPort: gm.localNode.Port, + Timestamp: time.Now().UnixMilli(), + Payload: json.RawMessage(`{"action":"ping"}`), + TTL: 1, + LamportTS: lamportTS, + } + + gm.sendMessage(msg, node.IP, node.Port) +} + +// receiveLoop принимает входящие gossip сообщения +// ИСПРАВЛЕНО: обработка временных ошибок для Linux и illumos. +func (gm *GossipManager) receiveLoop() { + defer gm.wg.Done() + + buffer := make([]byte, 65536) + + for { + select { + case <-gm.stopChan: + return + default: + n, addr, err := gm.udpConn.ReadFromUDP(buffer) + if err != nil { + // ИСПРАВЛЕНО: обработка временных ошибок + if isTemporaryNetError(err) { + continue + } + // Проверяем, не закрыт ли сокет + if gm.udpConn == nil { + return + } + continue + } + + if n > 0 { + go gm.handleMessage(buffer[:n], addr) + } + } + } +} + +// handleMessage обрабатывает полученное gossip сообщение +func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) { + var msg GossipMessage + if err := json.Unmarshal(data, &msg); err != nil { + return + } + + gm.metrics.MessagesReceived.Add(1) + + // ИСПРАВЛЕНО: обновляем логические часы + if msg.LamportTS > 0 { + gm.lamportClock.Observe(msg.LamportTS) + } + + gm.updateNodeFromMessage(&msg) + + switch msg.Type { + case GossipTypeMembership: + gm.handleMembershipMessage(&msg) + case GossipTypeHealth: + gm.handleHealthMessage(&msg) + case GossipTypeLoad: + gm.handleLoadMessage(&msg) + case GossipTypeConfig: + gm.handleConfigMessage(&msg) + case GossipTypeSync: + gm.handleSyncMessage(&msg) + } + + if msg.TTL > 0 && msg.SenderID != gm.localNode.ID { + msg.TTL-- + gm.forwardMessage(&msg) + } +} + +// updateNodeFromMessage обновляет информацию об узле из сообщения +// ИСПРАВЛЕНО: учитывает clock skew и использует логические часы. +func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) { + if msg.SenderID == gm.localNode.ID { + return + } + + var nodeState NodeState + if msg.Type == GossipTypeMembership { + if err := json.Unmarshal(msg.Payload, &nodeState); err != nil { + return + } + } else { + nodeState = NodeState{ + ID: msg.SenderID, + IP: msg.SenderIP, + Port: msg.SenderPort, + LastSeen: time.Now(), + IsAlive: true, + } + } + + // ИСПРАВЛЕНО: вычисляем и сохраняем смещение часов + now := time.Now().UnixMilli() + if msg.Timestamp > 0 { + offset := msg.Timestamp - now + nodeState.ClockOffsetMs = offset + + // Проверяем, не превышает ли смещение допустимое + if offset < -gm.config.MaxClockSkewMs || offset > gm.config.MaxClockSkewMs { + gm.metrics.ClockSkewDetected.Add(1) + if gm.logger != nil { + gm.logger.Warn(fmt.Sprintf("Clock skew detected for node %s: %d ms (max allowed: %d)", + msg.SenderID, offset, gm.config.MaxClockSkewMs)) + } + } + } + + nodeState.LamportTS = msg.LamportTS + + if existing, ok := gm.nodes.Load(msg.SenderID); ok { + existingNode := existing.(*NodeState) + + // ИСПРАВЛЕНО: используем логические часы для сравнения версий + if nodeState.LamportTS > existingNode.LamportTS || + (existingNode.LamportTS == 0 && nodeState.Incarnation > existingNode.Incarnation) { + existingNode.IP = nodeState.IP + existingNode.Port = nodeState.Port + existingNode.LastSeen = time.Now() + existingNode.Incarnation = nodeState.Incarnation + existingNode.IsAlive = true + existingNode.IsSuspect = false + existingNode.LamportTS = nodeState.LamportTS + existingNode.ClockOffsetMs = nodeState.ClockOffsetMs + if nodeState.Status != "" { + existingNode.Status = nodeState.Status + } + gm.nodes.Store(msg.SenderID, existingNode) + } else { + existingNode.LastSeen = time.Now() + if existingNode.IsSuspect { + existingNode.IsSuspect = false + } + // ИСПРАВЛЕНО: обновляем смещение часов + existingNode.ClockOffsetMs = nodeState.ClockOffsetMs + gm.nodes.Store(msg.SenderID, existingNode) + } + } else { + gm.nodes.Store(msg.SenderID, &nodeState) + gm.metrics.NodesDiscovered.Add(1) + + if gm.logger != nil { + gm.logger.Info(fmt.Sprintf("Discovered new node: %s (%s:%d)", msg.SenderID, msg.SenderIP, msg.SenderPort)) + } + } + + gm.updateMembership() +} + +// handleMembershipMessage обрабатывает сообщение о членстве +func (gm *GossipManager) handleMembershipMessage(msg *GossipMessage) { + // Обновление уже выполнено в updateNodeFromMessage +} + +// handleHealthMessage обрабатывает сообщение о здоровье +func (gm *GossipManager) handleHealthMessage(msg *GossipMessage) { + var payload map[string]string + if err := json.Unmarshal(msg.Payload, &payload); err == nil { + if action, ok := payload["action"]; ok && action == "ping" { + lamportTS := gm.lamportClock.Tick() + + response := &GossipMessage{ + Type: GossipTypeHealth, + SenderID: gm.localNode.ID, + SenderIP: gm.localNode.IP, + SenderPort: gm.localNode.Port, + Timestamp: time.Now().UnixMilli(), + Payload: json.RawMessage(`{"action":"pong"}`), + TTL: 1, + LamportTS: lamportTS, + } + gm.sendMessage(response, msg.SenderIP, msg.SenderPort) + } + } +} + +// handleLoadMessage обрабатывает сообщение о нагрузке +func (gm *GossipManager) handleLoadMessage(msg *GossipMessage) { + // TODO: Обновляем информацию о нагрузке узла +} + +// handleConfigMessage обрабатывает сообщение об изменении конфигурации +func (gm *GossipManager) handleConfigMessage(msg *GossipMessage) { + // TODO: Применяем изменения конфигурации +} + +// handleSyncMessage обрабатывает сообщение синхронизации +func (gm *GossipManager) handleSyncMessage(msg *GossipMessage) { + if msg.SenderID != gm.localNode.ID { + gm.sendFullState(msg.SenderIP, msg.SenderPort) + } +} + +// sendMessage отправляет gossip сообщение указанному узлу +func (gm *GossipManager) sendMessage(msg *GossipMessage, ip string, port int) { + if gm.udpConn == nil { + return + } + + data, err := json.Marshal(msg) + if err != nil { + return + } + + addr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", ip, port)) + if err != nil { + return + } + + gm.udpConn.WriteToUDP(data, addr) +} + +// forwardMessage пересылает сообщение дальше +func (gm *GossipManager) forwardMessage(msg *GossipMessage) { + if msg.TTL <= 0 { + return + } + + nodes := gm.getActiveNodes() + if len(nodes) == 0 { + return + } + + selected := gm.selectRandomNodes(nodes, gm.config.Fanout) + for _, node := range selected { + if node.ID != msg.SenderID && node.ID != gm.localNode.ID { + gm.sendMessage(msg, node.IP, node.Port) + } + } +} + +// sendFullState отправляет полное состояние кластера +func (gm *GossipManager) sendFullState(ip string, port int) { + nodes := gm.getAllNodes() + payload, _ := json.Marshal(nodes) + + lamportTS := gm.lamportClock.Tick() + + msg := &GossipMessage{ + Type: GossipTypeSync, + SenderID: gm.localNode.ID, + SenderIP: gm.localNode.IP, + SenderPort: gm.localNode.Port, + Timestamp: time.Now().UnixMilli(), + Payload: payload, + TTL: 1, + LamportTS: lamportTS, + } + + gm.sendMessage(msg, ip, port) +} + +// syncLoop периодически синхронизирует состояние с другими узлами +func (gm *GossipManager) syncLoop() { + defer gm.wg.Done() + + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + for { + select { + case <-gm.stopChan: + return + case <-ticker.C: + gm.syncWithCluster() + } + } +} + +// syncWithCluster синхронизирует состояние с кластером +func (gm *GossipManager) syncWithCluster() { + nodes := gm.getActiveNodes() + if len(nodes) == 0 { + return + } + + var target *NodeState + for _, node := range nodes { + if node.ID != gm.localNode.ID { + target = node + break + } + } + + if target == nil { + return + } + + gm.sendFullState(target.IP, target.Port) + gm.metrics.LastSync.Store(time.Now().UnixMilli()) +} + +// getActiveNodes возвращает список активных узлов +func (gm *GossipManager) getActiveNodes() []*NodeState { + nodes := make([]*NodeState, 0) + gm.nodes.Range(func(key, value interface{}) bool { + node := value.(*NodeState) + if node.IsAlive && node.ID != gm.localNode.ID { + nodes = append(nodes, node) + } + return true + }) + return nodes +} + +// getAllNodes возвращает список всех узлов +func (gm *GossipManager) getAllNodes() []*NodeState { + nodes := make([]*NodeState, 0) + gm.nodes.Range(func(key, value interface{}) bool { + nodes = append(nodes, value.(*NodeState)) + return true + }) + return nodes +} + +// selectRandomNodes выбирает случайные узлы из списка +func (gm *GossipManager) selectRandomNodes(nodes []*NodeState, count int) []*NodeState { + if len(nodes) == 0 { + return nil + } + + if count > len(nodes) { + count = len(nodes) + } + + shuffled := make([]*NodeState, len(nodes)) + copy(shuffled, nodes) + + for i := len(shuffled) - 1; i > 0; i-- { + j := rand.Intn(i + 1) + shuffled[i], shuffled[j] = shuffled[j], shuffled[i] + } + + return shuffled[:count] +} + +// updateMembership обновляет список членства +func (gm *GossipManager) updateMembership() { + nodes := gm.getAllNodes() + gm.membership.Store(nodes) +} + +// GetMembership возвращает текущий список членства +func (gm *GossipManager) GetMembership() []*NodeState { + return gm.membership.Load().([]*NodeState) +} + +// GetNodeByID возвращает узел по ID +func (gm *GossipManager) GetNodeByID(id string) *NodeState { + if val, ok := gm.nodes.Load(id); ok { + return val.(*NodeState) + } + return nil +} + +// AddSeed добавляет seed узел +func (gm *GossipManager) AddSeed(addr string) { + gm.mu.Lock() + defer gm.mu.Unlock() + + for _, existing := range gm.seeds { + if existing == addr { + return + } + } + gm.seeds = append(gm.seeds, addr) +} + +// Subscribe подписывается на события gossip +func (gm *GossipManager) Subscribe(id string) <-chan *GossipMessage { + gm.subscriberMu.Lock() + defer gm.subscriberMu.Unlock() + + ch := make(chan *GossipMessage, gm.config.BufferSize) + gm.subscribers[id] = ch + return ch +} + +// Unsubscribe отписывается от событий gossip +func (gm *GossipManager) Unsubscribe(id string) { + gm.subscriberMu.Lock() + defer gm.subscriberMu.Unlock() + + if ch, ok := gm.subscribers[id]; ok { + close(ch) + delete(gm.subscribers, id) + } +} + +// GetMetrics возвращает метрики gossip протокола +func (gm *GossipManager) GetMetrics() map[string]interface{} { + return map[string]interface{}{ + "messages_sent": gm.metrics.MessagesSent.Load(), + "messages_received": gm.metrics.MessagesReceived.Load(), + "messages_dropped": gm.metrics.MessagesDropped.Load(), + "nodes_discovered": gm.metrics.NodesDiscovered.Load(), + "nodes_removed": gm.metrics.NodesRemoved.Load(), + "last_broadcast": gm.metrics.LastBroadcast.Load(), + "last_sync": gm.metrics.LastSync.Load(), + "known_nodes": gm.getNodeCount(), + "active_nodes": len(gm.getActiveNodes()), + "clock_skew_detected": gm.metrics.ClockSkewDetected.Load(), + "lamport_ts": gm.lamportClock.Value(), + } +} + +// getNodeCount возвращает количество известных узлов +func (gm *GossipManager) getNodeCount() int { + count := 0 + gm.nodes.Range(func(key, value interface{}) bool { + count++ + return true + }) + return count +} + +// UpdateLocalNode обновляет локальное состояние узла +func (gm *GossipManager) UpdateLocalNode(status string, loadCPU float64, loadMemory float64) { + gm.mu.Lock() + defer gm.mu.Unlock() + + if gm.localNode != nil { + gm.localNode.Status = status + gm.localNode.LoadCPU = loadCPU + gm.localNode.LoadMemory = loadMemory + gm.localNode.LastHeartbeat = time.Now() + gm.localNode.Incarnation++ + gm.localNode.LamportTS = gm.lamportClock.Tick() + } +}