diff --git a/include/rpc/trace.h b/include/rpc/trace.h index 99e3833..af454d3 100644 --- a/include/rpc/trace.h +++ b/include/rpc/trace.h @@ -31,6 +31,18 @@ typedef enum rpc_trace_metric { RPC_TRACE_COUNT, } rpc_trace_metric; +typedef enum rpc_trace_worker_metric { + RPC_TRACE_WORKER_POLL_EVENTS, + RPC_TRACE_WORKER_ACCEPTS, + RPC_TRACE_WORKER_READS, + RPC_TRACE_WORKER_RPCS, + RPC_TRACE_WORKER_WRITES, + RPC_TRACE_WORKER_ACTIVE, + RPC_TRACE_WORKER_COUNT, +} rpc_trace_worker_metric; + +#define RPC_TRACE_MAX_WORKERS 128u + typedef struct rpc_trace_stat { const char *name; uint64_t count; @@ -45,6 +57,10 @@ void rpc_trace_set_enabled(int enabled); uint64_t rpc_trace_begin_slow(void); void rpc_trace_end_slow(rpc_trace_metric metric, uint64_t start_ns); void rpc_trace_add_slow(rpc_trace_metric metric, uint64_t value); +void rpc_trace_worker_end_slow(uint32_t worker, rpc_trace_worker_metric metric, + uint64_t start_ns); +void rpc_trace_worker_add_slow(uint32_t worker, rpc_trace_worker_metric metric, + uint64_t value); void rpc_trace_reset(void); void rpc_trace_snapshot(rpc_trace_stat out[RPC_TRACE_COUNT]); void rpc_trace_dump(FILE *out); @@ -80,6 +96,22 @@ static inline void rpc_trace_add(rpc_trace_metric metric, uint64_t value) { } } +static inline void rpc_trace_worker_end(uint32_t worker, + rpc_trace_worker_metric metric, + uint64_t start_ns) { + if (rpc_trace_is_enabled() && start_ns != 0) { + rpc_trace_worker_end_slow(worker, metric, start_ns); + } +} + +static inline void rpc_trace_worker_add(uint32_t worker, + rpc_trace_worker_metric metric, + uint64_t value) { + if (rpc_trace_is_enabled()) { + rpc_trace_worker_add_slow(worker, metric, value); + } +} + #ifdef __cplusplus } #endif diff --git a/src/server.c b/src/server.c index 2c9c6d6..df4da2b 100644 --- a/src/server.c +++ b/src/server.c @@ -47,6 +47,11 @@ typedef struct rpc_connection { struct rpc_connection *next; } rpc_connection; +typedef struct rpc_pending_fd { + int fd; + struct rpc_pending_fd *next; +} rpc_pending_fd; + struct rpc_worker { rpc_server *server; rpc_backend *backend; @@ -55,6 +60,10 @@ struct rpc_worker { rpc_listener listeners[RPC_MAX_LISTENERS]; size_t listener_count; rpc_connection *connections; + pthread_mutex_t pending_mutex; + rpc_pending_fd *pending_head; + rpc_pending_fd *pending_tail; + int pending_mutex_ready; pthread_t thread; int thread_started; uint32_t index; @@ -65,6 +74,7 @@ struct rpc_server { rpc_worker *workers; uint32_t worker_count; uint32_t workers_ready; + atomic_uint next_worker; atomic_bool stopping; int listening; int routes_ready; @@ -72,6 +82,7 @@ struct rpc_server { }; static void connection_write(rpc_connection *conn); +static int worker_add_connection(rpc_worker *worker, int fd); static int connection_set_interests(rpc_connection *conn, uint32_t interests) { if (conn->interests == interests) { @@ -274,6 +285,7 @@ op_disconnect: return 0; op_rpc: { + rpc_trace_worker_add(conn->worker->index, RPC_TRACE_WORKER_RPCS, 1); rpc_route route; uint64_t trace_route = rpc_trace_begin(); if (rpc_routes_lookup(&conn->worker->server->routes, header->proc_id, &route) != 0) { @@ -328,6 +340,7 @@ static void parse_available(rpc_connection *conn) { } static void connection_read(rpc_connection *conn) { + rpc_trace_worker_add(conn->worker->index, RPC_TRACE_WORKER_READS, 1); uint64_t trace = rpc_trace_begin(); for (;;) { if (conn->read_cap - conn->read_len < RPC_READ_CHUNK) { @@ -368,6 +381,7 @@ static void connection_read(rpc_connection *conn) { } static void connection_write(rpc_connection *conn) { + rpc_trace_worker_add(conn->worker->index, RPC_TRACE_WORKER_WRITES, 1); uint64_t trace = rpc_trace_begin(); while (conn->write_off < conn->write_len) { ssize_t n = send(conn->fd, conn->write_buf + conn->write_off, @@ -395,9 +409,66 @@ static void connection_write(rpc_connection *conn) { rpc_trace_end(RPC_TRACE_SERVER_WRITE, trace); } +static int worker_add_connection(rpc_worker *worker, int fd) { + rpc_connection *conn = rpc_fixed_arena_alloc(&worker->connection_arena); + if (!conn) { + close(fd); + return -1; + } + conn->fd = fd; + conn->worker = worker; + conn->next = worker->connections; + worker->connections = conn; + if (rpc_backend_register(worker->backend, fd, RPC_BACKEND_READ, + (uintptr_t)conn) != 0) { + connection_destroy(conn); + return -1; + } + conn->interests = RPC_BACKEND_READ; + rpc_trace_worker_add(worker->index, RPC_TRACE_WORKER_ACCEPTS, 1); + return 0; +} + +static void worker_enqueue_connection(rpc_worker *worker, int fd) { + rpc_pending_fd *pending = malloc(sizeof(*pending)); + if (!pending) { + close(fd); + return; + } + pending->fd = fd; + pending->next = NULL; + + pthread_mutex_lock(&worker->pending_mutex); + if (worker->pending_tail) { + worker->pending_tail->next = pending; + } else { + worker->pending_head = pending; + } + worker->pending_tail = pending; + pthread_mutex_unlock(&worker->pending_mutex); + + (void)rpc_backend_wake(worker->backend); +} + +static void worker_drain_pending(rpc_worker *worker) { + pthread_mutex_lock(&worker->pending_mutex); + rpc_pending_fd *pending = worker->pending_head; + worker->pending_head = NULL; + worker->pending_tail = NULL; + pthread_mutex_unlock(&worker->pending_mutex); + + while (pending) { + rpc_pending_fd *next = pending->next; + (void)worker_add_connection(worker, pending->fd); + free(pending); + pending = next; + } +} + static void accept_ready(rpc_listener *listener) { uint64_t trace = rpc_trace_begin(); rpc_worker *worker = listener->worker; + rpc_server *server = worker->server; for (;;) { struct sockaddr_storage addr; socklen_t addr_len = sizeof(addr); @@ -413,21 +484,11 @@ static void accept_ready(rpc_listener *listener) { close(fd); continue; } - rpc_connection *conn = rpc_fixed_arena_alloc(&worker->connection_arena); - if (!conn) { - close(fd); - continue; - } - conn->fd = fd; - conn->worker = worker; - conn->next = worker->connections; - worker->connections = conn; - if (rpc_backend_register(worker->backend, fd, RPC_BACKEND_READ, - (uintptr_t)conn) != 0) { - connection_destroy(conn); - } else { - conn->interests = RPC_BACKEND_READ; - } + uint32_t idx = + atomic_fetch_add_explicit(&server->next_worker, 1, + memory_order_relaxed) % + server->worker_count; + worker_enqueue_connection(&server->workers[idx], fd); } } @@ -438,6 +499,10 @@ static int worker_init(rpc_server *server, rpc_worker *worker, uint32_t index) { for (size_t i = 0; i < RPC_MAX_LISTENERS; ++i) { worker->listeners[i].fd = -1; } + if (pthread_mutex_init(&worker->pending_mutex, NULL) != 0) { + return -1; + } + worker->pending_mutex_ready = 1; if (rpc_backend_kqueue_create(&worker->backend) != 0) { return -1; } @@ -468,6 +533,20 @@ static void worker_destroy(rpc_worker *worker) { worker->listeners[i].fd = -1; } } + if (worker->pending_mutex_ready) { + pthread_mutex_lock(&worker->pending_mutex); + rpc_pending_fd *pending = worker->pending_head; + worker->pending_head = NULL; + worker->pending_tail = NULL; + pthread_mutex_unlock(&worker->pending_mutex); + while (pending) { + rpc_pending_fd *next = pending->next; + close(pending->fd); + free(pending); + pending = next; + } + pthread_mutex_destroy(&worker->pending_mutex); + } rpc_scheduler_destroy(worker->scheduler); rpc_fixed_arena_destroy(&worker->connection_arena); rpc_backend_destroy(worker->backend); @@ -508,6 +587,7 @@ int rpc_server_init(rpc_server **out_server) { return -1; } atomic_init(&server->stopping, false); + atomic_init(&server->next_worker, 0); server->worker_count = cpu_count(); if (rpc_routes_init(&server->routes) != 0) { rpc_server_destroy(server); @@ -549,63 +629,37 @@ int rpc_server_bind(rpc_server *server, const char *host, const char *port) { break; } - size_t added[RPC_MAX_WORKERS]; - memset(added, 0, sizeof(added)); - int bound_all = 1; - - for (uint32_t wi = 0; wi < server->worker_count; ++wi) { - rpc_worker *worker = &server->workers[wi]; - int fd = socket(ai->ai_family, ai->ai_socktype, ai->ai_protocol); - if (fd < 0) { - bound_all = 0; - break; - } - int yes = 1; - (void)setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &yes, sizeof(yes)); -#ifdef SO_REUSEPORT - (void)setsockopt(fd, SOL_SOCKET, SO_REUSEPORT, &yes, sizeof(yes)); -#endif - - struct sockaddr_storage addr; - memcpy(&addr, ai->ai_addr, ai->ai_addrlen); - if (server->port != 0) { - sockaddr_set_port((struct sockaddr *)&addr, server->port); - } - - if (set_nonblock(fd) != 0 || - bind(fd, (struct sockaddr *)&addr, ai->ai_addrlen) != 0) { - close(fd); - bound_all = 0; - break; - } + rpc_worker *worker = &server->workers[0]; + int fd = socket(ai->ai_family, ai->ai_socktype, ai->ai_protocol); + if (fd < 0) { + continue; + } + int yes = 1; + (void)setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &yes, sizeof(yes)); - if (server->port == 0) { - struct sockaddr_storage bound; - socklen_t bound_len = sizeof(bound); - if (getsockname(fd, (struct sockaddr *)&bound, &bound_len) == 0) { - server->port = sockaddr_get_port((struct sockaddr *)&bound); - } - } + struct sockaddr_storage addr; + memcpy(&addr, ai->ai_addr, ai->ai_addrlen); + if (server->port != 0) { + sockaddr_set_port((struct sockaddr *)&addr, server->port); + } - rpc_listener *listener = &worker->listeners[worker->listener_count++]; - listener->fd = fd; - listener->worker = worker; - added[wi] = 1; + if (set_nonblock(fd) != 0 || + bind(fd, (struct sockaddr *)&addr, ai->ai_addrlen) != 0) { + close(fd); + continue; } - if (!bound_all) { - for (uint32_t wi = 0; wi < server->worker_count; ++wi) { - if (!added[wi]) { - continue; - } - rpc_worker *worker = &server->workers[wi]; - rpc_listener *listener = &worker->listeners[--worker->listener_count]; - close(listener->fd); - listener->fd = -1; - listener->worker = NULL; + if (server->port == 0) { + struct sockaddr_storage bound; + socklen_t bound_len = sizeof(bound); + if (getsockname(fd, (struct sockaddr *)&bound, &bound_len) == 0) { + server->port = sockaddr_get_port((struct sockaddr *)&bound); } - continue; } + + rpc_listener *listener = &worker->listeners[worker->listener_count++]; + listener->fd = fd; + listener->worker = worker; ok = 0; } @@ -652,9 +706,12 @@ static int worker_run(rpc_worker *worker) { return -1; } rpc_trace_add(RPC_TRACE_SERVER_POLL_EVENTS, (uint64_t)n); + rpc_trace_worker_add(worker->index, RPC_TRACE_WORKER_POLL_EVENTS, + (uint64_t)n); uint64_t trace_active = rpc_trace_begin(); for (int i = 0; i < n; ++i) { if (events[i].events & RPC_BACKEND_WAKE) { + worker_drain_pending(worker); continue; } if (events[i].user & 1u) { @@ -674,6 +731,7 @@ static int worker_run(rpc_worker *worker) { uint64_t trace_schedule = rpc_trace_begin(); rpc_scheduler_run_ready(worker->scheduler); rpc_trace_end(RPC_TRACE_SERVER_SCHEDULE, trace_schedule); + rpc_trace_worker_end(worker->index, RPC_TRACE_WORKER_ACTIVE, trace_active); rpc_trace_end(RPC_TRACE_SERVER_LOOP_ACTIVE, trace_active); } return 0; diff --git a/src/trace.c b/src/trace.c index 05ad4f5..a6ffae3 100644 --- a/src/trace.c +++ b/src/trace.c @@ -61,6 +61,8 @@ static const char *trace_max_units[RPC_TRACE_COUNT] = { }; static rpc_trace_counter counters[RPC_TRACE_COUNT]; +static rpc_trace_counter worker_counters[RPC_TRACE_MAX_WORKERS] + [RPC_TRACE_WORKER_COUNT]; int rpc_trace_enabled = 0; @@ -110,12 +112,48 @@ void rpc_trace_add_slow(rpc_trace_metric metric, uint64_t value) { } } +static void trace_counter_add(rpc_trace_counter *counter, uint64_t value) { + atomic_fetch_add_explicit(&counter->count, 1, memory_order_relaxed); + atomic_fetch_add_explicit(&counter->total_ns, value, memory_order_relaxed); + + uint64_t old = atomic_load_explicit(&counter->max_ns, memory_order_relaxed); + while (old < value && + !atomic_compare_exchange_weak_explicit(&counter->max_ns, &old, value, + memory_order_relaxed, + memory_order_relaxed)) { + } +} + +void rpc_trace_worker_end_slow(uint32_t worker, rpc_trace_worker_metric metric, + uint64_t start_ns) { + rpc_trace_worker_add_slow(worker, metric, rpc_trace_begin_slow() - start_ns); +} + +void rpc_trace_worker_add_slow(uint32_t worker, rpc_trace_worker_metric metric, + uint64_t value) { + if (worker >= RPC_TRACE_MAX_WORKERS || + (unsigned)metric >= RPC_TRACE_WORKER_COUNT) { + return; + } + trace_counter_add(&worker_counters[worker][metric], value); +} + void rpc_trace_reset(void) { for (size_t i = 0; i < RPC_TRACE_COUNT; ++i) { atomic_store_explicit(&counters[i].count, 0, memory_order_relaxed); atomic_store_explicit(&counters[i].total_ns, 0, memory_order_relaxed); atomic_store_explicit(&counters[i].max_ns, 0, memory_order_relaxed); } + for (size_t worker = 0; worker < RPC_TRACE_MAX_WORKERS; ++worker) { + for (size_t metric = 0; metric < RPC_TRACE_WORKER_COUNT; ++metric) { + atomic_store_explicit(&worker_counters[worker][metric].count, 0, + memory_order_relaxed); + atomic_store_explicit(&worker_counters[worker][metric].total_ns, 0, + memory_order_relaxed); + atomic_store_explicit(&worker_counters[worker][metric].max_ns, 0, + memory_order_relaxed); + } + } } void rpc_trace_snapshot(rpc_trace_stat out[RPC_TRACE_COUNT]) { @@ -172,4 +210,70 @@ void rpc_trace_dump(FILE *out) { } } } + + int printed_workers = 0; + for (size_t worker = 0; worker < RPC_TRACE_MAX_WORKERS; ++worker) { + rpc_trace_counter *poll_events = + &worker_counters[worker][RPC_TRACE_WORKER_POLL_EVENTS]; + rpc_trace_counter *accepts = + &worker_counters[worker][RPC_TRACE_WORKER_ACCEPTS]; + rpc_trace_counter *reads = + &worker_counters[worker][RPC_TRACE_WORKER_READS]; + rpc_trace_counter *rpcs = + &worker_counters[worker][RPC_TRACE_WORKER_RPCS]; + rpc_trace_counter *writes = + &worker_counters[worker][RPC_TRACE_WORKER_WRITES]; + rpc_trace_counter *active = + &worker_counters[worker][RPC_TRACE_WORKER_ACTIVE]; + + uint64_t polls_count = + atomic_load_explicit(&poll_events->count, memory_order_relaxed); + uint64_t polls_total = + atomic_load_explicit(&poll_events->total_ns, memory_order_relaxed); + uint64_t polls_max = + atomic_load_explicit(&poll_events->max_ns, memory_order_relaxed); + uint64_t accept_total = + atomic_load_explicit(&accepts->total_ns, memory_order_relaxed); + uint64_t read_total = + atomic_load_explicit(&reads->total_ns, memory_order_relaxed); + uint64_t rpc_total = + atomic_load_explicit(&rpcs->total_ns, memory_order_relaxed); + uint64_t write_total = + atomic_load_explicit(&writes->total_ns, memory_order_relaxed); + uint64_t active_count = + atomic_load_explicit(&active->count, memory_order_relaxed); + uint64_t active_total = + atomic_load_explicit(&active->total_ns, memory_order_relaxed); + uint64_t active_max = + atomic_load_explicit(&active->max_ns, memory_order_relaxed); + + if (polls_count == 0 && accept_total == 0 && read_total == 0 && + rpc_total == 0 && write_total == 0 && active_count == 0) { + continue; + } + + if (!printed_workers) { + fprintf(out, "\nworker trace:\n"); + fprintf(out, + " %6s %8s %14s %10s %9s %9s %9s %9s %12s %12s\n", + "worker", "polls", "events/poll", "max_ev", "accepts", + "reads", "rpcs", "writes", "active_avg", "active_max"); + printed_workers = 1; + } + + char active_avg_buf[16]; + char active_max_buf[16]; + format_duration(active_count ? active_total / active_count : 0, + active_avg_buf); + format_duration(active_max, active_max_buf); + double events_per_poll = + polls_count ? (double)polls_total / (double)polls_count : 0.0; + + fprintf(out, + " %6zu %8llu %14.2f %10llu %9llu %9llu %9llu %9llu %12s %12s\n", + worker, (unsigned long long)polls_count, events_per_poll, + (unsigned long long)polls_max, (unsigned long long)accept_total, + (unsigned long long)read_total, (unsigned long long)rpc_total, + (unsigned long long)write_total, active_avg_buf, active_max_buf); + } }