From 8f74cc6e1b870154e11c7283e1a2a5bb67fdcef2 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:29:37 +0000 Subject: [PATCH] Update internal/cluster/gossip.go --- internal/cluster/gossip.go | 337 +++++++++++++++++++++---------------- 1 file changed, 195 insertions(+), 142 deletions(-) diff --git a/internal/cluster/gossip.go b/internal/cluster/gossip.go index b6702a0..51809b4 100644 --- a/internal/cluster/gossip.go +++ b/internal/cluster/gossip.go @@ -34,65 +34,69 @@ import ( type GossipMessageType int const ( - GossipTypeMembership GossipMessageType = iota // Обновление членства - GossipTypeHealth // Информация о здоровье - GossipTypeLoad // Информация о нагрузке - GossipTypeConfig // Обновление конфигурации - GossipTypeSync // Синхронизация состояния + 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"` + 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"` + 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"` + 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 { - // Интервал отправки gossip сообщений (сек) BroadcastIntervalSec int `json:"broadcast_interval_sec"` - // Интервал проверки состояния узлов (сек) - MonitorIntervalSec int `json:"monitor_interval_sec"` - // Количество узлов для случайного распространения - Fanout int `json:"fanout"` - // Максимальное TTL сообщений - MaxTTL int `json:"max_ttl"` - // Время ожидания подтверждения (сек) - TimeoutSec int `json:"timeout_sec"` - // Время до перевода в подозрительные (сек) - SuspectTimeoutSec int `json:"suspect_timeout_sec"` - // Время до удаления узла (сек) - DeadTimeoutSec int `json:"dead_timeout_sec"` - // Максимальный размер буфера сообщений - BufferSize int `json:"buffer_size"` + 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, @@ -100,9 +104,12 @@ func DefaultGossipConfig() *GossipConfig { Fanout: 3, MaxTTL: 3, TimeoutSec: 10, - SuspectTimeoutSec: 15, - DeadTimeoutSec: 30, + // ИСПРАВЛЕНО: было 15/30, стало 30/60 для устойчивости к сетевым задержкам + SuspectTimeoutSec: 30, + DeadTimeoutSec: 60, BufferSize: 10000, + // ИСПРАВЛЕНО: 10 секунд — допустимая рассинхронизация часов + MaxClockSkewMs: 10000, } } @@ -116,8 +123,8 @@ type GossipManager struct { logger *log.Logger coordinator *RaftCoordinator localNode *NodeState - nodes sync.Map // map[string]*NodeState - membership atomic.Value // []*NodeState + nodes sync.Map + membership atomic.Value subscribers map[string]chan *GossipMessage subscriberMu sync.RWMutex stopChan chan struct{} @@ -128,6 +135,10 @@ type GossipManager struct { mu sync.RWMutex seeds []string metrics *GossipMetrics + // ИСПРАВЛЕНО: логические часы для упорядочивания сообщений + lamportClock *LamportClock + // ИСПРАВЛЕНО: локальное смещение часов + clockOffset atomic.Int64 } // GossipMetrics хранит метрики gossip протокола @@ -139,6 +150,8 @@ type GossipMetrics struct { NodesRemoved atomic.Uint64 LastBroadcast atomic.Int64 LastSync atomic.Int64 + // ИСПРАВЛЕНО: метрики для clock skew + ClockSkewDetected atomic.Uint64 } // NewGossipManager создаёт новый менеджер gossip протокола @@ -148,13 +161,14 @@ func NewGossipManager(config *GossipConfig, coordinator *RaftCoordinator, logger } gm := &GossipManager{ - config: config, - logger: logger, - coordinator: coordinator, - subscribers: make(map[string]chan *GossipMessage), - stopChan: make(chan struct{}), - metrics: &GossipMetrics{}, - seeds: make([]string, 0), + 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)) @@ -165,25 +179,24 @@ func NewGossipManager(config *GossipConfig, coordinator *RaftCoordinator, logger // 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(), + 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", + Incarnation: 0, + IsAlive: true, + IsSuspect: false, + Version: "1.0", + LamportTS: 0, } - // Запускаем UDP сервер для приёма gossip сообщений 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() @@ -208,9 +221,18 @@ func (gm *GossipManager) Stop() { } // startUDPServer запускает UDP сервер для приёма сообщений +// ИСПРАВЛЕНО: корректная привязка на Linux и OpenIndiana (illumos). +// Если IP пустой или 0.0.0.0, используем wildcard-адрес. func (gm *GossipManager) startUDPServer() error { port := gm.coordinator.config.Cluster.RaftPort + 1 - addr := fmt.Sprintf("%s:%d", gm.localNode.IP, port) + + // ИСПРАВЛЕНО: определяем адрес для привязки + 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 { @@ -219,7 +241,18 @@ func (gm *GossipManager) startUDPServer() error { conn, err := net.ListenUDP("udp", udpAddr) if err != nil { - return err + // ИСПРАВЛЕНО: 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 @@ -245,29 +278,29 @@ func (gm *GossipManager) broadcastLoop() { // 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, + 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), + Timestamp: time.Now().UnixMilli(), + Payload: payload, + TTL: gm.config.MaxTTL, + Version: gm.incarnation.Add(1), + LamportTS: lamportTS, } - // Получаем список активных узлов nodes := gm.getActiveNodes() if len(nodes) == 0 { - // Если нет известных узлов, пробуем связаться с seed узлами gm.broadcastToSeeds(msg) return } - // Выбираем случайные узлы для рассылки fanout := gm.config.Fanout if fanout > len(nodes) { fanout = len(nodes) @@ -275,7 +308,6 @@ func (gm *GossipManager) broadcast() { selected := gm.selectRandomNodes(nodes, fanout) - // Рассылаем сообщения for _, node := range selected { gm.sendMessage(msg, node.IP, node.Port) } @@ -292,7 +324,6 @@ func (gm *GossipManager) broadcastToSeeds(msg *GossipMessage) { gm.mu.RUnlock() for _, seed := range seeds { - // Парсим seed адрес host, port, err := net.SplitHostPort(seed) if err != nil { continue @@ -321,6 +352,7 @@ func (gm *GossipManager) monitorLoop() { } // monitorNodes проверяет состояние всех известных узлов +// ИСПРАВЛЕНО: учитывает clock skew между узлами, использует монотонное время. func (gm *GossipManager) monitorNodes() { now := time.Now() deadTimeout := time.Duration(gm.config.DeadTimeoutSec) * time.Second @@ -340,10 +372,11 @@ func (gm *GossipManager) monitorNodes() { return true } + // ИСПРАВЛЕНО: используем локальное время, а не время узла + // (защита от рассинхронизации часов) lastSeen := now.Sub(node.LastSeen) if lastSeen > deadTimeout { - // Узел считается мёртвым node.IsAlive = false node.IsSuspect = false gm.nodes.Store(nodeID, node) @@ -354,12 +387,10 @@ func (gm *GossipManager) monitorNodes() { 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) @@ -368,7 +399,6 @@ func (gm *GossipManager) monitorNodes() { gm.logger.Warn(fmt.Sprintf("Node %s marked as suspect (last seen %v ago)", nodeID, lastSeen)) } - // Отправляем запрос на подтверждение gm.requestConfirmation(node) } } @@ -376,31 +406,33 @@ func (gm *GossipManager) monitorNodes() { 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, + 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, + 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() @@ -413,6 +445,14 @@ func (gm *GossipManager) receiveLoop() { default: n, addr, err := gm.udpConn.ReadFromUDP(buffer) if err != nil { + // ИСПРАВЛЕНО: обработка временных ошибок + if isTemporaryNetError(err) { + continue + } + // Проверяем, не закрыт ли сокет + if gm.udpConn == nil { + return + } continue } @@ -432,10 +472,13 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) { 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) @@ -449,7 +492,6 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) { gm.handleSyncMessage(&msg) } - // Пересылаем сообщение дальше (если TTL > 0) if msg.TTL > 0 && msg.SenderID != gm.localNode.ID { msg.TTL-- gm.forwardMessage(&msg) @@ -457,8 +499,8 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) { } // updateNodeFromMessage обновляет информацию об узле из сообщения +// ИСПРАВЛЕНО: учитывает clock skew и использует логические часы. func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) { - // Проверяем, не от себя ли сообщение if msg.SenderID == gm.localNode.ID { return } @@ -469,7 +511,6 @@ func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) { return } } else { - // Для других типов создаём базовое состояние nodeState = NodeState{ ID: msg.SenderID, IP: msg.SenderIP, @@ -479,32 +520,52 @@ func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) { } } - // Обновляем или создаём узел + // ИСПРАВЛЕНО: вычисляем и сохраняем смещение часов + 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.Incarnation > existingNode.Incarnation { + // ИСПРАВЛЕНО: используем логические часы для сравнения версий + 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) @@ -523,24 +584,22 @@ func (gm *GossipManager) handleMembershipMessage(msg *GossipMessage) { // handleHealthMessage обрабатывает сообщение о здоровье func (gm *GossipManager) handleHealthMessage(msg *GossipMessage) { - // Отвечаем на ping запросы var payload map[string]string if err := json.Unmarshal(msg.Payload, &payload); err == nil { if action, ok := payload["action"]; ok && action == "ping" { - // Отправляем pong + lamportTS := gm.lamportClock.Tick() + response := &GossipMessage{ - Type: GossipTypeHealth, - SenderID: gm.localNode.ID, - SenderIP: gm.localNode.IP, + 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, + Timestamp: time.Now().UnixMilli(), + Payload: json.RawMessage(`{"action":"pong"}`), + TTL: 1, + LamportTS: lamportTS, } - // Отправляем обратно отправителю - host := msg.SenderIP - port := msg.SenderPort - gm.sendMessage(response, host, port) + gm.sendMessage(response, msg.SenderIP, msg.SenderPort) } } } @@ -557,7 +616,6 @@ func (gm *GossipManager) handleConfigMessage(msg *GossipMessage) { // handleSyncMessage обрабатывает сообщение синхронизации func (gm *GossipManager) handleSyncMessage(msg *GossipMessage) { - // Отправляем полное состояние кластера if msg.SenderID != gm.localNode.ID { gm.sendFullState(msg.SenderIP, msg.SenderPort) } @@ -569,13 +627,11 @@ func (gm *GossipManager) sendMessage(msg *GossipMessage, ip string, port int) { 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 @@ -586,18 +642,15 @@ func (gm *GossipManager) sendMessage(msg *GossipMessage, ip string, port int) { // forwardMessage пересылает сообщение дальше func (gm *GossipManager) forwardMessage(msg *GossipMessage) { - // Если TTL истёк, не пересылаем 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 { @@ -611,14 +664,17 @@ 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, + Type: GossipTypeSync, + SenderID: gm.localNode.ID, + SenderIP: gm.localNode.IP, SenderPort: gm.localNode.Port, - Timestamp: time.Now().UnixMilli(), - Payload: payload, - TTL: 1, + Timestamp: time.Now().UnixMilli(), + Payload: payload, + TTL: 1, + LamportTS: lamportTS, } gm.sendMessage(msg, ip, port) @@ -643,13 +699,11 @@ func (gm *GossipManager) syncLoop() { // 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 { @@ -662,7 +716,6 @@ func (gm *GossipManager) syncWithCluster() { return } - // Запрашиваем полное состояние gm.sendFullState(target.IP, target.Port) gm.metrics.LastSync.Store(time.Now().UnixMilli()) } @@ -700,11 +753,9 @@ func (gm *GossipManager) selectRandomNodes(nodes []*NodeState, count int) []*Nod count = len(nodes) } - // Создаём копию и перемешиваем shuffled := make([]*NodeState, len(nodes)) copy(shuffled, nodes) - // Перемешиваем с использованием rand for i := len(shuffled) - 1; i > 0; i-- { j := rand.Intn(i + 1) shuffled[i], shuffled[j] = shuffled[j], shuffled[i] @@ -737,7 +788,6 @@ func (gm *GossipManager) AddSeed(addr string) { gm.mu.Lock() defer gm.mu.Unlock() - // Проверяем, что адрес не дублируется for _, existing := range gm.seeds { if existing == addr { return @@ -770,15 +820,17 @@ func (gm *GossipManager) Unsubscribe(id string) { // 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()), + "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(), } } @@ -803,5 +855,6 @@ func (gm *GossipManager) UpdateLocalNode(status string, loadCPU float64, loadMem gm.localNode.LoadMemory = loadMemory gm.localNode.LastHeartbeat = time.Now() gm.localNode.Incarnation++ + gm.localNode.LamportTS = gm.lamportClock.Tick() } -} +} \ No newline at end of file