diff --git a/src/etcp_api.h b/src/etcp_api.h index f20af532..32281049 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -32,6 +32,7 @@ extern "C" { #define ETCP_ID_ROUTE_ENTRY 0x01 // Элемент роутинг-таблицы (BGP) #define ETCP_ID_NAT_DETECTION 0x02 // NAT-детекция (STUN-like ping через третий узел) #define ETCP_ID_SVC_ROUTE 0x03 // Транспорт роутера (etcp_router) +#define ETCP_ID_NTP_TIME 0x04 // NTP time sync P2P // === Router-сервисы (etcp_router_bind/etcp_route_send) === #define ETCP_RT_ID_DATA 0x00 // routing.c — маршрутизация данных @@ -42,7 +43,6 @@ extern "C" { #define ETCP_RT_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit) #define ETCP_RT_ID_MSG_TRANSPORT 0x10 // msg_transport — локальный IPC транспорт сообщений #define ETCP_RT_ID_CONN_MGR 0x11 // Connection Manager — management connections -#define ETCP_RT_ID_NTP_TIME 0x12 // NTP time sync между узлами // Connection status events (instance-level callback) #define ETCP_CONN_STATUS_NEW 0 // соединение создано diff --git a/src/etcp_api_doc.md b/src/etcp_api_doc.md index 62047ded..c807dd25 100644 --- a/src/etcp_api_doc.md +++ b/src/etcp_api_doc.md @@ -65,7 +65,7 @@ etcp_set_new_conn_cbk(instance, on_new_conn, my_data); - `ETCP_RT_ID_ICMP_PROXY` (0x06) — ICMP-прокси - `ETCP_RT_ID_MSG_TRANSPORT` (0x10) — локальный IPC-транспорт сообщений - `ETCP_RT_ID_CONN_MGR` (0x11) — управление соединениями -- `ETCP_RT_ID_NTP_TIME` (0x12) — синхронизация времени +- `ETCP_ID_NTP_TIME` (0x04) — синхронизация времени (raw ETCP, P2P) ### Фоновые соединения (etcp_connect) diff --git a/src/ntp_node_time.c b/src/ntp_node_time.c index 621e3592..f438b2b7 100644 --- a/src/ntp_node_time.c +++ b/src/ntp_node_time.c @@ -4,7 +4,6 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" #include "etcp_api.h" -#include "etcp_router.h" #include "etcp.h" #include #include @@ -21,28 +20,48 @@ struct time_sync_msg { _Static_assert(sizeof(struct time_sync_msg) == 9, "time_sync_msg size mismatch"); +#define TIME_SYNC_PKT_SIZE (1 + sizeof(struct time_sync_msg)) // cmd(1) + msg(9) = 10 + static void ntp_node_on_conn_init(struct ETCP_CONN* conn, void* arg); static void ntp_node_on_new_conn(struct ETCP_CONN* conn, void* arg); static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); -static int send_time_sync(struct UTUN_INSTANCE* inst, uint64_t dst_node_id) { - struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, ETCP_RT_ID_NTP_TIME); - if (!rconn) return -1; +static int send_time_sync(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn) { + if (!inst || !conn) return -1; struct time_sync_msg msg; msg.flags = ntp_time_is_synced(inst) ? TIME_SYNC_FLAG_SYNCED : 0; - struct timeval tv; + int64_t now_us; + if (ntp_time_is_synced(inst)) { + now_us = ntp_time_get_us(inst); + } else { + struct timeval tv; #ifdef _WIN32 - utun_gettimeofday(&tv, NULL); + utun_gettimeofday(&tv, NULL); #else - gettimeofday(&tv, NULL); + gettimeofday(&tv, NULL); #endif - msg.t_sender_us = ntp_time_is_synced(inst) - ? ntp_time_get_us(inst) - : (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec; - - return etcp_router_conn_send(rconn, (const uint8_t*)&msg, sizeof(msg)); + now_us = (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec; + } + msg.t_sender_us = now_us; + + uint8_t* pkt = u_malloc(TIME_SYNC_PKT_SIZE); + if (!pkt) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: u_malloc failed for TIME_SYNC"); return -1; } + pkt[0] = ETCP_ID_NTP_TIME; + memcpy(pkt + 1, &msg, sizeof(msg)); + + struct ll_entry* e = queue_entry_new(0); + if (!e) { u_free(pkt); DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: queue_entry_new failed for TIME_SYNC"); return -1; } + e->dgram = pkt; + e->len = TIME_SYNC_PKT_SIZE; + + if (etcp_send(conn, e) != 0) { + u_free(pkt); queue_entry_free(e); + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: etcp_send failed to node %012llx", (unsigned long long)conn->peer_node_id); + return -1; + } + return 0; } static int64_t gettimeofday_us(void) { @@ -86,12 +105,9 @@ static void ntp_node_on_conn_init(struct ETCP_CONN* conn, void* arg) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; if (!inst || !conn) return; - if (send_time_sync(inst, conn->peer_node_id) == 0) { + if (send_time_sync(inst, conn) == 0) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: sent TIME_SYNC to node %012llx (synced=%d)", (unsigned long long)conn->peer_node_id, ntp_time_is_synced(inst)); - } else { - DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: failed to send TIME_SYNC to node %012llx", - (unsigned long long)conn->peer_node_id); } } @@ -103,17 +119,14 @@ static void ntp_node_on_new_conn(struct ETCP_CONN* conn, void* arg) { } static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!conn || !entry) return; + if (!conn || !entry) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } struct UTUN_INSTANCE* inst = conn->instance; - if (!inst) { queue_dgram_free(entry); return; } - - uint8_t svc_id = entry->dgram[offsetof(struct SVC_ROUTE_HDR, svc_id)]; - if (svc_id != ETCP_RT_ID_NTP_TIME) { queue_dgram_free(entry); return; } + if (!inst) { queue_dgram_free(entry); queue_entry_free(entry); return; } - size_t payload_len = entry->len < SVC_ROUTE_HDR_SIZE ? 0 : entry->len - SVC_ROUTE_HDR_SIZE; - if (payload_len < sizeof(struct time_sync_msg)) { queue_dgram_free(entry); return; } + if (entry->dgram[0] != ETCP_ID_NTP_TIME) { queue_dgram_free(entry); queue_entry_free(entry); return; } + if (entry->len < TIME_SYNC_PKT_SIZE) { queue_dgram_free(entry); queue_entry_free(entry); return; } - struct time_sync_msg* msg = (struct time_sync_msg*)(entry->dgram + SVC_ROUTE_HDR_SIZE); + struct time_sync_msg* msg = (struct time_sync_msg*)(entry->dgram + 1); int sender_synced = (msg->flags & TIME_SYNC_FLAG_SYNCED) != 0; int64_t t4 = gettimeofday_us(); @@ -129,18 +142,18 @@ static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { int was_unsynced = !ntp_time_is_synced(inst); if (was_unsynced && sender_synced) { - inst->ntp.offset_us = msg->t_sender_us - t4; - inst->ntp.synced = 1; + inst->ntp.offset_us = msg->t_sender_us - t4; + inst->ntp.synced = 1; inst->ntp.last_sync_tb = get_time_tb(); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: synced from node %012llx, offset=%lldus", (unsigned long long)conn->peer_node_id, (long long)inst->ntp.offset_us); } - queue_dgram_free(entry); + queue_dgram_free(entry); queue_entry_free(entry); // Ответный TIME_SYNC — всегда - send_time_sync(inst, conn->peer_node_id); + send_time_sync(inst, conn); // Если мы только что скорректировались → раздаём коррекцию всем peer'ам if (was_unsynced && ntp_time_is_synced(inst)) { @@ -155,7 +168,7 @@ void ntp_node_sync_peers(struct UTUN_INSTANCE* inst) { for (struct ll_entry* entry = inst->connections->head; entry; entry = entry->next) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; if (!ce->conn->initialized || !ce->conn->links_up) continue; - if (send_time_sync(inst, ce->conn->peer_node_id) == 0) sent++; + if (send_time_sync(inst, ce->conn) == 0) sent++; } DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: broadcast TIME_SYNC to %d peers", sent); @@ -166,10 +179,10 @@ int ntp_node_time_init(struct UTUN_INSTANCE* inst) { struct NTP_NODE_TIME* np = &inst->ntp_node; np->peer_count = 0; - etcp_router_bind(inst, ETCP_RT_ID_NTP_TIME, ntp_node_recv_cb); + etcp_bind(inst, ETCP_ID_NTP_TIME, ntp_node_recv_cb); etcp_add_new_conn_cbk(inst, ntp_node_on_new_conn, inst); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: initialized (RT_ID=0x%02X)", ETCP_RT_ID_NTP_TIME); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: initialized (raw ETCP, ID=0x%02X)", ETCP_ID_NTP_TIME); return 0; } @@ -177,7 +190,7 @@ void ntp_node_time_destroy(struct UTUN_INSTANCE* inst) { if (!inst) return; struct NTP_NODE_TIME* np = &inst->ntp_node; - etcp_router_unbind(inst, ETCP_RT_ID_NTP_TIME); + etcp_unbind(inst, ETCP_ID_NTP_TIME); etcp_remove_new_conn_cbk(inst, ntp_node_on_new_conn, inst); np->peer_count = 0; diff --git a/src/ntp_node_time_doc.md b/src/ntp_node_time_doc.md index 0bb2d820..b84c1962 100644 --- a/src/ntp_node_time_doc.md +++ b/src/ntp_node_time_doc.md @@ -4,7 +4,7 @@ Синхронизация времени между узлами uTun через ETCP-канал. Позволяет вычислять смещение часов относительно каждого соседа, обнаруживать дрейф и, если узел сам не синхронизирован с NTP, корректировать свои часы по времени синхронизированного соседа. Используется для точной координации событий между узлами (db_sync, timestamp-ы сообщений в чате). ## 2. Как пользоваться -`ntp_node_time_init()` регистрирует callback на новые ETCP-соединения и callback на приём данных через `etcp_router` (service ID `ETCP_RT_ID_NTP_TIME = 0x12`). При установке нового соединения автоматически вызывается `ntp_node_on_conn_ready` — отправка `TIME_SYNC` соседу. При получении `TIME_SYNC` вычисляется per-peer offset, проверяется дрейф, и если мы не синхронизированы а сосед синхронизирован — корректируем свои часы и рассылаем коррекцию всем остальным соседям. +`ntp_node_time_init()` регистрирует callback на новые ETCP-соединения и callback на приём данных через raw `etcp_bind` (packet ID `ETCP_ID_NTP_TIME = 0x04`, P2P без роутера). При установке нового соединения автоматически вызывается `ntp_node_on_conn_init` — отправка `TIME_SYNC` соседу. При получении `TIME_SYNC` вычисляется per-peer offset, проверяется дрейф, и если мы не синхронизированы а сосед синхронизирован — корректируем свои часы и рассылаем коррекцию всем остальным соседям. `ntp_node_sync_peers()` — широковещательная рассылка `TIME_SYNC` всем активным соседям. Вызывается автоматически после успешной NTP-синхронизации из `ntp_time`. @@ -20,6 +20,6 @@ - `struct NTP_NODE_TIME` — фиксированный массив `peers[32]` + счётчик `peer_count` **Функции:** -- `ntp_node_time_init(inst)` — регистрирует bind в `etcp_router` (RT_ID 0x12) и callback на новые соединения +- `ntp_node_time_init(inst)` — регистрирует bind в `etcp_bind` (ID 0x04) и callback на новые соединения - `ntp_node_time_destroy(inst)` — снимает регистрацию bind и callback, обнуляет peer_count - `ntp_node_sync_peers(inst)` — отправляет `TIME_SYNC` всем активным/initialized соседям (вызывается после NTP-синхронизации)