You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
207 lines
7.4 KiB
207 lines
7.4 KiB
#include "ntp_node_time.h" |
|
#include "ntp_time.h" |
|
#include "utun_instance.h" |
|
#include "../lib/debug_config.h" |
|
#include "../lib/mem.h" |
|
#include "etcp_api.h" |
|
#include "etcp.h" |
|
#include "etcp_connections.h" |
|
#include <string.h> |
|
#include <stdlib.h> |
|
#include <sys/time.h> |
|
|
|
#define TIME_SYNC_FLAG_SYNCED 0x01 |
|
|
|
#pragma pack(push, 1) |
|
struct time_sync_msg { |
|
uint8_t flags; // bit0 = sender_synced |
|
int64_t t_sender_us; // corrected time if synced, local gettimeofday otherwise |
|
}; |
|
#pragma pack(pop) |
|
|
|
_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, int event, 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, 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; |
|
|
|
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); |
|
#else |
|
gettimeofday(&tv, NULL); |
|
#endif |
|
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) { |
|
struct timeval tv; |
|
#ifdef _WIN32 |
|
utun_gettimeofday(&tv, NULL); |
|
#else |
|
gettimeofday(&tv, NULL); |
|
#endif |
|
return (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec; |
|
} |
|
|
|
static struct NTP_NODE_PEER* find_peer(struct NTP_NODE_TIME* np, uint64_t node_id) { |
|
for (int i = 0; i < np->peer_count; i++) { |
|
if (np->peers[i].node_id == node_id) return &np->peers[i]; |
|
} |
|
return NULL; |
|
} |
|
|
|
static struct NTP_NODE_PEER* add_peer(struct NTP_NODE_TIME* np, uint64_t node_id) { |
|
if (np->peer_count >= NTP_NODE_MAX_PEERS) return NULL; |
|
struct NTP_NODE_PEER* p = &np->peers[np->peer_count]; |
|
p->node_id = node_id; |
|
p->offset_us = 0; |
|
p->has_offset = 0; |
|
np->peer_count++; |
|
return p; |
|
} |
|
|
|
static void check_drift(uint64_t node_id, int64_t offset_delta_us) { |
|
int64_t abs_us = offset_delta_us < 0 ? -offset_delta_us : offset_delta_us; |
|
if (abs_us > NTP_NODE_DRIFT_ERROR_US) { |
|
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP_node: clock offset changed by >10s with node %012llx: %lldus", |
|
(unsigned long long)node_id, (long long)offset_delta_us); |
|
} else if (abs_us > NTP_NODE_DRIFT_WARN_US) { |
|
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP_node: clock offset changed by >2s with node %012llx: %lldus", |
|
(unsigned long long)node_id, (long long)offset_delta_us); |
|
} |
|
} |
|
|
|
static void ntp_node_on_conn_init(struct ETCP_CONN* conn, int event, void* arg) { (void)event; |
|
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; |
|
if (!inst || !conn) return; |
|
|
|
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)); |
|
} |
|
} |
|
|
|
static void ntp_node_on_new_conn(struct ETCP_CONN* conn, void* arg) { |
|
if (!conn) return; |
|
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; |
|
if (!inst) return; |
|
etcp_conn_add_cbk(conn, ntp_node_on_conn_init, inst, ETCP_CBK_EVENT_INIT); |
|
} |
|
|
|
static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { |
|
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); queue_entry_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 + 1); |
|
int sender_synced = (msg->flags & TIME_SYNC_FLAG_SYNCED) != 0; |
|
|
|
int64_t t4 = gettimeofday_us(); |
|
int64_t peer_offset = t4 - msg->t_sender_us; |
|
|
|
struct NTP_NODE_TIME* np = &inst->ntp_node; |
|
struct NTP_NODE_PEER* peer = find_peer(np, conn->peer_node_id); |
|
if (!peer) peer = add_peer(np, conn->peer_node_id); |
|
if (!peer) return; |
|
|
|
if (peer->has_offset) { |
|
int64_t offset_delta = peer_offset - peer->offset_us; |
|
check_drift(conn->peer_node_id, offset_delta); |
|
} |
|
peer->offset_us = peer_offset; |
|
peer->has_offset = 1; |
|
|
|
int was_unsynced = !ntp_time_is_synced(inst); |
|
|
|
if (was_unsynced && sender_synced) { |
|
int64_t rtt_us = 0, n = 0; |
|
for (struct ETCP_LINK* l = conn->links; l; l = l->next) { rtt_us += l->rtt_last * 100LL; n++; } |
|
inst->ntp.offset_us = t4 - msg->t_sender_us - (n ? rtt_us / n / 2 : 0); |
|
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_entry_free(entry); |
|
|
|
// Ответный TIME_SYNC — всегда |
|
send_time_sync(inst, conn); |
|
|
|
// Если мы только что скорректировались → раздаём коррекцию всем peer'ам |
|
if (was_unsynced && ntp_time_is_synced(inst)) { |
|
ntp_node_sync_peers(inst); |
|
} |
|
} |
|
|
|
void ntp_node_sync_peers(struct UTUN_INSTANCE* inst) { |
|
if (!inst || !ntp_time_is_synced(inst)) return; |
|
|
|
int sent = 0; |
|
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) == 0) sent++; |
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: broadcast TIME_SYNC to %d peers", sent); |
|
} |
|
|
|
int ntp_node_time_init(struct UTUN_INSTANCE* inst) { |
|
if (!inst) return -1; |
|
struct NTP_NODE_TIME* np = &inst->ntp_node; |
|
np->peer_count = 0; |
|
|
|
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 (raw ETCP, ID=0x%02X)", ETCP_ID_NTP_TIME); |
|
return 0; |
|
} |
|
|
|
void ntp_node_time_destroy(struct UTUN_INSTANCE* inst) { |
|
if (!inst) return; |
|
struct NTP_NODE_TIME* np = &inst->ntp_node; |
|
|
|
etcp_unbind(inst, ETCP_ID_NTP_TIME); |
|
etcp_remove_new_conn_cbk(inst, ntp_node_on_new_conn, inst); |
|
|
|
np->peer_count = 0; |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: destroyed"); |
|
}
|
|
|