Update internal/cluster/node.go

This commit is contained in:
gvsafronov committed 2026-09-17 21:28:20 +00:00
1 parent 0d7bad66bb
commit f047009ec8
1 file changed
+104 -20
+104 -20
View File
@@ -78,6 +78,9 @@ type NodeRequest struct {
FromNode string `json:"from_node"` // ID узла-отправителя FromNode string `json:"from_node"` // ID узла-отправителя
Data json.RawMessage `json:"data"` // Данные запроса в формате JSON Data json.RawMessage `json:"data"` // Данные запроса в формате JSON
Timestamp int64 `json:"timestamp"` // Временная метка запроса (Unix millis) Timestamp int64 `json:"timestamp"` // Временная метка запроса (Unix millis)
// ИСПРАВЛЕНО: логическая метка Лампорта для упорядочивания сообщений
// без зависимости от физических часов узлов.
LamportTS uint64 `json:"lamport_ts,omitempty"`
} }
// NodeInfo представляет информацию об узле в кластере. // NodeInfo представляет информацию об узле в кластере.
@@ -147,6 +150,7 @@ func SafeGoWithLogger(fn func(), logger *log.Logger, name string) {
// - mu защищает операции, изменяющие состояние узла // - mu защищает операции, изменяющие состояние узла
// - connPool использует sync.Map для потокобезопасного хранения соединений // - connPool использует sync.Map для потокобезопасного хранения соединений
// - routinesMu защищает доступ к recoverableRoutines // - routinesMu защищает доступ к recoverableRoutines
// - docVersions использует sync.Map для потокобезопасного хранения версий документов
type Node struct { type Node struct {
// Основная информация об узле // Основная информация об узле
ID string // Уникальный идентификатор узла (UUID) ID string // Уникальный идентификатор узла (UUID)
@@ -191,6 +195,16 @@ type Node struct {
panicRecoveryMgr *PanicRecoveryManager // Менеджер восстановления panicRecoveryMgr *PanicRecoveryManager // Менеджер восстановления
recoverableRoutines map[string]*RecoverableRoutine // Зарегистрированные восстанавливаемые горутины recoverableRoutines map[string]*RecoverableRoutine // Зарегистрированные восстанавливаемые горутины
routinesMu sync.RWMutex // Защита доступа к recoverableRoutines routinesMu sync.RWMutex // Защита доступа к recoverableRoutines
// ИСПРАВЛЕНО: логические часы Лампорта для eventual consistency.
// Позволяют упорядочивать события без зависимости от физических часов узлов,
// что защищает от проблемы last-write-wins при рассинхронизации часов.
lamportClock *LamportClock
// ИСПРАВЛЕНО: карта версий документов для разрешения конфликтов репликации.
// Ключ — ID документа, значение — логическая метка последней применённой версии.
// Используется для игнорирования устаревших репликаций при рассинхронизации.
docVersions sync.Map // map[string]uint64 (docID -> lamport timestamp)
} }
// NodeConfig представляет конфигурацию для создания узла. // NodeConfig представляет конфигурацию для создания узла.
@@ -262,6 +276,7 @@ func NewNodeWithRecovery(ip string, port int, store *storage.Storage, logger *lo
replicator: replicator, replicator: replicator,
panicRecoveryMgr: panicRecoveryMgr, panicRecoveryMgr: panicRecoveryMgr,
recoverableRoutines: make(map[string]*RecoverableRoutine), recoverableRoutines: make(map[string]*RecoverableRoutine),
lamportClock: NewLamportClock(ip), // ИСПРАВЛЕНО: инициализация логических часов
} }
// Устанавливаем начальный статус и временные метки // Устанавливаем начальный статус и временные метки
@@ -356,6 +371,8 @@ func (n *Node) stopRecoverableRoutine(name string) {
// - Автоматически перезапускается при панике // - Автоматически перезапускается при панике
// - Обрабатывает ошибки таймаута без остановки сервера // - Обрабатывает ошибки таймаута без остановки сервера
// - Передаёт соединения в канал incomingConn для асинхронной обработки // - Передаёт соединения в канал incomingConn для асинхронной обработки
// - ИСПРАВЛЕНО: обрабатывает временные ошибки (EINTR, EAGAIN, EWOULDBLOCK),
// которые могут возникать на Linux и OpenIndiana (illumos)
func (n *Node) startTCPServer() { func (n *Node) startTCPServer() {
// Защита от паники с автоматическим перезапуском // Защита от паники с автоматическим перезапуском
defer func() { defer func() {
@@ -407,8 +424,10 @@ func (n *Node) startTCPServer() {
// Принимаем новое соединение // Принимаем новое соединение
conn, err := listener.Accept() conn, err := listener.Accept()
if err != nil { if err != nil {
// Проверяем, не является ли ошибка таймаутом (это нормально) // ИСПРАВЛЕНО: используем isTemporaryNetError для обработки
if netErr, ok := err.(net.Error); ok && netErr.Timeout() { // EINTR/EAGAIN/EWOULDBLOCK, которые могут возникать на illumos
// при получении сигналов, а также стандартных таймаутов.
if isTemporaryNetError(err) {
continue continue
} }
if n.logger != nil { if n.logger != nil {
@@ -529,9 +548,16 @@ func (n *Node) handleNodeRequest(conn net.Conn) {
// Обновляем время последнего контакта // Обновляем время последнего контакта
n.lastSeen.Store(time.Now().UnixMilli()) n.lastSeen.Store(time.Now().UnixMilli())
// ИСПРАВЛЕНО: обновляем логические часы при получении удалённого запроса.
// Это гарантирует, что наши последующие метки будут больше метки отправителя,
// что обеспечивает корректное упорядочивание событий.
if req.LamportTS > 0 {
n.lamportClock.Observe(req.LamportTS)
}
if n.logger != nil { if n.logger != nil {
n.logger.Debug(fmt.Sprintf("Node %s received request type %s from %s at %s", n.logger.Debug(fmt.Sprintf("Node %s received request type %s from %s at %s (lamport=%d)",
n.ID, req.Type, req.FromNode, time.UnixMilli(req.Timestamp).Format("15:04:05.000"))) n.ID, req.Type, req.FromNode, time.UnixMilli(req.Timestamp).Format("15:04:05.000"), req.LamportTS))
} }
// Маршрутизация по типу запроса // Маршрутизация по типу запроса
@@ -566,6 +592,11 @@ func (n *Node) handleNodeRequest(conn net.Conn) {
// 3. Создание документа с полями // 3. Создание документа с полями
// 4. Вставка документа в коллекцию // 4. Вставка документа в коллекцию
// 5. Логирование результата // 5. Логирование результата
//
// ИСПРАВЛЕНО: использует логические часы и версии документов для разрешения
// конфликтов last-write-wins. Если локальная версия документа новее полученной,
// репликация игнорируется. Это защищает от ситуации, когда при рассинхронизации
// часов узлов устаревшая версия перезаписывает более новую.
func (n *Node) handleReplicateRequest(data []byte) { func (n *Node) handleReplicateRequest(data []byte) {
// Защита от паники // Защита от паники
defer func() { defer func() {
@@ -585,6 +616,10 @@ func (n *Node) handleReplicateRequest(data []byte) {
Document map[string]interface{} `json:"document"` // Данные документа Document map[string]interface{} `json:"document"` // Данные документа
SourceNode string `json:"source_node"` // Узел-источник SourceNode string `json:"source_node"` // Узел-источник
ReplicaID string `json:"replica_id"` // Уникальный ID репликации ReplicaID string `json:"replica_id"` // Уникальный ID репликации
// ИСПРАВЛЕНО: логическая метка Лампорта от узла-источника
LamportTS uint64 `json:"lamport_ts,omitempty"`
// ИСПРАВЛЕНО: версия документа для разрешения конфликтов
DocVersion uint64 `json:"doc_version,omitempty"`
} }
if err := json.Unmarshal(data, &repData); err != nil { if err := json.Unmarshal(data, &repData); err != nil {
@@ -621,6 +656,26 @@ func (n *Node) handleReplicateRequest(data []byte) {
return return
} }
// ИСПРАВЛЕНО: проверяем версию документа для last-write-wins.
// Если у нас уже есть более новая версия — игнорируем старую репликацию.
// Это защищает от перезаписи актуальных данных устаревшими при
// рассинхронизации часов между узлами.
if existingVersion, ok := n.docVersions.Load(docID); ok {
if existingVersion.(uint64) > repData.DocVersion {
if n.logger != nil {
n.logger.Debug(fmt.Sprintf("Node %s ignoring stale replication for doc %s (local=%d, remote=%d)",
n.ID, docID, existingVersion.(uint64), repData.DocVersion))
}
return
}
}
// ИСПРАВЛЕНО: обновляем логические часы от удалённого узла
if repData.LamportTS > 0 {
n.lamportClock.Observe(repData.LamportTS)
}
newVersion := n.lamportClock.Tick()
// Создаём документ с полученными полями // Создаём документ с полученными полями
doc := storage.NewDocumentWithID(docID) doc := storage.NewDocumentWithID(docID)
for k, v := range repData.Document { for k, v := range repData.Document {
@@ -633,11 +688,14 @@ func (n *Node) handleReplicateRequest(data []byte) {
n.logger.Error(fmt.Sprintf("Node %s failed to replicate document: %v", n.ID, err)) n.logger.Error(fmt.Sprintf("Node %s failed to replicate document: %v", n.ID, err))
} }
} else { } else {
// ИСПРАВЛЕНО: сохраняем версию документа для разрешения будущих конфликтов
n.docVersions.Store(docID, newVersion)
// Логируем успешную репликацию // Логируем успешную репликацию
duration := time.Now().UnixMilli() - startTime duration := time.Now().UnixMilli() - startTime
if n.logger != nil { if n.logger != nil {
n.logger.Debug(fmt.Sprintf("Node %s replicated document %s from %s (took %d ms)", n.logger.Debug(fmt.Sprintf("Node %s replicated document %s from %s (took %d ms, version=%d)",
n.ID, doc.ID, repData.SourceNode, duration)) n.ID, doc.ID, repData.SourceNode, duration, newVersion))
} }
} }
} }
@@ -820,6 +878,7 @@ func (n *Node) handleSyncRequest(data []byte, conn net.Conn) {
// - node_id: ID текущего узла // - node_id: ID текущего узла
// - timestamp: время обработки // - timestamp: время обработки
// - uptime_ms: время работы узла в миллисекундах // - uptime_ms: время работы узла в миллисекундах
// - lamport_ts: текущая логическая метка (ИСПРАВЛЕНО)
func (n *Node) handleHeartbeatRequest(req NodeRequest, conn net.Conn) { func (n *Node) handleHeartbeatRequest(req NodeRequest, conn net.Conn) {
// Защита от паники // Защита от паники
defer func() { defer func() {
@@ -833,12 +892,18 @@ func (n *Node) handleHeartbeatRequest(req NodeRequest, conn net.Conn) {
// Обновляем время последнего контакта // Обновляем время последнего контакта
n.lastSeen.Store(time.Now().UnixMilli()) n.lastSeen.Store(time.Now().UnixMilli())
// ИСПРАВЛЕНО: обновляем логические часы от отправителя
if req.LamportTS > 0 {
n.lamportClock.Observe(req.LamportTS)
}
// Формируем ответ // Формируем ответ
response := map[string]interface{}{ response := map[string]interface{}{
"status": "alive", "status": "alive",
"node_id": n.ID, "node_id": n.ID,
"timestamp": time.Now().UnixMilli(), "timestamp": time.Now().UnixMilli(),
"uptime_ms": time.Now().UnixMilli() - n.startedAt, "uptime_ms": time.Now().UnixMilli() - n.startedAt,
"lamport_ts": n.lamportClock.Value(), // ИСПРАВЛЕНО: возвращаем логическую метку
} }
// Отправляем ответ // Отправляем ответ
@@ -884,11 +949,12 @@ func (n *Node) handleStatusSyncRequest(req NodeRequest, conn net.Conn) {
// Формируем ответ // Формируем ответ
response := map[string]interface{}{ response := map[string]interface{}{
"status": "synced", "status": "synced",
"node_id": n.ID, "node_id": n.ID,
"term": n.coordinator.GetCurrentTerm(), "term": n.coordinator.GetCurrentTerm(),
"is_leader": n.coordinator.IsLeader(), "is_leader": n.coordinator.IsLeader(),
"timestamp": time.Now().UnixMilli(), "timestamp": time.Now().UnixMilli(),
"lamport_ts": n.lamportClock.Value(), // ИСПРАВЛЕНО: добавляем логическую метку
} }
// Отправляем ответ // Отправляем ответ
@@ -954,7 +1020,8 @@ func (n *Node) heartbeatLoop() {
n.coordinator.SendHeartbeat(n.ID) n.coordinator.SendHeartbeat(n.ID)
n.lastSeen.Store(time.Now().UnixMilli()) n.lastSeen.Store(time.Now().UnixMilli())
if n.logger != nil { if n.logger != nil {
n.logger.Debug(fmt.Sprintf("Node %s sent heartbeat at %s", n.ID, n.GetLastSeenStr())) n.logger.Debug(fmt.Sprintf("Node %s sent heartbeat at %s (lamport=%d)",
n.ID, n.GetLastSeenStr(), n.lamportClock.Value()))
} }
} }
} }
@@ -1235,6 +1302,7 @@ func (n *Node) GetAddress() string {
// GetStats возвращает полную статистику узла. // GetStats возвращает полную статистику узла.
// Включает информацию о статусе, времени, нагрузке и компонентах. // Включает информацию о статусе, времени, нагрузке и компонентах.
// ИСПРАВЛЕНО: добавлена информация о логических часах.
func (n *Node) GetStats() map[string]interface{} { func (n *Node) GetStats() map[string]interface{} {
stats := map[string]interface{}{ stats := map[string]interface{}{
"id": n.ID, "id": n.ID,
@@ -1249,6 +1317,7 @@ func (n *Node) GetStats() map[string]interface{} {
"request_count": n.requestCount.Load(), "request_count": n.requestCount.Load(),
"bytes_rx": n.bytesRx.Load(), "bytes_rx": n.bytesRx.Load(),
"bytes_tx": n.bytesTx.Load(), "bytes_tx": n.bytesTx.Load(),
"lamport_ts": n.lamportClock.Value(), // ИСПРАВЛЕНО: логические часы
} }
// Добавляем статистику пула воркеров, если доступен // Добавляем статистику пула воркеров, если доступен
@@ -1321,6 +1390,10 @@ func generateReplicationID() string {
// 4. Ожидание завершения всех репликаций с таймаутом // 4. Ожидание завершения всех репликаций с таймаутом
// 5. Возврат результата (успех/частичный успех/ошибка) // 5. Возврат результата (успех/частичный успех/ошибка)
// //
// ИСПРАВЛЕНО: передаёт логическую метку Лампорта и версию документа,
// что позволяет удалённым узлам корректно разрешать конфликты
// last-write-wins даже при рассинхронизации физических часов.
//
// Параметры: // Параметры:
// - database: имя базы данных // - database: имя базы данных
// - collection: имя коллекции // - collection: имя коллекции
@@ -1344,6 +1417,12 @@ func (n *Node) ReplicateDocument(database, collection string, doc *storage.Docum
return nil // Только текущий узел - репликация не требуется return nil // Только текущий узел - репликация не требуется
} }
// ИСПРАВЛЕНО: получаем текущую логическую метку и версию документа.
// Логическая метка гарантирует монотонность версий независимо от
// физических часов узлов.
lamportTS := n.lamportClock.Tick()
docVersion := lamportTS
// Формируем данные для репликации // Формируем данные для репликации
repData := struct { repData := struct {
Database string `json:"database"` Database string `json:"database"`
@@ -1351,12 +1430,16 @@ func (n *Node) ReplicateDocument(database, collection string, doc *storage.Docum
Document map[string]interface{} `json:"document"` Document map[string]interface{} `json:"document"`
SourceNode string `json:"source_node"` SourceNode string `json:"source_node"`
ReplicaID string `json:"replica_id"` ReplicaID string `json:"replica_id"`
LamportTS uint64 `json:"lamport_ts"`
DocVersion uint64 `json:"doc_version"`
}{ }{
Database: database, Database: database,
Collection: collection, Collection: collection,
Document: doc.GetFields(), Document: doc.GetFields(),
SourceNode: n.ID, SourceNode: n.ID,
ReplicaID: generateReplicationID(), ReplicaID: generateReplicationID(),
LamportTS: lamportTS,
DocVersion: docVersion,
} }
data, err := json.Marshal(repData) data, err := json.Marshal(repData)
@@ -1403,7 +1486,8 @@ func (n *Node) ReplicateDocument(database, collection string, doc *storage.Docum
successCount.Add(1) successCount.Add(1)
if n.logger != nil { if n.logger != nil {
n.logger.Debug(fmt.Sprintf("Successfully replicated document %s to node %s", doc.ID, targetNodeID)) n.logger.Debug(fmt.Sprintf("Successfully replicated document %s to node %s (version=%d)",
doc.ID, targetNodeID, docVersion))
} }
return nil return nil
}) })
@@ -1429,8 +1513,8 @@ func (n *Node) ReplicateDocument(database, collection string, doc *storage.Docum
case <-done: case <-done:
duration := time.Now().UnixMilli() - startTime duration := time.Now().UnixMilli() - startTime
if n.logger != nil { if n.logger != nil {
n.logger.Info(fmt.Sprintf("Replicated document %s to %d/%d nodes (took %d ms)", n.logger.Info(fmt.Sprintf("Replicated document %s to %d/%d nodes (took %d ms, version=%d)",
doc.ID, successCount.Load(), len(nodes)-1, duration)) doc.ID, successCount.Load(), len(nodes)-1, duration, docVersion))
} }
case <-time.After(30 * time.Second): case <-time.After(30 * time.Second):
if n.logger != nil { if n.logger != nil {
@@ -1548,4 +1632,4 @@ func (n *Node) Stop() {
if n.logger != nil { if n.logger != nil {
n.logger.Info(fmt.Sprintf("Node %s stopped at %s", n.ID, time.UnixMilli(n.stoppedAt).Format("2006-01-02 15:04:05.000"))) n.logger.Info(fmt.Sprintf("Node %s stopped at %s", n.ID, time.UnixMilli(n.stoppedAt).Format("2006-01-02 15:04:05.000")))
} }
} }