#define AL_LOG_SECTION "rpc" //#define AL_LOG_ENABLE_TRACE #include #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 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) { callback->callback(callback->userdata, conn, packet); al_array_remove_at(conn->callbacks, i); return; } } 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(); write_rpc_header(rpacket, -1, id); if (command->callback(command->userdata, conn, packet, 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); if (!--conn->outgoing && conn->flushing) { nn_packet_stream_disconnect(conn->stream); } } static bool stream_connection_callback(void *userdata, struct nn_packet_stream *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_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); 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; } static void stream_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata; al_assert(stream == conn->stream); conn->rpc->connection_closed_callback(conn->rpc->userdata, conn); } static inline void init_rpc_connection(struct nn_rpc *rpc, struct nn_rpc_connection *conn) { conn->rpc = rpc; 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) { struct nn_rpc_connection *conn = al_alloc_object(struct nn_rpc_connection); init_rpc_connection(rpc, conn); conn->stream = stream; conn->stream->connection_callback = stream_connection_callback; conn->stream->connection_closed_callback = stream_connection_closed_callback; conn->stream->userdata = conn; } 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); init_rpc_connection(rpc, conn); conn->stream = al_alloc_object(struct nn_packet_stream); nn_packet_stream_init(conn->stream, stream_connection_callback, stream_connection_closed_callback, conn); rpc->conn = conn; al_array_push(rpc->connections, conn); } bool nn_rpc_connect(struct nn_rpc *rpc, u8 id, u8 type, str *addr, u16 port) { al_assert(rpc->conn); return nn_packet_stream_connect(rpc->conn->stream, rpc->loop, id, type, addr, port); } bool nn_rpc_reconnect(struct nn_rpc *rpc, str *addr, u16 port) { al_assert(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); } struct nn_packet *nn_rpc_get_packet(struct nn_rpc *rpc, s8 op) { struct nn_packet *packet = nn_packet_create(); write_rpc_header(packet, op, (rpc->increment = al_u32_add_wrap(rpc->increment, 1, 0x7fffff))); return packet; } void nn_rpc_free(struct nn_rpc *rpc) { struct nn_rpc_connection *conn; al_array_foreach(rpc->connections, i, conn) { al_array_free(conn->callbacks); nn_packet_stream_free(conn->stream); al_free(conn->stream); al_free(conn); } al_array_free(rpc->connections); al_array_free(rpc->commands); } void nn_rpc_connection_command(struct nn_rpc_connection *conn, struct nn_packet *packet, void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata) { if (callback) { al_array_push(conn->callbacks, ((struct nn_rpc_callback){ .id = HEADER_ID(nn_packet_get_u32(packet, NNWT_PACKET_HEADER_LENGTH)), .callback = callback, .userdata = userdata })); } rpc_send_packet(conn, packet); } void nn_rpc_conn_flush(struct nn_rpc_connection *conn) { if (conn->outgoing != 0) { conn->flushing = true; } else { nn_packet_stream_disconnect(conn->stream); } } void nn_rpc_conn_disconnect(struct nn_rpc_connection *conn) { nn_packet_stream_disconnect(conn->stream); }