Browse Source

add NTP time synchronization module

topo_upd
Evgeny 3 months ago
parent
commit
1a0229fd4f
  1. 2
      src/Makefile.am
  2. 26
      src/config_parser.c
  3. 6
      src/config_parser.h
  4. 298
      src/ntp_time.c
  5. 40
      src/ntp_time.h
  6. 10
      src/utun_instance.c
  7. 4
      src/utun_instance.h
  8. 1
      tools/chatgui/libutun/CMakeLists.txt

2
src/Makefile.am

@ -44,6 +44,7 @@ utun_CORE_SOURCES = \
eim_nat.c \
nat_transport.c \
dummynet.c \
ntp_time.c \
proxy/tcp_proxy_client.c \
etcp_router.c \
proxy/tcp_proxy_server.c \
@ -97,6 +98,7 @@ libutun_a_SOURCES = \
eim_nat.c \
nat_transport.c \
dummynet.c \
ntp_time.c \
proxy/tcp_proxy_client.c \
etcp_router.c \
proxy/tcp_proxy_server.c \

26
src/config_parser.c

@ -48,7 +48,8 @@ typedef enum {
SECTION_TCP_PROXY_CLIENT,
SECTION_TCP_PROXY_SERVER,
SECTION_MSG_TRANSPORT,
SECTION_NETWORK
SECTION_NETWORK,
SECTION_NTP
} section_type_t;
static char* trim(char *str) {
@ -763,6 +764,7 @@ static section_type_t parse_section_header(const char *line, char *name, size_t
if (strcasecmp(section, "tcp_proxy_client") == 0) return SECTION_TCP_PROXY_CLIENT;
if (strcasecmp(section, "tcp_proxy_server") == 0) return SECTION_TCP_PROXY_SERVER;
if (strcasecmp(section, "msg_transport") == 0) return SECTION_MSG_TRANSPORT;
if (strcasecmp(section, "ntp") == 0) return SECTION_NTP;
char *colon = strchr(section, ':');
if (!colon) return SECTION_UNKNOWN;
@ -802,6 +804,9 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename)
cfg->global.control_allow_count = 0;
cfg->global.db_sync_enabled = 0;
cfg->global.db_sync_ttl = 86400;
cfg->global.ntp_enabled = 0;
cfg->global.ntp_server_count = 0;
cfg->global.ntp_resync_interval = 3600;
section_type_t cur_section = SECTION_UNKNOWN;
struct CFG_SERVER *cur_server = NULL;
@ -972,6 +977,25 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename)
// parse_network already printed the error
}
break;
case SECTION_NTP:
if (strcmp(key, "enabled") == 0) {
cfg->global.ntp_enabled = strcasecmp(value, "yes") == 0 || strcasecmp(value, "1") == 0 || strcasecmp(value, "true") == 0;
} else if (strcmp(key, "server") == 0) {
if (cfg->global.ntp_server_count >= 5) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Too many NTP servers (max 5)", filename, line_num);
} else {
int idx = cfg->global.ntp_server_count;
strncpy(cfg->global.ntp_servers[idx], value, sizeof(cfg->global.ntp_servers[idx]) - 1);
cfg->global.ntp_servers[idx][sizeof(cfg->global.ntp_servers[idx]) - 1] = '\0';
cfg->global.ntp_server_count++;
}
} else if (strcmp(key, "interval") == 0) {
cfg->global.ntp_resync_interval = atoi(value);
if (cfg->global.ntp_resync_interval < 60) cfg->global.ntp_resync_interval = 60;
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Unknown ntp option '%s'. Valid: enabled, server, interval", filename, line_num, key);
}
break;
default:
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "%s:%d: Key outside section: %s", filename, line_num, key);
break;

6
src/config_parser.h

