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