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.
 
 
 
 
 
 

352 lines
12 KiB

#include "ntp_time.h"
#include "ntp_node_time.h"
#include "utun_instance.h"
#include "../lib/debug_config.h"
#include "../lib/socket_compat.h"
#include "../lib/u_async.h"
#include "../lib/mem.h"
#include <stdlib.h>
#include <string.h>
#ifdef _WIN32
#include <winsock2.h>
#include <ws2tcpip.h>
#else
#include <sys/select.h>
#include <sys/time.h>
#include <netdb.h>
#include <arpa/inet.h>
#endif
#define NTP_DELTA 2208988800ULL // seconds from 1900 to 1970
#define NTP_PORT 123
#define NTP_TIMEOUT_SEC 2
#define NTP_MAX_TRIES 3
#pragma pack(push, 1)
struct ntp_packet {
uint8_t li_vn_mode; // LI(2) | VN(3) | Mode(3)
uint8_t stratum;
uint8_t poll;
uint8_t precision;
uint32_t root_delay;
uint32_t root_dispersion;
uint32_t ref_id;
uint64_t ref_ts;
uint64_t orig_ts;
uint64_t recv_ts;
uint64_t xmit_ts;
};
#pragma pack(pop)
_Static_assert(sizeof(struct ntp_packet) == 48, "NTP packet must be 48 bytes");
static void ntp_time_sync_cb(void* arg);
static uint64_t timeval_to_ntp(struct timeval *tv) {
uint64_t sec = (uint64_t)(tv->tv_sec + NTP_DELTA);
uint64_t frac = ((uint64_t)tv->tv_usec << 32) / 1000000ULL;
return (sec << 32) | frac;
}
static int64_t ntp64_to_us(uint64_t ntp) {
int64_t sec = (int64_t)(ntp >> 32) - (int64_t)NTP_DELTA;
int64_t frac = ((int64_t)(ntp & 0xFFFFFFFFULL) * 1000000ULL) >> 32;
return sec * 1000000LL + frac;
}
static int ntp_do_query(const struct sockaddr_in* sin, const char* server_name, int64_t* offset_us_out, int* stratum_out) {
socket_t sock = socket_create_udp(AF_INET);
if (sock == SOCKET_INVALID) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: socket() failed: %s", socket_strerror(socket_get_error()));
return -1;
}
struct ntp_packet request;
memset(&request, 0, sizeof(request));
request.li_vn_mode = (0 << 6) | (4 << 3) | 3;
struct timeval t1_tv;
#ifdef _WIN32
utun_gettimeofday(&t1_tv, NULL);
#else
gettimeofday(&t1_tv, NULL);
#endif
uint64_t t1_ntp = timeval_to_ntp(&t1_tv);
request.xmit_ts = htobe64(t1_ntp);
int send_ok = 0;
int tries;
for (tries = 0; tries < NTP_MAX_TRIES; tries++) {
ssize_t sent = sendto(sock, (const char*)&request, sizeof(request), 0,
(const struct sockaddr*)sin, sizeof(*sin));
if (sent < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: sendto(%s) failed: %s", server_name, socket_strerror(socket_get_error()));
break;
}
fd_set fds;
FD_ZERO(&fds);
FD_SET(sock, &fds);
struct timeval tv = {NTP_TIMEOUT_SEC, 0};
int r = select((int)(sock + 1), &fds, NULL, NULL, &tv);
if (r < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: select() failed: %s", socket_strerror(socket_get_error()));
break;
}
if (r > 0) { send_ok = 1; break; }
}
if (!send_ok) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP: no response from %s after %d tries", server_name, tries);
socket_close_wrapper(sock);
return -1;
}
struct ntp_packet reply;
struct sockaddr_in from;
socklen_t from_len = sizeof(from);
ssize_t n = recvfrom(sock, (char*)&reply, sizeof(reply), 0, (struct sockaddr*)&from, &from_len);
if (n < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: recvfrom(%s) failed: %s", server_name, socket_strerror(socket_get_error()));
socket_close_wrapper(sock);
return -1;
}
struct timeval t4_tv;
#ifdef _WIN32
utun_gettimeofday(&t4_tv, NULL);
#else
gettimeofday(&t4_tv, NULL);
#endif
socket_close_wrapper(sock);
if (n < (ssize_t)sizeof(reply)) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: short reply from %s (got %zd, expected %zu)", server_name, n, sizeof(reply));
return -1;
}
uint8_t mode = reply.li_vn_mode & 0x07;
if (mode != 4) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: unexpected mode %d from %s (expected 4)", mode, server_name);
return -1;
}
int stratum = reply.stratum;
if (stratum == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: kiss-o-death from %s", server_name);
return -1;
}
uint64_t reply_orig = be64toh(reply.orig_ts);
if (reply_orig != t1_ntp) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: originate timestamp mismatch from %s", server_name);
return -1;
}
int64_t t1_us = ntp64_to_us(t1_ntp);
int64_t t2_us = ntp64_to_us(be64toh(reply.recv_ts));
int64_t t3_us = ntp64_to_us(be64toh(reply.xmit_ts));
int64_t t4_us = ntp64_to_us(timeval_to_ntp(&t4_tv));
int64_t offset_us = ((t2_us - t1_us) + (t3_us - t4_us)) / 2;
int64_t rtt_us = (t4_us - t1_us) - (t3_us - t2_us);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: synced from %s offset=%lldus rtt=%lldus stratum=%d",
server_name, (long long)offset_us, (long long)rtt_us, stratum);
*offset_us_out = offset_us;
*stratum_out = stratum;
return 0;
}
static int ntp_query_server(const char* server, int64_t* offset_us_out, int* stratum_out) {
struct addrinfo hints, *result = NULL;
memset(&hints, 0, sizeof(hints));
hints.ai_family = AF_INET;
hints.ai_socktype = SOCK_DGRAM;
char port_str[8];
snprintf(port_str, sizeof(port_str), "%d", NTP_PORT);
int gai_err = getaddrinfo(server, port_str, &hints, &result);
if (gai_err != 0 || !result) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: failed to resolve %s: %s", server, gai_strerror(gai_err));
return -1;
}
int ret = ntp_do_query((const struct sockaddr_in*)result->ai_addr, server, offset_us_out, stratum_out);
freeaddrinfo(result);
return ret;
}
static void ntp_time_sync_cb(void* arg) {
struct UTUN_INSTANCE* instance = (struct UTUN_INSTANCE*)arg;
if (!instance) return;
struct NTP_TIME* ntp = &instance->ntp;
ntp->timer = NULL;
if (!ntp->enabled || ntp->server_count == 0) {
ntp->timer = uasync_set_timeout(instance->ua, ntp->resync_interval_sec * 10000,
instance, ntp_time_sync_cb, "ntp_sync");
return;
}
int64_t offset_us = 0;
int stratum = 0;
int synced = 0;
if (ntp->test_addr.sin_port != 0) {
if (ntp_do_query(&ntp->test_addr, "test-server", &offset_us, &stratum) == 0) {
ntp->offset_us = offset_us;
ntp->synced = 1;
ntp->last_sync_tb = get_time_tb();
synced = 1;
}
} else {
const char* server = ntp->servers[ntp->server_current];
if (ntp_query_server(server, &offset_us, &stratum) == 0) {
ntp->offset_us = offset_us;
ntp->synced = 1;
ntp->last_sync_tb = get_time_tb();
ntp->server_current = (ntp->server_current + 1) % ntp->server_count;
synced = 1;
} else {
for (int i = 0; i < ntp->server_count; i++) {
ntp->server_current = (ntp->server_current + 1) % ntp->server_count;
const char* next_server = ntp->servers[ntp->server_current];
if (strcmp(next_server, server) == 0) break;
if (ntp_query_server(next_server, &offset_us, &stratum) == 0) {
ntp->offset_us = offset_us;
ntp->synced = 1;
ntp->last_sync_tb = get_time_tb();
synced = 1;
break;
}
}
}
}
if (ntp->synced) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: time corrected, offset=%lldus",
(long long)ntp->offset_us);
ntp_node_sync_peers(instance);
} else {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP: all servers unreachable, retry in %ds",
ntp->resync_interval_sec);
}
ntp->timer = uasync_set_timeout(instance->ua, ntp->resync_interval_sec * 10000,
instance, ntp_time_sync_cb, "ntp_sync");
}
int ntp_time_init(struct UTUN_INSTANCE* instance) {
if (!instance) return -1;
struct NTP_TIME* ntp = &instance->ntp;
struct global_config* g = &instance->config->global;
ntp->enabled = g->ntp_enabled;
ntp->synced = 0;
ntp->offset_us = 0;
ntp->last_sync_tb = 0;
ntp->timer = NULL;
ntp->servers = NULL;
ntp->server_count = g->ntp_server_count;
ntp->server_current = 0;
ntp->resync_interval_sec = g->ntp_resync_interval;
memset(&ntp->test_addr, 0, sizeof(ntp->test_addr));
if (g->ntp_server_count > 0) {
ntp->servers = u_malloc(g->ntp_server_count * sizeof(char*));
if (!ntp->servers) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "NTP: failed to allocate servers array");
ntp->server_count = 0;
return -1;
}
struct CFG_NTP_SERVER *ns = g->ntp_servers;
int i = 0;
while (ns && i < g->ntp_server_count) {
ntp->servers[i] = u_strdup(ns->name);
if (!ntp->servers[i]) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "NTP: failed to duplicate server name");
while (--i >= 0) u_free(ntp->servers[i]);
u_free(ntp->servers);
ntp->servers = NULL;
ntp->server_count = 0;
return -1;
}
i++;
ns = ns->next;
}
}
if (!ntp->enabled || ntp->server_count == 0) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: disabled or no servers configured");
return 0;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: initialized, %d servers, interval=%ds, first sync immediately",
ntp->server_count, ntp->resync_interval_sec);
ntp->timer = uasync_set_timeout(instance->ua, NTP_FIRST_SYNC_DELAY_SEC * 10000,
instance, ntp_time_sync_cb, "ntp_sync");
if (!ntp->timer) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: failed to start sync timer");
return -1;
}
return 0;
}
void ntp_time_destroy(struct UTUN_INSTANCE* instance) {
if (!instance) return;
struct NTP_TIME* ntp = &instance->ntp;
if (ntp->timer && instance->ua) {
uasync_cancel_timeout(instance->ua, ntp->timer);
ntp->timer = NULL;
}
if (ntp->servers) {
for (int i = 0; i < ntp->server_count; i++) u_free(ntp->servers[i]);
u_free(ntp->servers);
ntp->servers = NULL;
}
ntp->enabled = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: destroyed (synced=%d, offset=%lldus)",
ntp->synced, (long long)ntp->offset_us);
}
int64_t ntp_time_get_us(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->ntp.synced) {
struct timeval tv;
#ifdef _WIN32
utun_gettimeofday(&tv, NULL);
#else
gettimeofday(&tv, NULL);
#endif
return (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec;
}
struct timeval tv;
#ifdef _WIN32
utun_gettimeofday(&tv, NULL);
#else
gettimeofday(&tv, NULL);
#endif
int64_t local_us = (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec;
return local_us - instance->ntp.offset_us;
}
time_t ntp_time_get_seconds(struct UTUN_INSTANCE* instance) {
return (time_t)(ntp_time_get_us(instance) / 1000000LL);
}
int ntp_time_is_synced(struct UTUN_INSTANCE* instance) {
return instance && instance->ntp.synced;
}
void ntp_time_set_test_addr(struct NTP_TIME* ntp, const char* ip, uint16_t port) {
if (!ntp || !ip) return;
memset(&ntp->test_addr, 0, sizeof(ntp->test_addr));
ntp->test_addr.sin_family = AF_INET;
ntp->test_addr.sin_port = htons(port);
if (inet_pton(AF_INET, ip, &ntp->test_addr.sin_addr) != 1) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: test addr parse failed for %s", ip);
memset(&ntp->test_addr, 0, sizeof(ntp->test_addr));
}
}