From 9970a4a827dbf27c2f8fa7170ad4d137e1bd7e65 Mon Sep 17 00:00:00 2001 From: theMackabu Date: Fri, 1 May 2026 15:08:51 -0700 Subject: [PATCH] use name-based procedure calls --- demos/async/client.c | 2 +- demos/async/server.c | 2 +- demos/bench.c | 103 ++++++++++++++++++++++++++++++++++++--- demos/client.c | 2 +- demos/server.c | 4 +- include/proc.h | 8 +++ include/routes.h | 11 +++-- include/rpc/client.h | 6 +-- include/rpc/protocol.h | 4 +- include/rpc/server.h | 9 ++-- include/scheduler.h | 6 +-- src/client.c | 24 +++++++-- src/protocol.c | 24 ++++++--- src/routes.c | 95 ++++++++++++++++++++++++++++++------ src/scheduler.c | 6 +-- src/server.c | 24 +++++++-- tests/test_integration.c | 12 ++--- tests/test_protocol.c | 2 +- 18 files changed, 273 insertions(+), 71 deletions(-) create mode 100644 include/proc.h diff --git a/demos/async/client.c b/demos/async/client.c index 3b77070..363fe57 100644 --- a/demos/async/client.c +++ b/demos/async/client.c @@ -22,7 +22,7 @@ int main(int argc, char **argv) { rpc_value *result = NULL; size_t result_count = 0; - if (rpc_client_call(client, 3, &args, &result, &result_count) != 0) { + if (rpc_client_call_name(client, "add", &args, &result, &result_count) != 0) { fprintf(stderr, "async RPC error: %s\n", rpc_client_error(client)); rpc_writer_free(&args); rpc_client_close(client); diff --git a/demos/async/server.c b/demos/async/server.c index 3d8fbe0..1215bb9 100644 --- a/demos/async/server.c +++ b/demos/async/server.c @@ -130,7 +130,7 @@ int main(int argc, char **argv) { } if (rpc_server_init(&g_server) != 0 || rpc_server_set_workers(g_server, workers) != 0 || - rpc_server_add_async_route(g_server, 3, async_add, &g_queue) != 0 || rpc_server_bind(g_server, host, port) != 0 || + rpc_server_add_async_route_name(g_server, "add", async_add, &g_queue) != 0 || rpc_server_bind(g_server, host, port) != 0 || rpc_server_listen(g_server) != 0) { fprintf(stderr, "failed to start async RPC server\n"); rpc_server_destroy(g_server); diff --git a/demos/bench.c b/demos/bench.c index 6573f0f..d62af65 100644 --- a/demos/bench.c +++ b/demos/bench.c @@ -9,8 +9,19 @@ #include #include #include +#include #include +#ifdef __APPLE__ +#include +#endif + +typedef struct bench_memory { + uint64_t rss; + uint64_t virtual_size; + uint64_t footprint; +} bench_memory; + typedef struct bench_server { rpc_server *rpc; pthread_t thread; @@ -32,6 +43,7 @@ typedef struct worker_args { uint64_t requests; uint64_t warmup; uint64_t pipeline; + const char *proc_name; uint64_t latency_offset; uint64_t *latencies_ns; bench_gate *gate; @@ -50,6 +62,40 @@ static uint64_t now_ns(void) { return (uint64_t)ts.tv_sec * 1000000000ull + (uint64_t)ts.tv_nsec; } +static uint64_t timeval_ns(struct timeval tv) { + return (uint64_t)tv.tv_sec * 1000000000ull + (uint64_t)tv.tv_usec * 1000ull; +} + +static void format_bytes(uint64_t bytes, char out[16]) { + static const char *units[] = {"B", "KiB", "MiB", "GiB", "TiB"}; + double value = (double)bytes; + size_t unit = 0; + while (value >= 1024.0 && unit + 1 < sizeof(units) / sizeof(*units)) { + value /= 1024.0; + unit++; + } + if (unit == 0) { + snprintf(out, 16, "%" PRIu64 " %s", bytes, units[unit]); + } else { + snprintf(out, 16, "%.2f %s", value, units[unit]); + } +} + +static bench_memory current_memory(void) { + bench_memory memory = {0}; +#ifdef __APPLE__ + task_vm_info_data_t info; + mach_msg_type_number_t count = TASK_VM_INFO_COUNT; + if (task_info(mach_task_self(), TASK_VM_INFO, (task_info_t)&info, &count) == + KERN_SUCCESS) { + memory.rss = info.resident_size; + memory.virtual_size = info.virtual_size; + memory.footprint = info.phys_footprint; + } +#endif + return memory; +} + static int add_handler(rpc_ctx *ctx, const rpc_value *args, size_t argc, rpc_writer *out, void *user_data) { (void)ctx; (void)user_data; @@ -66,7 +112,7 @@ static void *server_main(void *arg) { static int start_server(bench_server *server, char port[16]) { if (rpc_server_init(&server->rpc) != 0 || (server->workers != 0 && rpc_server_set_workers(server->rpc, server->workers) != 0) || - rpc_server_add_route(server->rpc, 1, add_handler, NULL) != 0 || + rpc_server_add_route_name(server->rpc, "add", add_handler, NULL) != 0 || rpc_server_bind(server->rpc, "127.0.0.1", "0") != 0 || rpc_server_listen(server->rpc) != 0) { return -1; } @@ -88,10 +134,10 @@ static void stop_server(bench_server *server) { server->rpc = NULL; } -static int run_one_call(rpc_client *client, rpc_writer *payload, int64_t expected) { +static int run_one_call(rpc_client *client, const char *proc_name, rpc_writer *payload, int64_t expected) { rpc_value *values = NULL; size_t count = 0; - int rc = rpc_client_call(client, 1, payload, &values, &count); + int rc = rpc_client_call_name(client, proc_name, payload, &values, &count); if (rc == 0 && (count != 1 || values[0].type != RPC_TYPE_I64 || values[0].as.i64 != expected)) { rc = -1; } rpc_values_free(values); return rc; @@ -119,7 +165,7 @@ static void *worker_main(void *arg) { } for (uint64_t i = 0; i < worker->warmup; ++i) { - if (run_one_call(client, &payload, 42) != 0) { + if (run_one_call(client, worker->proc_name, &payload, 42) != 0) { worker->failed = 1; goto done; } @@ -134,7 +180,7 @@ static void *worker_main(void *arg) { if (worker->pipeline <= 1) { for (uint64_t i = 0; i < worker->requests; ++i) { uint64_t start = now_ns(); - if (run_one_call(client, &payload, 42) != 0) { + if (run_one_call(client, worker->proc_name, &payload, 42) != 0) { worker->failed = 1; goto done; } @@ -157,7 +203,7 @@ static void *worker_main(void *arg) { while (sent < worker->requests && sent - received < worker->pipeline) { uint64_t start = now_ns(); uint64_t call_id = 0; - if (rpc_client_send_call(client, 1, &payload, &call_id) != 0) { + if (rpc_client_send_call_name(client, worker->proc_name, &payload, &call_id) != 0) { worker->failed = 1; goto pipeline_done; } @@ -190,7 +236,7 @@ static void *worker_main(void *arg) { if (sent < worker->requests) { uint64_t start = now_ns(); uint64_t next_call_id = 0; - if (rpc_client_send_call(client, 1, &payload, &next_call_id) != 0) { + if (rpc_client_send_call_name(client, worker->proc_name, &payload, &next_call_id) != 0) { worker->failed = 1; goto pipeline_done; } @@ -360,6 +406,7 @@ int main(int argc, char **argv) { .requests = count, .warmup = warmup_base + (i < warmup_rem ? 1u : 0u), .pipeline = pipeline, + .proc_name = "add", .latency_offset = offset, .latencies_ns = latencies, .gate = &gate, @@ -377,6 +424,8 @@ int main(int argc, char **argv) { int failed = 0; int bench_started = 0; uint64_t bench_start = 0; + struct rusage usage_start; + memset(&usage_start, 0, sizeof(usage_start)); if (gate_wait_ready(&gate, started) != 0) { fprintf(stderr, "benchmark workers did not become ready\n"); failed = 1; @@ -391,6 +440,7 @@ int main(int argc, char **argv) { } else { rpc_trace_reset(); rpc_trace_set_enabled(trace_enabled); + getrusage(RUSAGE_SELF, &usage_start); bench_start = now_ns(); bench_started = 1; gate_start(&gate); @@ -403,6 +453,10 @@ int main(int argc, char **argv) { } rpc_trace_set_enabled(0); uint64_t bench_elapsed = bench_started ? now_ns() - bench_start : 0; + struct rusage usage_end; + memset(&usage_end, 0, sizeof(usage_end)); + getrusage(RUSAGE_SELF, &usage_end); + bench_memory memory = current_memory(); stop_server(&server); if (failed) { @@ -427,12 +481,35 @@ int main(int argc, char **argv) { char p95_buf[16]; char p99_buf[16]; char max_buf[16]; + char user_cpu_buf[16]; + char sys_cpu_buf[16]; + char cpu_total_buf[16]; + char rss_buf[16]; + char max_rss_buf[16]; + char virt_buf[16]; + char footprint_buf[16]; + uint64_t user_cpu = timeval_ns(usage_end.ru_utime) - timeval_ns(usage_start.ru_utime); + uint64_t sys_cpu = timeval_ns(usage_end.ru_stime) - timeval_ns(usage_start.ru_stime); + uint64_t cpu_total = user_cpu + sys_cpu; + double cpu_percent = bench_elapsed ? ((double)cpu_total / (double)bench_elapsed) * 100.0 : 0.0; +#ifdef __APPLE__ + uint64_t max_rss = (uint64_t)usage_end.ru_maxrss; +#else + uint64_t max_rss = (uint64_t)usage_end.ru_maxrss * 1024ull; +#endif format_duration(bench_elapsed, elapsed_buf); format_duration(sum / total_requests, avg_buf); format_duration(percentile(latencies, total_requests, 50), p50_buf); format_duration(percentile(latencies, total_requests, 95), p95_buf); format_duration(percentile(latencies, total_requests, 99), p99_buf); format_duration(latencies[total_requests - 1], max_buf); + format_duration(user_cpu, user_cpu_buf); + format_duration(sys_cpu, sys_cpu_buf); + format_duration(cpu_total, cpu_total_buf); + format_bytes(memory.rss, rss_buf); + format_bytes(max_rss, max_rss_buf); + format_bytes(memory.virtual_size, virt_buf); + format_bytes(memory.footprint, footprint_buf); printf("requests: %" PRIu64 "\n", total_requests); printf("clients: %" PRIu64 "\n", clients); @@ -447,6 +524,18 @@ int main(int argc, char **argv) { printf("latency p95: %s\n", p95_buf); printf("latency p99: %s\n", p99_buf); printf("latency max: %s\n", max_buf); + printf("\nresources:\n"); + printf(" cpu user: %s\n", user_cpu_buf); + printf(" cpu sys: %s\n", sys_cpu_buf); + printf(" cpu total: %s (%.0f%%)\n", cpu_total_buf, cpu_percent); + printf(" rss current: %s\n", rss_buf); + printf(" rss max: %s\n", max_rss_buf); + if (memory.virtual_size != 0) { printf(" virtual size: %s\n", virt_buf); } + if (memory.footprint != 0) { printf(" footprint: %s\n", footprint_buf); } + printf(" minor faults: %ld\n", usage_end.ru_minflt - usage_start.ru_minflt); + printf(" major faults: %ld\n", usage_end.ru_majflt - usage_start.ru_majflt); + printf(" voluntary csw: %ld\n", usage_end.ru_nvcsw - usage_start.ru_nvcsw); + printf(" involuntary csw:%ld\n", usage_end.ru_nivcsw - usage_start.ru_nivcsw); if (trace_enabled) { rpc_trace_dump(stdout); } pthread_cond_destroy(&gate.cond); diff --git a/demos/client.c b/demos/client.c index 41900dc..6efbfbd 100644 --- a/demos/client.c +++ b/demos/client.c @@ -26,7 +26,7 @@ int main(int argc, char **argv) { rpc_value *result = NULL; size_t result_count = 0; - if (rpc_client_call(client, 1, &args, &result, &result_count) != 0) { + if (rpc_client_call_name(client, "add", &args, &result, &result_count) != 0) { fprintf(stderr, "RPC error: %s\n", rpc_client_error(client)); rpc_writer_free(&args); rpc_client_close(client); diff --git a/demos/server.c b/demos/server.c index 318b578..488374b 100644 --- a/demos/server.c +++ b/demos/server.c @@ -32,8 +32,8 @@ int main(int argc, char **argv) { uint32_t workers = argc > 3 ? (uint32_t)strtoul(argv[3], NULL, 10) : 0; if (rpc_server_init(&g_server) != 0 || (workers != 0 && rpc_server_set_workers(g_server, workers) != 0) || - rpc_server_add_route(g_server, 1, add_i64, NULL) != 0 || - rpc_server_add_async_route(g_server, 2, echo_string, NULL) != 0 || rpc_server_bind(g_server, host, port) != 0 || + rpc_server_add_route_name(g_server, "add", add_i64, NULL) != 0 || + rpc_server_add_async_route_name(g_server, "echo", echo_string, NULL) != 0 || rpc_server_bind(g_server, host, port) != 0 || rpc_server_listen(g_server) != 0) { fprintf(stderr, "failed to start RPC server\n"); rpc_server_destroy(g_server); diff --git a/include/proc.h b/include/proc.h new file mode 100644 index 0000000..9e98e4a --- /dev/null +++ b/include/proc.h @@ -0,0 +1,8 @@ +#ifndef RPC_PROC_H +#define RPC_PROC_H + +#include + +uint64_t rpc_proc_id(const char *name); + +#endif diff --git a/include/routes.h b/include/routes.h index 174006f..8c003c7 100644 --- a/include/routes.h +++ b/include/routes.h @@ -12,10 +12,11 @@ #define RPC_ROUTE_ROOT_SIZE (1u << (32u - RPC_ROUTE_PAGE_BITS)) typedef struct rpc_route { - uint32_t proc_id; + uint64_t proc_id; rpc_handler_fn handler; void *user_data; int is_async; + struct rpc_route *next; } rpc_route; typedef struct rpc_retired_route { @@ -36,9 +37,9 @@ typedef struct rpc_routes { int rpc_routes_init(rpc_routes *routes); void rpc_routes_destroy(rpc_routes *routes); -int rpc_routes_add(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handler, void *user_data); -int rpc_routes_add_ex(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handler, void *user_data, int is_async); -int rpc_routes_remove(rpc_routes *routes, uint32_t proc_id); -int rpc_routes_lookup(rpc_routes *routes, uint32_t proc_id, rpc_route *out); +int rpc_routes_add(rpc_routes *routes, uint64_t proc_id, rpc_handler_fn handler, void *user_data); +int rpc_routes_add_ex(rpc_routes *routes, uint64_t proc_id, rpc_handler_fn handler, void *user_data, int is_async); +int rpc_routes_remove(rpc_routes *routes, uint64_t proc_id); +int rpc_routes_lookup(rpc_routes *routes, uint64_t proc_id, rpc_route *out); #endif diff --git a/include/rpc/client.h b/include/rpc/client.h index 82d8eda..beed626 100644 --- a/include/rpc/client.h +++ b/include/rpc/client.h @@ -13,9 +13,9 @@ int rpc_client_connect(rpc_client **out_client, const char *host, const char *po void rpc_client_close(rpc_client *client); int rpc_client_ping(rpc_client *client); -int rpc_client_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args, rpc_value **out_values, - size_t *out_count); -int rpc_client_send_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args, uint64_t *out_call_id); +int rpc_client_call_name(rpc_client *client, const char *proc_name, const rpc_writer *args, rpc_value **out_values, + size_t *out_count); +int rpc_client_send_call_name(rpc_client *client, const char *proc_name, const rpc_writer *args, uint64_t *out_call_id); int rpc_client_recv_response(rpc_client *client, uint64_t *out_call_id, rpc_value **out_values, size_t *out_count); const char *rpc_client_error(const rpc_client *client); diff --git a/include/rpc/protocol.h b/include/rpc/protocol.h index 89f7785..438af56 100644 --- a/include/rpc/protocol.h +++ b/include/rpc/protocol.h @@ -9,7 +9,7 @@ extern "C" { #endif -#define RPC_HEADER_SIZE 18u +#define RPC_HEADER_SIZE 22u #define RPC_MAX_PAYLOAD_SIZE (1024u * 1024u) typedef enum rpc_op { @@ -56,7 +56,7 @@ typedef struct rpc_value { typedef struct rpc_header { rpc_op op; uint8_t flags; - uint32_t proc_id; + uint64_t proc_id; uint32_t size; uint64_t call_id; } rpc_header; diff --git a/include/rpc/server.h b/include/rpc/server.h index 9944c30..8ecb67c 100644 --- a/include/rpc/server.h +++ b/include/rpc/server.h @@ -13,7 +13,7 @@ typedef struct rpc_server rpc_server; typedef int (*rpc_handler_fn)(rpc_ctx *ctx, const rpc_value *args, size_t argc, rpc_writer *out, void *user_data); uint64_t rpc_ctx_call_id(const rpc_ctx *ctx); -uint32_t rpc_ctx_proc_id(const rpc_ctx *ctx); +uint64_t rpc_ctx_proc_id(const rpc_ctx *ctx); void rpc_ctx_yield(rpc_ctx *ctx); int rpc_server_init(rpc_server **out_server); @@ -25,9 +25,10 @@ int rpc_server_run(rpc_server *server); void rpc_server_stop(rpc_server *server); void rpc_server_destroy(rpc_server *server); -int rpc_server_add_route(rpc_server *server, uint32_t proc_id, rpc_handler_fn handler, void *user_data); -int rpc_server_add_async_route(rpc_server *server, uint32_t proc_id, rpc_handler_fn handler, void *user_data); -int rpc_server_remove_route(rpc_server *server, uint32_t proc_id); +int rpc_server_add_route_name(rpc_server *server, const char *proc_name, rpc_handler_fn handler, void *user_data); +int rpc_server_add_async_route_name(rpc_server *server, const char *proc_name, rpc_handler_fn handler, + void *user_data); +int rpc_server_remove_route_name(rpc_server *server, const char *proc_name); #ifdef __cplusplus } diff --git a/include/scheduler.h b/include/scheduler.h index c32ab21..67a04ff 100644 --- a/include/scheduler.h +++ b/include/scheduler.h @@ -11,13 +11,13 @@ typedef void (*rpc_call_done_fn)(rpc_call *call, void *user_data); struct rpc_ctx { uint64_t call_id; - uint32_t proc_id; + uint64_t proc_id; void *server; }; int rpc_scheduler_init(rpc_scheduler **out); void rpc_scheduler_destroy(rpc_scheduler *scheduler); -int rpc_scheduler_submit(rpc_scheduler *scheduler, uint64_t call_id, uint32_t proc_id, rpc_handler_fn handler, +int rpc_scheduler_submit(rpc_scheduler *scheduler, uint64_t call_id, uint64_t proc_id, rpc_handler_fn handler, void *handler_data, const uint8_t *payload, size_t payload_len, rpc_call_done_fn done, void *done_data); void rpc_scheduler_run_ready(rpc_scheduler *scheduler); @@ -26,6 +26,6 @@ int rpc_call_result(const rpc_call *call); const rpc_writer *rpc_call_response(const rpc_call *call); const char *rpc_call_error(const rpc_call *call); uint64_t rpc_call_id(const rpc_call *call); -uint32_t rpc_call_proc_id(const rpc_call *call); +uint64_t rpc_call_proc_id(const rpc_call *call); #endif diff --git a/src/client.c b/src/client.c index 90fc3ef..e140b39 100644 --- a/src/client.c +++ b/src/client.c @@ -1,6 +1,8 @@ #include "rpc/client.h" #include "rpc/trace.h" +#include "proc.h" + #include #include #include @@ -100,7 +102,7 @@ static int client_read_fill(rpc_client *client, size_t need) { return 0; } -static int send_packet(rpc_client *client, rpc_op op, uint32_t proc_id, uint64_t call_id, const rpc_writer *payload) { +static int send_packet(rpc_client *client, rpc_op op, uint64_t proc_id, uint64_t call_id, const rpc_writer *payload) { uint64_t trace = rpc_trace_begin(); rpc_header header = { .op = op, @@ -174,6 +176,8 @@ static int recv_packet(rpc_client *client, rpc_header *header, uint8_t **body) { return 0; } +static int client_send_call_id(rpc_client *client, uint64_t proc_id, const rpc_writer *args, uint64_t *out_call_id); + int rpc_client_connect(rpc_client **out_client, const char *host, const char *port) { if (!out_client || !port) { return -1; } @@ -244,8 +248,8 @@ int rpc_client_ping(rpc_client *client) { return 0; } -int rpc_client_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args, rpc_value **out_values, - size_t *out_count) { +static int client_call_id(rpc_client *client, uint64_t proc_id, const rpc_writer *args, rpc_value **out_values, + size_t *out_count) { uint64_t trace_call = rpc_trace_begin(); if (!client || client->fd < 0 || !out_values || !out_count) { rpc_trace_end(RPC_TRACE_CLIENT_CALL, trace_call); @@ -255,7 +259,7 @@ int rpc_client_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args *out_count = 0; uint64_t call_id = 0; - if (rpc_client_send_call(client, proc_id, args, &call_id) != 0) { + if (client_send_call_id(client, proc_id, args, &call_id) != 0) { rpc_trace_end(RPC_TRACE_CLIENT_CALL, trace_call); return -1; } @@ -274,7 +278,12 @@ int rpc_client_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args return rc; } -int rpc_client_send_call(rpc_client *client, uint32_t proc_id, const rpc_writer *args, uint64_t *out_call_id) { +int rpc_client_call_name(rpc_client *client, const char *proc_name, const rpc_writer *args, rpc_value **out_values, + size_t *out_count) { + return client_call_id(client, rpc_proc_id(proc_name), args, out_values, out_count); +} + +static int client_send_call_id(rpc_client *client, uint64_t proc_id, const rpc_writer *args, uint64_t *out_call_id) { if (!client || client->fd < 0 || !out_call_id) { return -1; } uint64_t call_id = client->next_call_id++; if (send_packet(client, RPC_OP_RPC, proc_id, call_id, args) != 0) { return -1; } @@ -282,6 +291,11 @@ int rpc_client_send_call(rpc_client *client, uint32_t proc_id, const rpc_writer return 0; } +int rpc_client_send_call_name(rpc_client *client, const char *proc_name, const rpc_writer *args, + uint64_t *out_call_id) { + return client_send_call_id(client, rpc_proc_id(proc_name), args, out_call_id); +} + int rpc_client_recv_response(rpc_client *client, uint64_t *out_call_id, rpc_value **out_values, size_t *out_count) { if (!client || client->fd < 0 || !out_call_id || !out_values || !out_count) { return -1; } *out_call_id = 0; diff --git a/src/protocol.c b/src/protocol.c index 4133c40..6b40120 100644 --- a/src/protocol.c +++ b/src/protocol.c @@ -1,7 +1,19 @@ #include "rpc/protocol.h" +#include "proc.h" + #include +uint64_t rpc_proc_id(const char *name) { + uint64_t hash = 14695981039346656037ull; + if (!name) { return 0; } + while (*name) { + hash ^= (uint8_t)*name++; + hash *= 1099511628211ull; + } + return hash ? hash : 1u; +} + static void put_u32(uint8_t *out, uint32_t value) { out[0] = (uint8_t)(value >> 24u); out[1] = (uint8_t)(value >> 16u); @@ -37,9 +49,9 @@ int rpc_header_encode(const rpc_header *header, uint8_t out[RPC_HEADER_SIZE]) { out[0] = (uint8_t)header->op; out[1] = header->flags; - put_u32(out + 2, header->proc_id); - put_u32(out + 6, header->size); - put_u64(out + 10, header->call_id); + put_u64(out + 2, header->proc_id); + put_u32(out + 10, header->size); + put_u64(out + 14, header->call_id); return 0; } @@ -49,9 +61,9 @@ int rpc_header_decode(const uint8_t in[RPC_HEADER_SIZE], rpc_header *out) { memset(out, 0, sizeof(*out)); out->op = (rpc_op)in[0]; out->flags = in[1]; - out->proc_id = get_u32(in + 2); - out->size = get_u32(in + 6); - out->call_id = get_u64(in + 10); + out->proc_id = get_u64(in + 2); + out->size = get_u32(in + 10); + out->call_id = get_u64(in + 14); if (out->size > RPC_MAX_PAYLOAD_SIZE) { return -1; } return 0; } diff --git a/src/routes.c b/src/routes.c index 2ed4a5c..6a17f90 100644 --- a/src/routes.c +++ b/src/routes.c @@ -15,6 +15,23 @@ static rpc_route_slot *route_page_alloc(const rpc_routes *routes) { return route_mmap(routes->page_bytes); } +static void route_chain_free(rpc_route *route) { + while (route) { + rpc_route *next = route->next; + free(route); + route = next; + } +} + +static uint32_t route_index(uint64_t proc_id) { + proc_id ^= proc_id >> 33u; + proc_id *= 0xff51afd7ed558ccdull; + proc_id ^= proc_id >> 33u; + proc_id *= 0xc4ceb9fe1a85ec53ull; + proc_id ^= proc_id >> 33u; + return (uint32_t)proc_id; +} + static void free_retired(rpc_routes *routes) { if (atomic_load_explicit(&routes->active_readers, memory_order_acquire) != 0) return; @@ -23,7 +40,7 @@ static void free_retired(rpc_routes *routes) { while (node) { rpc_retired_route *next = node->next; - free(node->route); + route_chain_free(node->route); free(node); node = next; } @@ -76,7 +93,7 @@ void rpc_routes_destroy(rpc_routes *routes) { for (size_t i = 0; i < RPC_ROUTE_PAGE_SIZE; ++i) { rpc_route *route = atomic_load_explicit(&page[i], memory_order_relaxed); - free(route); + route_chain_free(route); } munmap(page, routes->page_bytes); } @@ -87,7 +104,7 @@ retired: while (node) { rpc_retired_route *next = node->next; - free(node->route); + route_chain_free(node->route); free(node); node = next; } @@ -97,11 +114,11 @@ retired: pthread_mutex_destroy(&routes->mutate_lock); } -int rpc_routes_add(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handler, void *user_data) { +int rpc_routes_add(rpc_routes *routes, uint64_t proc_id, rpc_handler_fn handler, void *user_data) { return rpc_routes_add_ex(routes, proc_id, handler, user_data, 0); } -int rpc_routes_add_ex(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handler, void *user_data, int is_async) { +int rpc_routes_add_ex(rpc_routes *routes, uint64_t proc_id, rpc_handler_fn handler, void *user_data, int is_async) { if (!routes || !handler) return -1; rpc_route *route = calloc(1, sizeof(*route)); @@ -110,8 +127,9 @@ int rpc_routes_add_ex(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handl *route = (rpc_route){.proc_id = proc_id, .handler = handler, .user_data = user_data, .is_async = is_async ? 1 : 0}; pthread_mutex_lock(&routes->mutate_lock); - uint32_t page_idx = proc_id >> RPC_ROUTE_PAGE_BITS; - uint32_t slot_idx = proc_id & RPC_ROUTE_PAGE_MASK; + uint32_t index = route_index(proc_id); + uint32_t page_idx = index >> RPC_ROUTE_PAGE_BITS; + uint32_t slot_idx = index & RPC_ROUTE_PAGE_MASK; rpc_route_slot *page = atomic_load_explicit(&routes->pages[page_idx], memory_order_acquire); if (!page) { page = route_page_alloc(routes); @@ -122,7 +140,23 @@ int rpc_routes_add_ex(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handl } atomic_store_explicit(&routes->pages[page_idx], page, memory_order_release); } - rpc_route *old = atomic_exchange_explicit(&page[slot_idx], route, memory_order_acq_rel); + rpc_route *old = atomic_load_explicit(&page[slot_idx], memory_order_acquire); + rpc_route *copy_head = route; + rpc_route **copy_tail = &route->next; + for (rpc_route *it = old; it; it = it->next) { + if (it->proc_id == proc_id) continue; + rpc_route *copy = calloc(1, sizeof(*copy)); + if (!copy) { + route_chain_free(copy_head); + pthread_mutex_unlock(&routes->mutate_lock); + return -1; + } + *copy = *it; + copy->next = NULL; + *copy_tail = copy; + copy_tail = ©->next; + } + atomic_store_explicit(&page[slot_idx], copy_head, memory_order_release); int rc = retire_route(routes, old); pthread_mutex_unlock(&routes->mutate_lock); @@ -130,22 +164,47 @@ int rpc_routes_add_ex(rpc_routes *routes, uint32_t proc_id, rpc_handler_fn handl return rc; } -int rpc_routes_remove(rpc_routes *routes, uint32_t proc_id) { +int rpc_routes_remove(rpc_routes *routes, uint64_t proc_id) { if (!routes) return -1; pthread_mutex_lock(&routes->mutate_lock); - uint32_t page_idx = proc_id >> RPC_ROUTE_PAGE_BITS; - uint32_t slot_idx = proc_id & RPC_ROUTE_PAGE_MASK; + uint32_t index = route_index(proc_id); + uint32_t page_idx = index >> RPC_ROUTE_PAGE_BITS; + uint32_t slot_idx = index & RPC_ROUTE_PAGE_MASK; rpc_route_slot *page = atomic_load_explicit(&routes->pages[page_idx], memory_order_acquire); - rpc_route *old = page ? atomic_exchange_explicit(&page[slot_idx], NULL, memory_order_acq_rel) : NULL; + rpc_route *old = page ? atomic_load_explicit(&page[slot_idx], memory_order_acquire) : NULL; + rpc_route *copy_head = NULL; + rpc_route **copy_tail = ©_head; + int removed = 0; + for (rpc_route *it = old; it; it = it->next) { + if (it->proc_id == proc_id) { + removed = 1; + continue; + } + rpc_route *copy = calloc(1, sizeof(*copy)); + if (!copy) { + route_chain_free(copy_head); + pthread_mutex_unlock(&routes->mutate_lock); + return -1; + } + *copy = *it; + copy->next = NULL; + *copy_tail = copy; + copy_tail = ©->next; + } + if (page && removed) { + atomic_store_explicit(&page[slot_idx], copy_head, memory_order_release); + } else { + route_chain_free(copy_head); + } - int rc = old ? retire_route(routes, old) : -1; + int rc = removed ? retire_route(routes, old) : -1; pthread_mutex_unlock(&routes->mutate_lock); return rc; } -int rpc_routes_lookup(rpc_routes *routes, uint32_t proc_id, rpc_route *out) { +int rpc_routes_lookup(rpc_routes *routes, uint64_t proc_id, rpc_route *out) { uint64_t trace = rpc_trace_begin(); if (!routes || !out) { rpc_trace_end(RPC_TRACE_ROUTE_LOOKUP, trace); @@ -153,10 +212,14 @@ int rpc_routes_lookup(rpc_routes *routes, uint32_t proc_id, rpc_route *out) { } atomic_fetch_add_explicit(&routes->active_readers, 1u, memory_order_acquire); - uint32_t page_idx = proc_id >> RPC_ROUTE_PAGE_BITS; - uint32_t slot_idx = proc_id & RPC_ROUTE_PAGE_MASK; + uint32_t index = route_index(proc_id); + uint32_t page_idx = index >> RPC_ROUTE_PAGE_BITS; + uint32_t slot_idx = index & RPC_ROUTE_PAGE_MASK; rpc_route_slot *page = atomic_load_explicit(&routes->pages[page_idx], memory_order_acquire); rpc_route *route = page ? atomic_load_explicit(&page[slot_idx], memory_order_acquire) : NULL; + while (route && route->proc_id != proc_id) { + route = route->next; + } if (route) { *out = *route; } unsigned old = atomic_fetch_sub_explicit(&routes->active_readers, 1u, memory_order_release); if (old == 1u) { diff --git a/src/scheduler.c b/src/scheduler.c index cab968a..ec19f03 100644 --- a/src/scheduler.c +++ b/src/scheduler.c @@ -117,7 +117,7 @@ void rpc_scheduler_destroy(rpc_scheduler *scheduler) { free(scheduler); } -int rpc_scheduler_submit(rpc_scheduler *scheduler, uint64_t call_id, uint32_t proc_id, rpc_handler_fn handler, +int rpc_scheduler_submit(rpc_scheduler *scheduler, uint64_t call_id, uint64_t proc_id, rpc_handler_fn handler, void *handler_data, const uint8_t *payload, size_t payload_len, rpc_call_done_fn done, void *done_data) { uint64_t trace_submit = rpc_trace_begin(); @@ -221,7 +221,7 @@ uint64_t rpc_call_id(const rpc_call *call) { return call ? call->ctx.call_id : 0; } -uint32_t rpc_call_proc_id(const rpc_call *call) { +uint64_t rpc_call_proc_id(const rpc_call *call) { return call ? call->ctx.proc_id : 0; } @@ -229,7 +229,7 @@ uint64_t rpc_ctx_call_id(const rpc_ctx *ctx) { return ctx ? ctx->call_id : 0; } -uint32_t rpc_ctx_proc_id(const rpc_ctx *ctx) { +uint64_t rpc_ctx_proc_id(const rpc_ctx *ctx) { return ctx ? ctx->proc_id : 0; } diff --git a/src/server.c b/src/server.c index 2fe9383..a836845 100644 --- a/src/server.c +++ b/src/server.c @@ -3,6 +3,7 @@ #include "arena.h" #include "backend.h" +#include "proc.h" #include "routes.h" #include "scheduler.h" @@ -155,7 +156,7 @@ static int reserve_bytes(uint8_t **buf, size_t *cap, size_t need) { return 0; } -static int queue_packet(rpc_connection *conn, rpc_op op, uint8_t flags, uint32_t proc_id, uint64_t call_id, +static int queue_packet(rpc_connection *conn, rpc_op op, uint8_t flags, uint64_t proc_id, uint64_t call_id, const uint8_t *payload, size_t payload_len) { if (payload_len > RPC_MAX_PAYLOAD_SIZE) { return -1; } uint8_t header_buf[RPC_HEADER_SIZE]; @@ -176,7 +177,7 @@ static int queue_packet(rpc_connection *conn, rpc_op op, uint8_t flags, uint32_t return 0; } -static int queue_string_error(rpc_connection *conn, uint32_t proc_id, uint64_t call_id, const char *message) { +static int queue_string_error(rpc_connection *conn, uint64_t proc_id, uint64_t call_id, const char *message) { rpc_writer writer; rpc_writer_init(&writer); int rc = rpc_writer_string(&writer, message, (uint32_t)strlen(message)); @@ -740,14 +741,27 @@ void rpc_server_destroy(rpc_server *server) { free(server); } -int rpc_server_add_route(rpc_server *server, uint32_t proc_id, rpc_handler_fn handler, void *user_data) { +static int server_add_route_id(rpc_server *server, uint64_t proc_id, rpc_handler_fn handler, void *user_data) { return server ? rpc_routes_add(&server->routes, proc_id, handler, user_data) : -1; } -int rpc_server_add_async_route(rpc_server *server, uint32_t proc_id, rpc_handler_fn handler, void *user_data) { +int rpc_server_add_route_name(rpc_server *server, const char *proc_name, rpc_handler_fn handler, void *user_data) { + return server_add_route_id(server, rpc_proc_id(proc_name), handler, user_data); +} + +static int server_add_async_route_id(rpc_server *server, uint64_t proc_id, rpc_handler_fn handler, void *user_data) { return server ? rpc_routes_add_ex(&server->routes, proc_id, handler, user_data, 1) : -1; } -int rpc_server_remove_route(rpc_server *server, uint32_t proc_id) { +int rpc_server_add_async_route_name(rpc_server *server, const char *proc_name, rpc_handler_fn handler, + void *user_data) { + return server_add_async_route_id(server, rpc_proc_id(proc_name), handler, user_data); +} + +static int server_remove_route_id(rpc_server *server, uint64_t proc_id) { return server ? rpc_routes_remove(&server->routes, proc_id) : -1; } + +int rpc_server_remove_route_name(rpc_server *server, const char *proc_name) { + return server_remove_route_id(server, rpc_proc_id(proc_name)); +} diff --git a/tests/test_integration.c b/tests/test_integration.c index fff947c..b906353 100644 --- a/tests/test_integration.c +++ b/tests/test_integration.c @@ -43,7 +43,7 @@ static void *client_thread(void *arg) { rpc_value *values = NULL; size_t count = 0; - assert(rpc_client_call(client, 7, &payload, &values, &count) == 0); + assert(rpc_client_call_name(client, "add", &payload, &values, &count) == 0); assert(count == 1); assert(values[0].as.i64 == job->index + 10); @@ -56,7 +56,7 @@ static void *client_thread(void *arg) { int main(void) { rpc_server *server = NULL; assert(rpc_server_init(&server) == 0); - assert(rpc_server_add_async_route(server, 7, add_handler, NULL) == 0); + assert(rpc_server_add_async_route_name(server, "add", add_handler, NULL) == 0); assert(rpc_server_bind(server, "127.0.0.1", "0") == 0); assert(rpc_server_listen(server) == 0); uint16_t port = rpc_server_port(server); @@ -79,12 +79,12 @@ int main(void) { rpc_value *values = NULL; size_t count = 0; - assert(rpc_client_call(client, 7, &payload, &values, &count) == 0); + assert(rpc_client_call_name(client, "add", &payload, &values, &count) == 0); assert(count == 1 && values[0].type == RPC_TYPE_I64 && values[0].as.i64 == 11); rpc_values_free(values); rpc_writer_reset(&payload); - assert(rpc_client_call(client, 999, &payload, &values, &count) != 0); + assert(rpc_client_call_name(client, "missing", &payload, &values, &count) != 0); enum { CLIENTS = 4 }; pthread_t clients[CLIENTS]; @@ -97,8 +97,8 @@ int main(void) { assert(pthread_join(clients[i], NULL) == 0); } - assert(rpc_server_remove_route(server, 7) == 0); - assert(rpc_client_call(client, 7, &payload, &values, &count) != 0); + assert(rpc_server_remove_route_name(server, "add") == 0); + assert(rpc_client_call_name(client, "add", &payload, &values, &count) != 0); rpc_writer_free(&payload); rpc_client_close(client); diff --git a/tests/test_protocol.c b/tests/test_protocol.c index 2aaf0a3..60f7dd7 100644 --- a/tests/test_protocol.c +++ b/tests/test_protocol.c @@ -9,7 +9,7 @@ int main(void) { rpc_header h = { .op = RPC_OP_RPC, .flags = RPC_FLAG_MORE, - .proc_id = 0x11223344u, + .proc_id = 0x1122334455667788ull, .size = 9, .call_id = 0x0102030405060708ull, }; -- 2.51.2