summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-07 14:30:40 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-07 14:30:40 -0400
commit36f09a0c35657f5c7379624c8a2463aafb953bd4 (patch)
tree7255f0cf39b2dbe680533277a5cfcf13d241453c /src
parent0d8f4879eedadc6cca1cf454c3e4535359a487a7 (diff)
downloadlibnaunet-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.c273
-rw-r--r--src/rpc2.h27
-rw-r--r--src/util/timer/timer.h2
3 files changed, 272 insertions, 30 deletions
diff --git a/src/rpc.c b/src/rpc.c
index 7938341..c48a0e3 100644
--- a/src/rpc.c
+++ b/src/rpc.c
@@ -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)
diff --git a/src/rpc2.h b/src/rpc2.h
index bed8822..52e14ec 100644
--- a/src/rpc2.h
+++ b/src/rpc2.h
@@ -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);