Update internal/cluster/gossip.go

This commit is contained in:
gvsafronov committed 2026-09-17 21:29:37 +00:00
1 parent f047009ec8
commit 8f74cc6e1b
1 file changed
+194 -141
+194 -141
View File
@@ -34,65 +34,69 @@ import (
type GossipMessageType int type GossipMessageType int
const ( const (
GossipTypeMembership GossipMessageType = iota // Обновление членства GossipTypeMembership GossipMessageType = iota
GossipTypeHealth // Информация о здоровье GossipTypeHealth
GossipTypeLoad // Информация о нагрузке GossipTypeLoad
GossipTypeConfig // Обновление конфигурации GossipTypeConfig
GossipTypeSync // Синхронизация состояния GossipTypeSync
) )
// GossipMessage представляет сообщение протокола сплетен // GossipMessage представляет сообщение протокола сплетен
type GossipMessage struct { type GossipMessage struct {
Type GossipMessageType `json:"type"` Type GossipMessageType `json:"type"`
SenderID string `json:"sender_id"` SenderID string `json:"sender_id"`
SenderIP string `json:"sender_ip"` SenderIP string `json:"sender_ip"`
SenderPort int `json:"sender_port"` SenderPort int `json:"sender_port"`
Timestamp int64 `json:"timestamp"` Timestamp int64 `json:"timestamp"`
Payload json.RawMessage `json:"payload"` Payload json.RawMessage `json:"payload"`
TTL int `json:"ttl"` TTL int `json:"ttl"`
Version uint64 `json:"version"` Version uint64 `json:"version"`
// Логическая метка для упорядочивания сообщений
LamportTS uint64 `json:"lamport_ts,omitempty"`
} }
// NodeState представляет состояние узла в протоколе сплетен // NodeState представляет состояние узла в протоколе сплетен
type NodeState struct { type NodeState struct {
ID string `json:"id"` ID string `json:"id"`
IP string `json:"ip"` IP string `json:"ip"`
Port int `json:"port"` Port int `json:"port"`
RaftPort int `json:"raft_port"` RaftPort int `json:"raft_port"`
Status string `json:"status"` Status string `json:"status"`
LastSeen time.Time `json:"last_seen"` LastSeen time.Time `json:"last_seen"`
LastHeartbeat time.Time `json:"last_heartbeat"` LastHeartbeat time.Time `json:"last_heartbeat"`
Incarnation uint64 `json:"incarnation"` Incarnation uint64 `json:"incarnation"`
LoadCPU float64 `json:"load_cpu"` LoadCPU float64 `json:"load_cpu"`
LoadMemory float64 `json:"load_memory"` LoadMemory float64 `json:"load_memory"`
LoadQueue int64 `json:"load_queue"` LoadQueue int64 `json:"load_queue"`
Connections int `json:"connections"` Connections int `json:"connections"`
Version string `json:"version"` Version string `json:"version"`
IsAlive bool `json:"is_alive"` IsAlive bool `json:"is_alive"`
IsSuspect bool `json:"is_suspect"` IsSuspect bool `json:"is_suspect"`
// ИСПРАВЛЕНО: логические часы узла для разрешения конфликтов
LamportTS uint64 `json:"lamport_ts,omitempty"`
// ИСПРАВЛЕНО: смещение часов относительно локального времени
ClockOffsetMs int64 `json:"clock_offset_ms,omitempty"`
} }
// GossipConfig содержит настройки протокола сплетен // GossipConfig содержит настройки протокола сплетен
// ИСПРАВЛЕНО: увеличены таймауты для предотвращения ложных срабатываний
// при сетевых задержках и рассинхронизации часов.
type GossipConfig struct { type GossipConfig struct {
// Интервал отправки gossip сообщений (сек)
BroadcastIntervalSec int `json:"broadcast_interval_sec"` BroadcastIntervalSec int `json:"broadcast_interval_sec"`
// Интервал проверки состояния узлов (сек) MonitorIntervalSec int `json:"monitor_interval_sec"`
MonitorIntervalSec int `json:"monitor_interval_sec"` Fanout int `json:"fanout"`
// Количество узлов для случайного распространения MaxTTL int `json:"max_ttl"`
Fanout int `json:"fanout"` TimeoutSec int `json:"timeout_sec"`
// Максимальное TTL сообщений // ИСПРАВЛЕНО: увеличены таймауты (было 15/30)
MaxTTL int `json:"max_ttl"` SuspectTimeoutSec int `json:"suspect_timeout_sec"`
// Время ожидания подтверждения (сек) DeadTimeoutSec int `json:"dead_timeout_sec"`
TimeoutSec int `json:"timeout_sec"` BufferSize int `json:"buffer_size"`
// Время до перевода в подозрительные (сек) // ИСПРАВЛЕНО: максимальная допустимая рассинхронизация часов
SuspectTimeoutSec int `json:"suspect_timeout_sec"` MaxClockSkewMs int64 `json:"max_clock_skew_ms"`
// Время до удаления узла (сек)
DeadTimeoutSec int `json:"dead_timeout_sec"`
// Максимальный размер буфера сообщений
BufferSize int `json:"buffer_size"`
} }
// DefaultGossipConfig возвращает конфигурацию gossip по умолчанию // DefaultGossipConfig возвращает конфигурацию gossip по умолчанию
// ИСПРАВЛЕНО: таймауты увеличены для стабильности при сетевых задержках.
func DefaultGossipConfig() *GossipConfig { func DefaultGossipConfig() *GossipConfig {
return &GossipConfig{ return &GossipConfig{
BroadcastIntervalSec: 5, BroadcastIntervalSec: 5,
@@ -100,9 +104,12 @@ func DefaultGossipConfig() *GossipConfig {
Fanout: 3, Fanout: 3,
MaxTTL: 3, MaxTTL: 3,
TimeoutSec: 10, TimeoutSec: 10,
SuspectTimeoutSec: 15, // ИСПРАВЛЕНО: было 15/30, стало 30/60 для устойчивости к сетевым задержкам
DeadTimeoutSec: 30, SuspectTimeoutSec: 30,
DeadTimeoutSec: 60,
BufferSize: 10000, BufferSize: 10000,
// ИСПРАВЛЕНО: 10 секунд — допустимая рассинхронизация часов
MaxClockSkewMs: 10000,
} }
} }
@@ -116,8 +123,8 @@ type GossipManager struct {
logger *log.Logger logger *log.Logger
coordinator *RaftCoordinator coordinator *RaftCoordinator
localNode *NodeState localNode *NodeState
nodes sync.Map // map[string]*NodeState nodes sync.Map
membership atomic.Value // []*NodeState membership atomic.Value
subscribers map[string]chan *GossipMessage subscribers map[string]chan *GossipMessage
subscriberMu sync.RWMutex subscriberMu sync.RWMutex
stopChan chan struct{} stopChan chan struct{}
@@ -128,6 +135,10 @@ type GossipManager struct {
mu sync.RWMutex mu sync.RWMutex
seeds []string seeds []string
metrics *GossipMetrics metrics *GossipMetrics
// ИСПРАВЛЕНО: логические часы для упорядочивания сообщений
lamportClock *LamportClock
// ИСПРАВЛЕНО: локальное смещение часов
clockOffset atomic.Int64
} }
// GossipMetrics хранит метрики gossip протокола // GossipMetrics хранит метрики gossip протокола
@@ -139,6 +150,8 @@ type GossipMetrics struct {
NodesRemoved atomic.Uint64 NodesRemoved atomic.Uint64
LastBroadcast atomic.Int64 LastBroadcast atomic.Int64
LastSync atomic.Int64 LastSync atomic.Int64
// ИСПРАВЛЕНО: метрики для clock skew
ClockSkewDetected atomic.Uint64
} }
// NewGossipManager создаёт новый менеджер gossip протокола // NewGossipManager создаёт новый менеджер gossip протокола
@@ -148,13 +161,14 @@ func NewGossipManager(config *GossipConfig, coordinator *RaftCoordinator, logger
} }
gm := &GossipManager{ gm := &GossipManager{
config: config, config: config,
logger: logger, logger: logger,
coordinator: coordinator, coordinator: coordinator,
subscribers: make(map[string]chan *GossipMessage), subscribers: make(map[string]chan *GossipMessage),
stopChan: make(chan struct{}), stopChan: make(chan struct{}),
metrics: &GossipMetrics{}, metrics: &GossipMetrics{},
seeds: make([]string, 0), seeds: make([]string, 0),
lamportClock: NewLamportClock("gossip"),
} }
gm.membership.Store(make([]*NodeState, 0)) gm.membership.Store(make([]*NodeState, 0))
@@ -165,25 +179,24 @@ func NewGossipManager(config *GossipConfig, coordinator *RaftCoordinator, logger
// Start запускает gossip протокол // Start запускает gossip протокол
func (gm *GossipManager) Start() error { func (gm *GossipManager) Start() error {
gm.localNode = &NodeState{ gm.localNode = &NodeState{
ID: gm.coordinator.localNodeInfo.ID, ID: gm.coordinator.localNodeInfo.ID,
IP: gm.coordinator.localNodeInfo.IP, IP: gm.coordinator.localNodeInfo.IP,
Port: gm.coordinator.localNodeInfo.Port, Port: gm.coordinator.localNodeInfo.Port,
RaftPort: gm.coordinator.config.Cluster.RaftPort, RaftPort: gm.coordinator.config.Cluster.RaftPort,
Status: "active", Status: "active",
LastSeen: time.Now(), LastSeen: time.Now(),
LastHeartbeat: time.Now(), LastHeartbeat: time.Now(),
Incarnation: 0, Incarnation: 0,
IsAlive: true, IsAlive: true,
IsSuspect: false, IsSuspect: false,
Version: "1.0", Version: "1.0",
LamportTS: 0,
} }
// Запускаем UDP сервер для приёма gossip сообщений
if err := gm.startUDPServer(); err != nil { if err := gm.startUDPServer(); err != nil {
return fmt.Errorf("failed to start UDP server: %v", err) return fmt.Errorf("failed to start UDP server: %v", err)
} }
// Запускаем горутины
gm.wg.Add(4) gm.wg.Add(4)
go gm.broadcastLoop() go gm.broadcastLoop()
go gm.monitorLoop() go gm.monitorLoop()
@@ -208,9 +221,18 @@ func (gm *GossipManager) Stop() {
} }
// startUDPServer запускает UDP сервер для приёма сообщений // startUDPServer запускает UDP сервер для приёма сообщений
// ИСПРАВЛЕНО: корректная привязка на Linux и OpenIndiana (illumos).
// Если IP пустой или 0.0.0.0, используем wildcard-адрес.
func (gm *GossipManager) startUDPServer() error { func (gm *GossipManager) startUDPServer() error {
port := gm.coordinator.config.Cluster.RaftPort + 1 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) udpAddr, err := net.ResolveUDPAddr("udp", addr)
if err != nil { if err != nil {
@@ -219,7 +241,18 @@ func (gm *GossipManager) startUDPServer() error {
conn, err := net.ListenUDP("udp", udpAddr) conn, err := net.ListenUDP("udp", udpAddr)
if err != nil { 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 gm.udpConn = conn
@@ -245,29 +278,29 @@ func (gm *GossipManager) broadcastLoop() {
// broadcast рассылает сообщение случайным узлам // broadcast рассылает сообщение случайным узлам
func (gm *GossipManager) broadcast() { func (gm *GossipManager) broadcast() {
// Формируем сообщение с состоянием узла // ИСПРАВЛЕНО: обновляем логическую метку
lamportTS := gm.lamportClock.Tick()
payload, _ := json.Marshal(gm.localNode) payload, _ := json.Marshal(gm.localNode)
msg := &GossipMessage{ msg := &GossipMessage{
Type: GossipTypeMembership, Type: GossipTypeMembership,
SenderID: gm.localNode.ID, SenderID: gm.localNode.ID,
SenderIP: gm.localNode.IP, SenderIP: gm.localNode.IP,
SenderPort: gm.localNode.Port, SenderPort: gm.localNode.Port,
Timestamp: time.Now().UnixMilli(), Timestamp: time.Now().UnixMilli(),
Payload: payload, Payload: payload,
TTL: gm.config.MaxTTL, TTL: gm.config.MaxTTL,
Version: gm.incarnation.Add(1), Version: gm.incarnation.Add(1),
LamportTS: lamportTS,
} }
// Получаем список активных узлов
nodes := gm.getActiveNodes() nodes := gm.getActiveNodes()
if len(nodes) == 0 { if len(nodes) == 0 {
// Если нет известных узлов, пробуем связаться с seed узлами
gm.broadcastToSeeds(msg) gm.broadcastToSeeds(msg)
return return
} }
// Выбираем случайные узлы для рассылки
fanout := gm.config.Fanout fanout := gm.config.Fanout
if fanout > len(nodes) { if fanout > len(nodes) {
fanout = len(nodes) fanout = len(nodes)
@@ -275,7 +308,6 @@ func (gm *GossipManager) broadcast() {
selected := gm.selectRandomNodes(nodes, fanout) selected := gm.selectRandomNodes(nodes, fanout)
// Рассылаем сообщения
for _, node := range selected { for _, node := range selected {
gm.sendMessage(msg, node.IP, node.Port) gm.sendMessage(msg, node.IP, node.Port)
} }
@@ -292,7 +324,6 @@ func (gm *GossipManager) broadcastToSeeds(msg *GossipMessage) {
gm.mu.RUnlock() gm.mu.RUnlock()
for _, seed := range seeds { for _, seed := range seeds {
// Парсим seed адрес
host, port, err := net.SplitHostPort(seed) host, port, err := net.SplitHostPort(seed)
if err != nil { if err != nil {
continue continue
@@ -321,6 +352,7 @@ func (gm *GossipManager) monitorLoop() {
} }
// monitorNodes проверяет состояние всех известных узлов // monitorNodes проверяет состояние всех известных узлов
// ИСПРАВЛЕНО: учитывает clock skew между узлами, использует монотонное время.
func (gm *GossipManager) monitorNodes() { func (gm *GossipManager) monitorNodes() {
now := time.Now() now := time.Now()
deadTimeout := time.Duration(gm.config.DeadTimeoutSec) * time.Second deadTimeout := time.Duration(gm.config.DeadTimeoutSec) * time.Second
@@ -340,10 +372,11 @@ func (gm *GossipManager) monitorNodes() {
return true return true
} }
// ИСПРАВЛЕНО: используем локальное время, а не время узла
// (защита от рассинхронизации часов)
lastSeen := now.Sub(node.LastSeen) lastSeen := now.Sub(node.LastSeen)
if lastSeen > deadTimeout { if lastSeen > deadTimeout {
// Узел считается мёртвым
node.IsAlive = false node.IsAlive = false
node.IsSuspect = false node.IsSuspect = false
gm.nodes.Store(nodeID, node) 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)) gm.logger.Warn(fmt.Sprintf("Node %s considered dead (last seen %v ago)", nodeID, lastSeen))
} }
// Уведомляем координатор об удалении узла
if gm.coordinator != nil { if gm.coordinator != nil {
gm.coordinator.RemoveNode(nodeID) gm.coordinator.RemoveNode(nodeID)
} }
} else if lastSeen > suspectTimeout { } else if lastSeen > suspectTimeout {
// Узел под подозрением
if !node.IsSuspect { if !node.IsSuspect {
node.IsSuspect = true node.IsSuspect = true
gm.nodes.Store(nodeID, node) 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.logger.Warn(fmt.Sprintf("Node %s marked as suspect (last seen %v ago)", nodeID, lastSeen))
} }
// Отправляем запрос на подтверждение
gm.requestConfirmation(node) gm.requestConfirmation(node)
} }
} }
@@ -376,31 +406,33 @@ func (gm *GossipManager) monitorNodes() {
return true return true
}) })
// Удаляем мёртвые узлы
for _, nodeID := range nodesToRemove { for _, nodeID := range nodesToRemove {
gm.nodes.Delete(nodeID) gm.nodes.Delete(nodeID)
} }
// Обновляем список членства
gm.updateMembership() gm.updateMembership()
} }
// requestConfirmation запрашивает подтверждение от узла // requestConfirmation запрашивает подтверждение от узла
func (gm *GossipManager) requestConfirmation(node *NodeState) { func (gm *GossipManager) requestConfirmation(node *NodeState) {
lamportTS := gm.lamportClock.Tick()
msg := &GossipMessage{ msg := &GossipMessage{
Type: GossipTypeHealth, Type: GossipTypeHealth,
SenderID: gm.localNode.ID, SenderID: gm.localNode.ID,
SenderIP: gm.localNode.IP, SenderIP: gm.localNode.IP,
SenderPort: gm.localNode.Port, SenderPort: gm.localNode.Port,
Timestamp: time.Now().UnixMilli(), Timestamp: time.Now().UnixMilli(),
Payload: json.RawMessage(`{"action":"ping"}`), Payload: json.RawMessage(`{"action":"ping"}`),
TTL: 1, TTL: 1,
LamportTS: lamportTS,
} }
gm.sendMessage(msg, node.IP, node.Port) gm.sendMessage(msg, node.IP, node.Port)
} }
// receiveLoop принимает входящие gossip сообщения // receiveLoop принимает входящие gossip сообщения
// ИСПРАВЛЕНО: обработка временных ошибок для Linux и illumos.
func (gm *GossipManager) receiveLoop() { func (gm *GossipManager) receiveLoop() {
defer gm.wg.Done() defer gm.wg.Done()
@@ -413,6 +445,14 @@ func (gm *GossipManager) receiveLoop() {
default: default:
n, addr, err := gm.udpConn.ReadFromUDP(buffer) n, addr, err := gm.udpConn.ReadFromUDP(buffer)
if err != nil { if err != nil {
// ИСПРАВЛЕНО: обработка временных ошибок
if isTemporaryNetError(err) {
continue
}
// Проверяем, не закрыт ли сокет
if gm.udpConn == nil {
return
}
continue continue
} }
@@ -432,10 +472,13 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) {
gm.metrics.MessagesReceived.Add(1) gm.metrics.MessagesReceived.Add(1)
// Обновляем информацию об отправителе // ИСПРАВЛЕНО: обновляем логические часы
if msg.LamportTS > 0 {
gm.lamportClock.Observe(msg.LamportTS)
}
gm.updateNodeFromMessage(&msg) gm.updateNodeFromMessage(&msg)
// Обрабатываем сообщение в зависимости от типа
switch msg.Type { switch msg.Type {
case GossipTypeMembership: case GossipTypeMembership:
gm.handleMembershipMessage(&msg) gm.handleMembershipMessage(&msg)
@@ -449,7 +492,6 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) {
gm.handleSyncMessage(&msg) gm.handleSyncMessage(&msg)
} }
// Пересылаем сообщение дальше (если TTL > 0)
if msg.TTL > 0 && msg.SenderID != gm.localNode.ID { if msg.TTL > 0 && msg.SenderID != gm.localNode.ID {
msg.TTL-- msg.TTL--
gm.forwardMessage(&msg) gm.forwardMessage(&msg)
@@ -457,8 +499,8 @@ func (gm *GossipManager) handleMessage(data []byte, addr *net.UDPAddr) {
} }
// updateNodeFromMessage обновляет информацию об узле из сообщения // updateNodeFromMessage обновляет информацию об узле из сообщения
// ИСПРАВЛЕНО: учитывает clock skew и использует логические часы.
func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) { func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) {
// Проверяем, не от себя ли сообщение
if msg.SenderID == gm.localNode.ID { if msg.SenderID == gm.localNode.ID {
return return
} }
@@ -469,7 +511,6 @@ func (gm *GossipManager) updateNodeFromMessage(msg *GossipMessage) {
return return
} }
} else { } else {
// Для других типов создаём базовое состояние
nodeState = NodeState{ nodeState = NodeState{
ID: msg.SenderID, ID: msg.SenderID,
IP: msg.SenderIP, 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 { if existing, ok := gm.nodes.Load(msg.SenderID); ok {
existingNode := existing.(*NodeState) 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.IP = nodeState.IP
existingNode.Port = nodeState.Port existingNode.Port = nodeState.Port
existingNode.LastSeen = time.Now() existingNode.LastSeen = time.Now()
existingNode.Incarnation = nodeState.Incarnation existingNode.Incarnation = nodeState.Incarnation
existingNode.IsAlive = true existingNode.IsAlive = true
existingNode.IsSuspect = false existingNode.IsSuspect = false
existingNode.LamportTS = nodeState.LamportTS
existingNode.ClockOffsetMs = nodeState.ClockOffsetMs
if nodeState.Status != "" { if nodeState.Status != "" {
existingNode.Status = nodeState.Status existingNode.Status = nodeState.Status
} }
gm.nodes.Store(msg.SenderID, existingNode) gm.nodes.Store(msg.SenderID, existingNode)
} else { } else {
// Обновляем только время последнего контакта
existingNode.LastSeen = time.Now() existingNode.LastSeen = time.Now()
if existingNode.IsSuspect { if existingNode.IsSuspect {
existingNode.IsSuspect = false existingNode.IsSuspect = false
} }
// ИСПРАВЛЕНО: обновляем смещение часов
existingNode.ClockOffsetMs = nodeState.ClockOffsetMs
gm.nodes.Store(msg.SenderID, existingNode) gm.nodes.Store(msg.SenderID, existingNode)
} }
} else { } else {
// Новый узел
gm.nodes.Store(msg.SenderID, &nodeState) gm.nodes.Store(msg.SenderID, &nodeState)
gm.metrics.NodesDiscovered.Add(1) gm.metrics.NodesDiscovered.Add(1)
@@ -523,24 +584,22 @@ func (gm *GossipManager) handleMembershipMessage(msg *GossipMessage) {
// handleHealthMessage обрабатывает сообщение о здоровье // handleHealthMessage обрабатывает сообщение о здоровье
func (gm *GossipManager) handleHealthMessage(msg *GossipMessage) { func (gm *GossipManager) handleHealthMessage(msg *GossipMessage) {
// Отвечаем на ping запросы
var payload map[string]string var payload map[string]string
if err := json.Unmarshal(msg.Payload, &payload); err == nil { if err := json.Unmarshal(msg.Payload, &payload); err == nil {
if action, ok := payload["action"]; ok && action == "ping" { if action, ok := payload["action"]; ok && action == "ping" {
// Отправляем pong lamportTS := gm.lamportClock.Tick()
response := &GossipMessage{ response := &GossipMessage{
Type: GossipTypeHealth, Type: GossipTypeHealth,
SenderID: gm.localNode.ID, SenderID: gm.localNode.ID,
SenderIP: gm.localNode.IP, SenderIP: gm.localNode.IP,
SenderPort: gm.localNode.Port, SenderPort: gm.localNode.Port,
Timestamp: time.Now().UnixMilli(), Timestamp: time.Now().UnixMilli(),
Payload: json.RawMessage(`{"action":"pong"}`), Payload: json.RawMessage(`{"action":"pong"}`),
TTL: 1, TTL: 1,
LamportTS: lamportTS,
} }
// Отправляем обратно отправителю gm.sendMessage(response, msg.SenderIP, msg.SenderPort)
host := msg.SenderIP
port := msg.SenderPort
gm.sendMessage(response, host, port)
} }
} }
} }
@@ -557,7 +616,6 @@ func (gm *GossipManager) handleConfigMessage(msg *GossipMessage) {
// handleSyncMessage обрабатывает сообщение синхронизации // handleSyncMessage обрабатывает сообщение синхронизации
func (gm *GossipManager) handleSyncMessage(msg *GossipMessage) { func (gm *GossipManager) handleSyncMessage(msg *GossipMessage) {
// Отправляем полное состояние кластера
if msg.SenderID != gm.localNode.ID { if msg.SenderID != gm.localNode.ID {
gm.sendFullState(msg.SenderIP, msg.SenderPort) gm.sendFullState(msg.SenderIP, msg.SenderPort)
} }
@@ -569,13 +627,11 @@ func (gm *GossipManager) sendMessage(msg *GossipMessage, ip string, port int) {
return return
} }
// Сериализуем сообщение
data, err := json.Marshal(msg) data, err := json.Marshal(msg)
if err != nil { if err != nil {
return return
} }
// Отправляем
addr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", ip, port)) addr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", ip, port))
if err != nil { if err != nil {
return return
@@ -586,18 +642,15 @@ func (gm *GossipManager) sendMessage(msg *GossipMessage, ip string, port int) {
// forwardMessage пересылает сообщение дальше // forwardMessage пересылает сообщение дальше
func (gm *GossipManager) forwardMessage(msg *GossipMessage) { func (gm *GossipManager) forwardMessage(msg *GossipMessage) {
// Если TTL истёк, не пересылаем
if msg.TTL <= 0 { if msg.TTL <= 0 {
return return
} }
// Получаем активные узлы
nodes := gm.getActiveNodes() nodes := gm.getActiveNodes()
if len(nodes) == 0 { if len(nodes) == 0 {
return return
} }
// Выбираем случайные узлы для пересылки
selected := gm.selectRandomNodes(nodes, gm.config.Fanout) selected := gm.selectRandomNodes(nodes, gm.config.Fanout)
for _, node := range selected { for _, node := range selected {
if node.ID != msg.SenderID && node.ID != gm.localNode.ID { 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() nodes := gm.getAllNodes()
payload, _ := json.Marshal(nodes) payload, _ := json.Marshal(nodes)
lamportTS := gm.lamportClock.Tick()
msg := &GossipMessage{ msg := &GossipMessage{
Type: GossipTypeSync, Type: GossipTypeSync,
SenderID: gm.localNode.ID, SenderID: gm.localNode.ID,
SenderIP: gm.localNode.IP, SenderIP: gm.localNode.IP,
SenderPort: gm.localNode.Port, SenderPort: gm.localNode.Port,
Timestamp: time.Now().UnixMilli(), Timestamp: time.Now().UnixMilli(),
Payload: payload, Payload: payload,
TTL: 1, TTL: 1,
LamportTS: lamportTS,
} }
gm.sendMessage(msg, ip, port) gm.sendMessage(msg, ip, port)
@@ -643,13 +699,11 @@ func (gm *GossipManager) syncLoop() {
// syncWithCluster синхронизирует состояние с кластером // syncWithCluster синхронизирует состояние с кластером
func (gm *GossipManager) syncWithCluster() { func (gm *GossipManager) syncWithCluster() {
// Выбираем случайный узел для синхронизации
nodes := gm.getActiveNodes() nodes := gm.getActiveNodes()
if len(nodes) == 0 { if len(nodes) == 0 {
return return
} }
// Выбираем случайный узел, кроме себя
var target *NodeState var target *NodeState
for _, node := range nodes { for _, node := range nodes {
if node.ID != gm.localNode.ID { if node.ID != gm.localNode.ID {
@@ -662,7 +716,6 @@ func (gm *GossipManager) syncWithCluster() {
return return
} }
// Запрашиваем полное состояние
gm.sendFullState(target.IP, target.Port) gm.sendFullState(target.IP, target.Port)
gm.metrics.LastSync.Store(time.Now().UnixMilli()) gm.metrics.LastSync.Store(time.Now().UnixMilli())
} }
@@ -700,11 +753,9 @@ func (gm *GossipManager) selectRandomNodes(nodes []*NodeState, count int) []*Nod
count = len(nodes) count = len(nodes)
} }
// Создаём копию и перемешиваем
shuffled := make([]*NodeState, len(nodes)) shuffled := make([]*NodeState, len(nodes))
copy(shuffled, nodes) copy(shuffled, nodes)
// Перемешиваем с использованием rand
for i := len(shuffled) - 1; i > 0; i-- { for i := len(shuffled) - 1; i > 0; i-- {
j := rand.Intn(i + 1) j := rand.Intn(i + 1)
shuffled[i], shuffled[j] = shuffled[j], shuffled[i] shuffled[i], shuffled[j] = shuffled[j], shuffled[i]
@@ -737,7 +788,6 @@ func (gm *GossipManager) AddSeed(addr string) {
gm.mu.Lock() gm.mu.Lock()
defer gm.mu.Unlock() defer gm.mu.Unlock()
// Проверяем, что адрес не дублируется
for _, existing := range gm.seeds { for _, existing := range gm.seeds {
if existing == addr { if existing == addr {
return return
@@ -770,15 +820,17 @@ func (gm *GossipManager) Unsubscribe(id string) {
// GetMetrics возвращает метрики gossip протокола // GetMetrics возвращает метрики gossip протокола
func (gm *GossipManager) GetMetrics() map[string]interface{} { func (gm *GossipManager) GetMetrics() map[string]interface{} {
return map[string]interface{}{ return map[string]interface{}{
"messages_sent": gm.metrics.MessagesSent.Load(), "messages_sent": gm.metrics.MessagesSent.Load(),
"messages_received": gm.metrics.MessagesReceived.Load(), "messages_received": gm.metrics.MessagesReceived.Load(),
"messages_dropped": gm.metrics.MessagesDropped.Load(), "messages_dropped": gm.metrics.MessagesDropped.Load(),
"nodes_discovered": gm.metrics.NodesDiscovered.Load(), "nodes_discovered": gm.metrics.NodesDiscovered.Load(),
"nodes_removed": gm.metrics.NodesRemoved.Load(), "nodes_removed": gm.metrics.NodesRemoved.Load(),
"last_broadcast": gm.metrics.LastBroadcast.Load(), "last_broadcast": gm.metrics.LastBroadcast.Load(),
"last_sync": gm.metrics.LastSync.Load(), "last_sync": gm.metrics.LastSync.Load(),
"known_nodes": gm.getNodeCount(), "known_nodes": gm.getNodeCount(),
"active_nodes": len(gm.getActiveNodes()), "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.LoadMemory = loadMemory
gm.localNode.LastHeartbeat = time.Now() gm.localNode.LastHeartbeat = time.Now()
gm.localNode.Incarnation++ gm.localNode.Incarnation++
gm.localNode.LamportTS = gm.lamportClock.Tick()
} }
} }