@ -182,6 +182,12 @@ struct global_config {
// TCP proxy server (exit node) configuration ([tcp_proxy_server] section)
int tcp_proxy_server_enabled;
int tcp_recv_buf; // SO_RCVBUF для TCP сокетов (0=default ОС)
// NTP configuration ([ntp] section)
int ntp_enabled;
char ntp_servers[5][256]; // up to 5 NTP server hostnames
int ntp_server_count;
int ntp_resync_interval; // seconds, default 3600
};
struct utun_config {

298
src/ntp_time.c

@ -0,0 +1,298 @@
#include "ntp_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>
#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_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;
}
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()));
freeaddrinfo(result);
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);
struct sockaddr_in* sin = (struct sockaddr_in*)result->ai_addr;
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,
(struct sockaddr*)sin, sizeof(*sin));
if (sent < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: sendto(%s) failed: %s", server, 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, tries);
socket_close_wrapper(sock);
freeaddrinfo(result);
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, socket_strerror(socket_get_error()));
socket_close_wrapper(sock);
freeaddrinfo(result);
return -1;
}
struct timeval t4_tv;
#ifdef _WIN32
utun_gettimeofday(&t4_tv, NULL);
#else
gettimeofday(&t4_tv, NULL);
#endif
socket_close_wrapper(sock);
freeaddrinfo(result);
if (n < (ssize_t)sizeof(reply)) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: short reply from %s (got %zd, expected %zu)", server, 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);
return -1;
}
int stratum = reply.stratum;
if (stratum == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP: kiss-o-death from %s", server);
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);
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, (long long)offset_us, (long long)rtt_us, stratum);
*offset_us_out = offset_us;
*stratum_out = stratum;
return 0;
}
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;
}
const char* server = ntp->servers[ntp->server_current];
int64_t offset_us = 0;
int stratum = 0;
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;
} 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();
break;
}
}
}
if (ntp->synced) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: time corrected, offset=%lldus",
(long long)ntp->offset_us);
} 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->server_count = g->ntp_server_count;
ntp->server_current = 0;
ntp->resync_interval_sec = g->ntp_resync_interval;
for (int i = 0; i < g->ntp_server_count && i < NTP_MAX_SERVERS; i++) {
strncpy(ntp->servers[i], g->ntp_servers[i], sizeof(ntp->servers[i]) - 1);
ntp->servers[i][sizeof(ntp->servers[i]) - 1] = '\0';
}
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;
}
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;
}

40
src/ntp_time.h

@ -0,0 +1,40 @@
#ifndef NTP_TIME_H
#define NTP_TIME_H
#include <stdint.h>
#include <time.h>
#ifdef __cplusplus
extern "C" {
#endif
struct UTUN_INSTANCE;
#define NTP_MAX_SERVERS 5
#define NTP_DEFAULT_INTERVAL_SEC 3600
#define NTP_FIRST_SYNC_DELAY_SEC 0
struct NTP_TIME {
int enabled;
int synced;
int64_t offset_us; // local_time - ntp_time (positive = local clock ahead)
uint64_t last_sync_tb; // monotonic time of last successful sync (0.1ms units)
void* timer; // uasync timer handle
char servers[NTP_MAX_SERVERS][256];
int server_count;
int server_current;
int resync_interval_sec;
};
int ntp_time_init(struct UTUN_INSTANCE* instance);
void ntp_time_destroy(struct UTUN_INSTANCE* instance);
int64_t ntp_time_get_us(struct UTUN_INSTANCE* instance);
time_t ntp_time_get_seconds(struct UTUN_INSTANCE* instance);
int ntp_time_is_synced(struct UTUN_INSTANCE* instance);
#ifdef __cplusplus
}
#endif
#endif // NTP_TIME_H

10
src/utun_instance.c

@ -369,6 +369,9 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
// Stop running if not already
instance->running = 0;
// Cancel NTP timer
ntp_time_destroy(instance);
// Shutdown message transport server
if (instance->msg_t) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Shutting down message transport server");
@ -608,7 +611,12 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
// Start the main loop
instance->running = 1;
// Initialize NTP time sync (non-fatal if fails)
if (ntp_time_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP time sync init failed (non-fatal)");
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully");
return 0;
}

4
src/utun_instance.h

@ -20,6 +20,7 @@ extern "C" {
#include "proxy/tcp_proxy_client.h"
#include "etcp_router.h"
#include "proxy/tcp_proxy_server.h"
#include "ntp_time.h"
#include "etcp_api.h"
#include "config_parser.h"
@ -145,6 +146,9 @@ struct UTUN_INSTANCE {
// Pending background connections (etcp_connect API)
struct ETCP_CONNECT* pending_connects;
uint32_t etcp_connect_timeout_tb; // Initial timeout in 0.1ms units, default 20000 (2s)
// NTP time synchronization
struct NTP_TIME ntp;
};
// Functions

1
tools/chatgui/libutun/CMakeLists.txt

@ -92,6 +92,7 @@ set(UTUN_COMMON_SOURCES
${SRC_DIR}/eim_nat.c
${SRC_DIR}/nat_transport.c
${SRC_DIR}/dummynet.c
${SRC_DIR}/ntp_time.c
${SRC_DIR}/etcp_router.c
${SRC_DIR}/proxy/tcp_proxy_client.c
${SRC_DIR}/proxy/tcp_proxy_server.c

Loading…
Cancel
Save