diff options
| author | 2026-09-07 14:30:40 -0400 | |
|---|---|---|
| committer | 2026-09-07 14:30:40 -0400 | |
| commit | 36f09a0c35657f5c7379624c8a2463aafb953bd4 (patch) | |
| tree | 7255f0cf39b2dbe680533277a5cfcf13d241453c /src | |
| parent | 0d8f4879eedadc6cca1cf454c3e4535359a487a7 (diff) | |
| download | libnaunet-36f09a0c35657f5c7379624c8a2463aafb953bd4.tar.gz libnaunet-36f09a0c35657f5c7379624c8a2463aafb953bd4.tar.bz2 libnaunet-36f09a0c35657f5c7379624c8a2463aafb953bd4.zip | |
Add function to query rpc clients' clock skew
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
| -rw-r--r-- | src/rpc.c | 273 | ||||
| -rw-r--r-- | src/rpc2.h | 27 | ||||
| -rw-r--r-- | src/util/timer/timer.h | 2 |
3 files changed, 272 insertions, 30 deletions
@@ -1,32 +1,204 @@ +#define AL_LOG_SECTION "rpc" +//#define AL_LOG_ENABLE_TRACE +#include <al/log.h> + +#include "util/timer/timer.h" + #include "rpc2.h" +#define HEADER_ID(u) ((u32)((u) & 0x7fffff)) +#define HEADER_OP(u) ((s8)(((u) >> 24) & 0xff)) + +#define SKEW_REPEAT 14u +AL_STATIC_ASSERT(skew_repeat, SKEW_REPEAT, >=, 8u); +#define SKEW_Q1 ((SKEW_REPEAT / 4) - 1) +#define SKEW_Q3 (((SKEW_REPEAT + 3) / 4) * 3 - 1) // Round up + +#define rpc_send_packet(conn, packet) do { \ + conn->outgoing++; \ + nn_packet_stream_send_packet(conn->stream, packet); \ +} while (0) + +// Internal OPs. +enum { + RPC_SKEW_RESPONSE = -5, + RPC_SKEW_REQUEST = -4, + RPC_SKEW_BEGIN = -3, + RPC_SKEW_INIT_QUERY = -2, + RPC_RESPONSE = -1 +}; + +static inline void write_rpc_header(struct nn_packet *packet, s8 op, u32 id) +{ + al_assert(id <= 0x7fffff); + nn_packet_write_u32(packet, (op << 24) | id); +} + bool nn_rpc_init(struct nn_rpc *rpc, struct nn_event_loop *loop, void (*connection_callback)(void *, struct nn_rpc_connection *), + void (*ready_callback)(void *, struct nn_rpc_connection *), void (*connection_closed_callback)(void *, struct nn_rpc_connection *), void *userdata) { rpc->loop = loop; rpc->increment = 0; al_array_init(rpc->commands); + rpc->do_query_skew = false; rpc->conn = NULL; al_array_init(rpc->connections); rpc->connection_callback = connection_callback; + rpc->ready_callback = ready_callback; rpc->connection_closed_callback = connection_closed_callback; rpc->userdata = userdata; return true; } +void nn_rpc_query_clock_skews(struct nn_rpc *rpc, bool do_query_skew) +{ + rpc->do_query_skew = do_query_skew; +} + void nn_rpc_add_command(struct nn_rpc *rpc, struct nn_rpc_command *command) { al_array_push(rpc->commands, *command); } +static s32 skew_compare(const void *a, const void *b) +{ + struct nn_skew *sa = (struct nn_skew *)a; + struct nn_skew *sb = (struct nn_skew *)b; + if (sa->ts > sb->ts) { + return (sa->direction == NNWT_SKEW_NEGATIVE) ? -1 : 1; + } else if (sb->ts > sa->ts) { + return (sb->direction == NNWT_SKEW_NEGATIVE) ? 1 : -1; + } else { + return (sb->direction == NNWT_SKEW_NEGATIVE) - (sa->direction == NNWT_SKEW_NEGATIVE); + } +} + +// ~10-100ms of error is expected. +// https://www.eecis.udel.edu/~mills/exec.html Under 3. Precision and Accuracy. +static void filter_skew_results(struct nn_rpc_connection *conn, struct nn_skew *out) +{ + al_array_sort(conn->skew.results, struct nn_skew, skew_compare); + al_assert(SKEW_Q1 < conn->skew.results.count && SKEW_Q3 < conn->skew.results.count); + struct nn_skew q1 = al_array_at(conn->skew.results, SKEW_Q1); + s64 q1s = q1.ts; + if (q1.direction == NNWT_SKEW_NEGATIVE) { + q1s = -q1s; + } + struct nn_skew q3 = al_array_at(conn->skew.results, SKEW_Q3); + s64 q3s = q3.ts; + if (q3.direction == NNWT_SKEW_NEGATIVE) { + q3s = -q3s; + } + struct nn_skew *result; + al_array_foreach_ptr(conn->skew.results, i, result) { + s64 n = result->ts; + if (result->direction == NNWT_SKEW_NEGATIVE) { + n = -n; + } + if ((i < SKEW_Q1 && n < q1s - llabs(q1s)) || + (i > SKEW_Q3 && n > q3s + llabs(q3s))) { + al_array_remove_at_iter(conn->skew.results, i); + log_trace("skew_outlier: (%lld)", n); + } else { + log_trace("skew_result(%u): %lld", i, n); + } + } + al_assert(conn->skew.results.count > 0); + s64 avg = 0, rem = 0; + u32 count = conn->skew.results.count; + al_array_foreach_ptr_rev(conn->skew.results, i, result) { + al_assert(result->ts <= INT64_MAX); + s64 n = (s64)result->ts; + if (result->direction == NNWT_SKEW_NEGATIVE) { + n = -n; + } + avg += n / count; + rem += n % count; + } + avg += rem / count; + if (avg < 0) { + out->ts = (u64)(-avg); + out->direction = NNWT_SKEW_NEGATIVE; + } else { + out->ts = (u64)avg; + out->direction = NNWT_SKEW_POSITIVE; + } + log_debug("skew: %s%llu", (out->direction == NNWT_SKEW_POSITIVE) ? "" : "-", out->ts); +} + +static void query_clock_skew(struct nn_rpc_connection *conn) +{ + al_assert(conn->rpc->do_query_skew); + struct nn_packet *packet = nn_packet_create(); + write_rpc_header(packet, RPC_SKEW_REQUEST, 0); + conn->skew.query_active = RPC_SKEW_REQUEST; + rpc_send_packet(conn, packet); +} + static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; al_assert(conn->stream == stream); - u32 id = nn_packet_read_u32(packet); - s8 op = nn_packet_read_s8(packet); - if (op == -1) { // Response + u32 u = nn_packet_read_u32(packet); + s8 op = HEADER_OP(u); + u32 id = HEADER_ID(u); + switch (op) { + case RPC_SKEW_RESPONSE: { + u64 now = nn_get_timestamp(); + u64 skew_ts = conn->skew.s.ts; + conn->skew.s.ts = NNWT_TIMESTAMP_INVALID; + al_assert(skew_ts != NNWT_TIMESTAMP_INVALID && skew_ts < now); + u64 ts = nn_packet_read_u64(packet); + ts += (now - skew_ts) / 2; + struct nn_skew result; + if (ts >= now) { + result.ts = ts - now; + result.direction = NNWT_SKEW_POSITIVE; + } else { + result.ts = now - ts; + result.direction = NNWT_SKEW_NEGATIVE; + } + al_array_push(conn->skew.results, result); + al_assert(conn->skew.repeat > 0); + if (!--conn->skew.repeat) { + filter_skew_results(conn, &conn->skew.s); + al_array_free(conn->skew.results); + conn->rpc->ready_callback(conn->rpc->userdata, conn); + } else { + query_clock_skew(conn); + } + break; + } + case RPC_SKEW_REQUEST: { + struct nn_packet *rpacket = nn_packet_create(); + write_rpc_header(rpacket, RPC_SKEW_RESPONSE, 0); + conn->skew.s.ts = nn_get_timestamp(); + conn->skew.query_active = RPC_SKEW_RESPONSE; + rpc_send_packet(conn, rpacket); + al_assert(conn->skew.repeat > 0); + if (!--conn->skew.repeat) { + conn->rpc->ready_callback(conn->rpc->userdata, conn); + } + break; + } + case RPC_SKEW_BEGIN: { + query_clock_skew(conn); + break; + } + case RPC_SKEW_INIT_QUERY: { + conn->skew.repeat = nn_packet_read_u32(packet); + if (!conn->skew.repeat) { + conn->rpc->ready_callback(conn->rpc->userdata, conn); + } else { + struct nn_packet *rpacket = nn_packet_create(); + write_rpc_header(rpacket, RPC_SKEW_BEGIN, 0); + rpc_send_packet(conn, rpacket); + } + break; + } + case RPC_RESPONSE: { struct nn_rpc_callback *callback; al_array_foreach_ptr(conn->callbacks, i, callback) { if (callback->id == id) { @@ -35,33 +207,51 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str return; } } - } else { // Command + break; + } + default: { // User command. struct nn_rpc_command *command; al_array_foreach_ptr(conn->rpc->commands, i, command) { if (command->op == op) { struct nn_packet *rpacket = nn_packet_create(); - nn_packet_write_u32(rpacket, id); - nn_packet_write_s8(rpacket, -1); + write_rpc_header(rpacket, -1, id); if (command->callback(command->userdata, conn, packet, rpacket)) { - conn->outgoing++; - nn_packet_stream_send_packet(stream, rpacket); + rpc_send_packet(conn, rpacket); } else { nn_packet_free(rpacket); } return; } } + break; + } } nn_packet_stream_return_packet(stream, packet); } +static void packet_dequeued_callback(void *userdata, struct nn_packet *packet) +{ + struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; + if (conn->skew.query_active == RPC_SKEW_RESPONSE) { + u64 now = nn_get_timestamp(); + now -= ((now - conn->skew.s.ts) + 1) / 2; // Round up because we round down server-side. + nn_packet_write_u64(packet, now); + conn->skew.s.ts = NNWT_TIMESTAMP_INVALID; + conn->skew.query_active = 0; + } + nn_packet_write_size(packet); +} + static void packet_sent_callback(void *userdata, struct nn_packet *packet) { struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; nn_packet_free(packet); + if (conn->skew.query_active == RPC_SKEW_REQUEST) { + conn->skew.s.ts = nn_get_timestamp(); + conn->skew.query_active = 0; + } al_assert(conn->outgoing > 0); - conn->outgoing--; - if (conn->flushing && conn->outgoing == 0) { + if (!--conn->outgoing && conn->flushing) { nn_packet_stream_disconnect(conn->stream); } } @@ -71,14 +261,40 @@ static bool stream_connection_callback(void *userdata, struct nn_packet_stream * struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; struct nn_rpc *rpc = conn->rpc; if (stream->sock.type == NNWT_SOCKET_TCP) { - nn_packet_stream_set_nodelay(stream, 1); + nn_socket_set_nodelay(&stream->sock, 1); + struct nn_socket *sock = &stream->sock; +#ifdef AL_LOG_ENABLE_TRACE + u32 sndbuf = nn_socket_get_send_buf(sock); + u32 rcvbuf = nn_socket_get_recv_buf(sock); +#endif + nn_socket_set_send_buf(sock, KB(8)); + nn_socket_set_recv_buf(sock, KB(8)); +#ifdef AL_LOG_ENABLE_TRACE + // These log_trace()'s can influence the result of the first skew query (at least on Wine). + log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock)); + log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock)); +#endif } stream->packet_callback = packet_callback; + stream->packet_dequeued_callback = packet_dequeued_callback; stream->packet_sent_callback = packet_sent_callback; stream->userdata = conn; rpc->connection_callback(rpc->userdata, conn); - // Only add connections when acting as a server. - if (!rpc->conn) al_array_push(rpc->connections, conn); + conn->skew.repeat = 0; + if (!rpc->conn) { // Acting as a server. + al_array_push(rpc->connections, conn); + if (conn->rpc->do_query_skew) { + conn->skew.repeat = SKEW_REPEAT; + al_array_init(conn->skew.results); + } else { + rpc->ready_callback(rpc->userdata, conn); + } + struct nn_packet *packet = nn_packet_create(); + write_rpc_header(packet, RPC_SKEW_INIT_QUERY, 0); + // Repeat of 0 means there will be no query. + nn_packet_write_u32(packet, conn->skew.repeat); + rpc_send_packet(conn, packet); + } return true; } @@ -95,6 +311,8 @@ static inline void init_rpc_connection(struct nn_rpc *rpc, struct nn_rpc_connect al_array_init(conn->callbacks); conn->outgoing = 0; conn->flushing = false; + conn->skew.s.ts = NNWT_TIMESTAMP_INVALID; + conn->skew.query_active = 0; } void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream) @@ -107,7 +325,7 @@ void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream) conn->stream->userdata = conn; } -struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc) +void nn_rpc_prepare_client(struct nn_rpc *rpc) { al_assert(!rpc->conn); struct nn_rpc_connection *conn = al_alloc_object(struct nn_rpc_connection); @@ -116,28 +334,30 @@ struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc) nn_packet_stream_init(conn->stream, stream_connection_callback, stream_connection_closed_callback, conn); rpc->conn = conn; al_array_push(rpc->connections, conn); - return rpc->conn; } -void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port) +bool nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port) { al_assert(rpc->conn); - nn_packet_stream_connect(rpc->conn->stream, rpc->loop, id, type, addr, port); + return nn_packet_stream_connect(rpc->conn->stream, rpc->loop, id, type, addr, port); } -struct nn_rpc_connection *nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) +bool nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) { al_assert(rpc->conn); - nn_packet_stream_reconnect(rpc->conn->stream, addr, port); - return rpc->conn; + return nn_packet_stream_reconnect(rpc->conn->stream, addr, port); +} + +void nn_rpc_disconnect(struct nn_rpc *rpc) +{ + al_assert(rpc->conn); + nn_packet_stream_disconnect(rpc->conn->stream); } -// @TODO: Pack opcode into a u32, reduce ID by 8 bits (make struct with bitmask) struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op) { struct nn_packet *packet = nn_packet_create(); - nn_packet_write_u32(packet, (rpc->increment = al_u32_inc_wrap(rpc->increment))); - nn_packet_write_s8(packet, op); + write_rpc_header(packet, op, (rpc->increment = al_u32_add_wrap(rpc->increment, 1, 0x7fffff))); return packet; } @@ -145,8 +365,8 @@ void nn_rpc_free(struct nn_rpc *rpc) { struct nn_rpc_connection *conn; al_array_foreach(rpc->connections, i, conn) { - nn_packet_stream_free(conn->stream); al_array_free(conn->callbacks); + nn_packet_stream_free(conn->stream); al_free(conn->stream); al_free(conn); } @@ -159,13 +379,12 @@ void nn_rpc_connection_command(struct nn_rpc_connection *conn, struct nn_packet { if (callback) { al_array_push(conn->callbacks, ((struct nn_rpc_callback){ - .id = nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH), + .id = HEADER_ID(nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH)), .callback = callback, .userdata = userdata })); } - conn->outgoing++; - nn_packet_stream_send_packet(conn->stream, packet); + rpc_send_packet(conn, packet); } void nn_rpc_conn_flush(struct nn_rpc_connection *conn) @@ -13,11 +13,27 @@ struct nn_rpc_callback { void *userdata; }; +enum { + NNWT_SKEW_POSITIVE = 0, + NNWT_SKEW_NEGATIVE +}; + +struct nn_skew { + u64 ts; + bool direction; +}; + struct nn_rpc_connection { struct nn_packet_stream *stream; array(struct nn_rpc_callback) callbacks; u32 outgoing; bool flushing; + struct { + struct nn_skew s; + s8 query_active; + u32 repeat; + array(struct nn_skew) results; + } skew; struct nn_rpc *rpc; }; @@ -31,21 +47,26 @@ struct nn_rpc { struct nn_event_loop *loop; u32 increment; array(struct nn_rpc_command) commands; + bool do_query_skew; struct nn_rpc_connection *conn; // client array(struct nn_rpc_connection *) connections; void (*connection_callback)(void *, struct nn_rpc_connection *); + void (*ready_callback)(void *, struct nn_rpc_connection *); void (*connection_closed_callback)(void *, struct nn_rpc_connection *); void *userdata; }; bool nn_rpc_init(struct nn_rpc *rpc, struct nn_event_loop *loop, void (*connection_callback)(void *, struct nn_rpc_connection *), + void (*ready_callback)(void *, struct nn_rpc_connection *), void (*connection_closed_callback)(void *, struct nn_rpc_connection *), void *userdata); +void nn_rpc_query_clock_skews(struct nn_rpc *rpc, bool do_query_skew); void nn_rpc_add_command(struct nn_rpc *rpc, struct nn_rpc_command *command); void nn_rpc_add_stream(struct nn_rpc *rpc, struct nn_packet_stream *stream); -struct nn_rpc_connection *nn_rpc_prepare_client(struct nn_rpc *rpc); -void nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port); - struct nn_rpc_connection *nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port); +void nn_rpc_prepare_client(struct nn_rpc *rpc); +bool nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port); +bool nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port); +void nn_rpc_disconnect(struct nn_rpc *rpc); struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op); void nn_rpc_free(struct nn_rpc *rpc); diff --git a/src/util/timer/timer.h b/src/util/timer/timer.h index cd2daf7..2a53c08 100644 --- a/src/util/timer/timer.h +++ b/src/util/timer/timer.h @@ -2,5 +2,7 @@ #include <al/types.h> +#define NNWT_TIMESTAMP_INVALID ((u64)-1) + u64 nn_get_timestamp(void); f64 nn_get_tick(void); |