diff --git a/demos/bench.c b/demos/bench.c new file mode 100644 index 0000000..ab90243 --- /dev/null +++ b/demos/bench.c @@ -0,0 +1,249 @@ +#include "rpc/client.h" +#include "rpc/server.h" + +#include +#include +#include +#include +#include +#include +#include +#include + +typedef struct bench_server { + rpc_server *rpc; + pthread_t thread; +} bench_server; + +typedef struct worker_args { + const char *host; + char port[16]; + uint64_t requests; + uint64_t warmup; + uint64_t latency_offset; + uint64_t *latencies_ns; + int failed; +} worker_args; + +static uint64_t now_ns(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (uint64_t)ts.tv_sec * 1000000000ull + (uint64_t)ts.tv_nsec; +} + +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; + if (argc != 2 || args[0].type != RPC_TYPE_I64 || args[1].type != RPC_TYPE_I64) { + return -1; + } + return rpc_writer_i64(out, args[0].as.i64 + args[1].as.i64); +} + +static void *server_main(void *arg) { + bench_server *server = arg; + (void)rpc_server_run(server->rpc); + return NULL; +} + +static int start_server(bench_server *server, char port[16]) { + memset(server, 0, sizeof(*server)); + if (rpc_server_init(&server->rpc) != 0 || + rpc_server_add_route(server->rpc, 1, add_handler, NULL) != 0 || + rpc_server_bind(server->rpc, "127.0.0.1", "0") != 0 || + rpc_server_listen(server->rpc) != 0) { + return -1; + } + + snprintf(port, 16, "%u", rpc_server_port(server->rpc)); + if (pthread_create(&server->thread, NULL, server_main, server) != 0) { + rpc_server_destroy(server->rpc); + server->rpc = NULL; + return -1; + } + return 0; +} + +static void stop_server(bench_server *server) { + if (!server->rpc) { + return; + } + rpc_server_stop(server->rpc); + pthread_join(server->thread, NULL); + rpc_server_destroy(server->rpc); + server->rpc = NULL; +} + +static int run_one_call(rpc_client *client, rpc_writer *payload, int64_t expected) { + rpc_value *values = NULL; + size_t count = 0; + int rc = rpc_client_call(client, 1, 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; +} + +static void *worker_main(void *arg) { + worker_args *worker = arg; + rpc_client *client = NULL; + if (rpc_client_connect(&client, worker->host, worker->port) != 0) { + worker->failed = 1; + return NULL; + } + + rpc_writer payload; + rpc_writer_init(&payload); + if (rpc_writer_i64(&payload, 20) != 0 || rpc_writer_i64(&payload, 22) != 0) { + worker->failed = 1; + goto done; + } + + for (uint64_t i = 0; i < worker->warmup; ++i) { + if (run_one_call(client, &payload, 42) != 0) { + worker->failed = 1; + goto done; + } + } + + for (uint64_t i = 0; i < worker->requests; ++i) { + uint64_t start = now_ns(); + if (run_one_call(client, &payload, 42) != 0) { + worker->failed = 1; + goto done; + } + worker->latencies_ns[worker->latency_offset + i] = now_ns() - start; + } + +done: + rpc_writer_free(&payload); + rpc_client_close(client); + return NULL; +} + +static int cmp_u64(const void *a, const void *b) { + uint64_t x = *(const uint64_t *)a; + uint64_t y = *(const uint64_t *)b; + return (x > y) - (x < y); +} + +static double ns_to_us(uint64_t ns) { return (double)ns / 1000.0; } + +static uint64_t percentile(const uint64_t *values, uint64_t count, double pct) { + if (count == 0) { + return 0; + } + uint64_t idx = (uint64_t)((pct / 100.0) * (double)(count - 1)); + return values[idx]; +} + +static uint64_t parse_u64_arg(const char *text, uint64_t fallback) { + if (!text || !*text) { + return fallback; + } + char *end = NULL; + errno = 0; + unsigned long long value = strtoull(text, &end, 10); + if (errno != 0 || !end || *end != '\0' || value == 0) { + return fallback; + } + return (uint64_t)value; +} + +int main(int argc, char **argv) { + uint64_t total_requests = argc > 1 ? parse_u64_arg(argv[1], 50000) : 50000; + uint64_t clients = argc > 2 ? parse_u64_arg(argv[2], 1) : 1; + uint64_t warmup_total = argc > 3 ? parse_u64_arg(argv[3], 1000) : 1000; + if (clients > total_requests) { + clients = total_requests; + } + + uint64_t *latencies = calloc(total_requests, sizeof(*latencies)); + pthread_t *threads = calloc(clients, sizeof(*threads)); + worker_args *workers = calloc(clients, sizeof(*workers)); + if (!latencies || !threads || !workers) { + fprintf(stderr, "allocation failed\n"); + free(latencies); + free(threads); + free(workers); + return 1; + } + + bench_server server; + char port[16]; + if (start_server(&server, port) != 0) { + fprintf(stderr, "failed to start benchmark server\n"); + free(latencies); + free(threads); + free(workers); + return 1; + } + + uint64_t offset = 0; + uint64_t warmup_base = warmup_total / clients; + uint64_t warmup_rem = warmup_total % clients; + uint64_t request_base = total_requests / clients; + uint64_t request_rem = total_requests % clients; + + uint64_t bench_start = now_ns(); + for (uint64_t i = 0; i < clients; ++i) { + uint64_t count = request_base + (i < request_rem ? 1u : 0u); + workers[i] = (worker_args){ + .host = "127.0.0.1", + .requests = count, + .warmup = warmup_base + (i < warmup_rem ? 1u : 0u), + .latency_offset = offset, + .latencies_ns = latencies, + }; + snprintf(workers[i].port, sizeof(workers[i].port), "%s", port); + offset += count; + if (pthread_create(&threads[i], NULL, worker_main, &workers[i]) != 0) { + workers[i].failed = 1; + } + } + + int failed = 0; + for (uint64_t i = 0; i < clients; ++i) { + pthread_join(threads[i], NULL); + failed |= workers[i].failed; + } + uint64_t bench_elapsed = now_ns() - bench_start; + stop_server(&server); + + if (failed) { + fprintf(stderr, "benchmark failed\n"); + free(latencies); + free(threads); + free(workers); + return 1; + } + + qsort(latencies, total_requests, sizeof(*latencies), cmp_u64); + uint64_t sum = 0; + for (uint64_t i = 0; i < total_requests; ++i) { + sum += latencies[i]; + } + + double seconds = (double)bench_elapsed / 1000000000.0; + double throughput = (double)total_requests / seconds; + double avg_us = ns_to_us(sum / total_requests); + + printf("requests: %" PRIu64 "\n", total_requests); + printf("clients: %" PRIu64 "\n", clients); + printf("warmup: %" PRIu64 "\n", warmup_total); + printf("elapsed: %.3f s\n", seconds); + printf("throughput: %.0f req/s\n", throughput); + printf("latency avg: %.2f us\n", avg_us); + printf("latency p50: %.2f us\n", ns_to_us(percentile(latencies, total_requests, 50))); + printf("latency p95: %.2f us\n", ns_to_us(percentile(latencies, total_requests, 95))); + printf("latency p99: %.2f us\n", ns_to_us(percentile(latencies, total_requests, 99))); + printf("latency max: %.2f us\n", ns_to_us(latencies[total_requests - 1])); + + free(latencies); + free(threads); + free(workers); + return 0; +} diff --git a/demos/client.c b/demos/client.c new file mode 100644 index 0000000..4c25ee5 --- /dev/null +++ b/demos/client.c @@ -0,0 +1,43 @@ +#include "rpc/client.h" + +#include +#include + +int main(int argc, char **argv) { + const char *host = argc > 1 ? argv[1] : "127.0.0.1"; + const char *port = argc > 2 ? argv[2] : "7000"; + int64_t a = argc > 3 ? strtoll(argv[3], NULL, 10) : 20; + int64_t b = argc > 4 ? strtoll(argv[4], NULL, 10) : 22; + + rpc_client *client = NULL; + if (rpc_client_connect(&client, host, port) != 0) { + fprintf(stderr, "connect failed\n"); + return 1; + } + + rpc_writer args; + rpc_writer_init(&args); + rpc_writer_i64(&args, a); + rpc_writer_i64(&args, b); + + rpc_value *result = NULL; + size_t result_count = 0; + if (rpc_client_call(client, 1, &args, &result, &result_count) != 0) { + fprintf(stderr, "RPC error: %s\n", rpc_client_error(client)); + rpc_writer_free(&args); + rpc_client_close(client); + return 1; + } + + if (result_count == 1 && result[0].type == RPC_TYPE_I64) { + printf("%lld + %lld = %lld\n", (long long)a, (long long)b, + (long long)result[0].as.i64); + } else { + fprintf(stderr, "unexpected response\n"); + } + + rpc_values_free(result); + rpc_writer_free(&args); + rpc_client_close(client); + return 0; +} diff --git a/demos/server.c b/demos/server.c new file mode 100644 index 0000000..66ad272 --- /dev/null +++ b/demos/server.c @@ -0,0 +1,57 @@ +#include "rpc/server.h" + +#include +#include +#include +#include + +static rpc_server *g_server; + +static void on_signal(int signo) { + (void)signo; + if (g_server) { + rpc_server_stop(g_server); + } +} + +static int add_i64(rpc_ctx *ctx, const rpc_value *args, size_t argc, + rpc_writer *out, void *user_data) { + (void)ctx; + (void)user_data; + if (argc != 2 || args[0].type != RPC_TYPE_I64 || args[1].type != RPC_TYPE_I64) { + return -1; + } + return rpc_writer_i64(out, args[0].as.i64 + args[1].as.i64); +} + +static int echo_string(rpc_ctx *ctx, const rpc_value *args, size_t argc, + rpc_writer *out, void *user_data) { + (void)user_data; + if (argc != 1 || args[0].type != RPC_TYPE_STRING) { + return -1; + } + rpc_ctx_yield(ctx); + return rpc_writer_string(out, args[0].as.string.data, args[0].as.string.len); +} + +int main(int argc, char **argv) { + const char *host = argc > 1 ? argv[1] : "127.0.0.1"; + const char *port = argc > 2 ? argv[2] : "7000"; + + if (rpc_server_init(&g_server) != 0 || + rpc_server_add_route(g_server, 1, add_i64, NULL) != 0 || + rpc_server_add_route(g_server, 2, 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); + return 1; + } + + signal(SIGINT, on_signal); + signal(SIGTERM, on_signal); + printf("rpc demo server listening on %s:%u\n", host, rpc_server_port(g_server)); + int rc = rpc_server_run(g_server); + rpc_server_destroy(g_server); + return rc == 0 ? 0 : 1; +}