From 58c66ebff2bf18180e7cec997b6a0390bc3df6de Mon Sep 17 00:00:00 2001 From: Jordi Posthumus Date: Fri, 4 Sep 2026 14:58:01 -0300 Subject: [PATCH] server: separate resident sessions from active requests --- README.md | 11 ++++ ds4_help.c | 1 + ds4_server.c | 161 +++++++++++++++++++++++++++++++++++++++++++++++++-- 3 files changed, 168 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index 503f00c63..195343006 100644 --- a/README.md +++ b/README.md @@ -1273,6 +1273,17 @@ request is never evicted. Choose `N` and `--ctx` so all resident KV allocations fit in GPU memory. Without this option, inference retains the original single-session behavior. +`--max-active-requests N` separates resident KV capacity from compute +concurrency. It requires `--batched-session`, may not exceed the resident +session count, and defaults to that count so existing behavior is unchanged. +For example, `--batched-session 10 --max-active-requests 1` allocates ten +resident session slots while processing complete requests one at a time; +additional requests wait in the server queue instead of occupying another +active slot. This avoids mixed prefill/decode contention without reducing +resident KV capacity. Request-to-slot selection is unchanged: this option +controls admission concurrency, not which checkpoint the existing router +chooses for a request. + While generation is active, prefill yields every 128 tokens by default. `--mixed-prefill-quantum N` changes that interval for testing; larger values reduce scheduling handoffs but can make active decoders wait longer. diff --git a/ds4_help.c b/ds4_help.c index 961ca8cb0..9372e480e 100644 --- a/ds4_help.c +++ b/ds4_help.c @@ -345,6 +345,7 @@ static void print_server_api(FILE *fp, const help_colors *c) { opt(fp, c, "--cors", "Add Access-Control-Allow-* headers for browser JS clients."); opt(fp, c, "--trace FILE", "Write prompts, cache decisions, output, and tool calls."); opt(fp, c, "--batched-session N", "Keep N resident sessions and batch decode-ready requests."); + opt(fp, c, "--max-active-requests N", "Limit concurrently assigned requests without reducing resident slot count."); opt(fp, c, "--mixed-prefill-quantum N", "Prefill chunk while generations are active. Default: 128; GLM-5.3 minimum: 1024"); para(fp, c, "Endpoints: /v1/chat/completions, /v1/responses, /v1/completions, and /v1/messages."); para(fp, c, "Model endpoint aliases include deepseek-v4-flash and deepseek-v4-pro; both serve the loaded GGUF."); diff --git a/ds4_server.c b/ds4_server.c index 50efbe6bc..cf56a37af 100644 --- a/ds4_server.c +++ b/ds4_server.c @@ -9099,6 +9099,10 @@ struct server { ds4_tp *tp_leader; server_slot *slots; int slot_count; + /* Resident slot allocation and compute concurrency are intentionally + * independent. busy covers both assigned and running jobs, so this cap + * serializes complete requests. Slot selection is unchanged. */ + int max_active_requests; int ctx_size; bool batched_mode; pthread_t *slot_threads; @@ -13230,9 +13234,28 @@ static int job_slot_score(server *s, server_slot *slot, const job *j, return common; } +/* Requires s->mu. A busy slot owns an assigned or running request for its + * whole lifetime, including prefill, decode, streaming, and cancellation. */ +static int active_request_count_locked(const server *s) { + if (!s) return 0; + int active = 0; + for (int i = 0; i < s->slot_count; i++) { + if (s->slots[i].busy) active++; + } + return active; +} + +static int active_request_limit(const server *s) { + if (!s || s->slot_count <= 0) return 0; + return s->max_active_requests > 0 + ? s->max_active_requests : s->slot_count; +} + static void dispatch_jobs_locked(server *s) { if (!s || !s->batched_mode) return; - for (;;) { + int active = active_request_count_locked(s); + const int active_limit = active_request_limit(s); + while (active < active_limit) { job *chosen = NULL; job *chosen_prev = NULL; server_slot *chosen_slot = NULL; @@ -13269,6 +13292,7 @@ static void dispatch_jobs_locked(server *s) { chosen->next = NULL; chosen_slot->assigned = chosen; chosen_slot->busy = true; + active++; pthread_cond_broadcast(&s->cv); } } @@ -13771,6 +13795,7 @@ typedef struct { int tool_memory_max_ids; bool enable_cors; int batched_sessions; + int max_active_requests; int mixed_prefill_quantum; } server_config; @@ -14005,6 +14030,9 @@ static server_config parse_options(int argc, char **argv) { c.trace_path = need_arg(&i, argc, argv, arg); } else if (!strcmp(arg, "--batched-session")) { c.batched_sessions = parse_int_arg(need_arg(&i, argc, argv, arg), arg); + } else if (!strcmp(arg, "--max-active-requests")) { + c.max_active_requests = + parse_int_arg(need_arg(&i, argc, argv, arg), arg); } else if (!strcmp(arg, "--mixed-prefill-quantum")) { c.mixed_prefill_quantum = parse_int_arg(need_arg(&i, argc, argv, arg), arg); @@ -14120,6 +14148,16 @@ static server_config parse_options(int argc, char **argv) { "ds4-server: --kv-cache-cold-max-tokens must be 0 or >= --kv-cache-min-tokens"); exit(2); } + if (c.max_active_requests > 0 && c.batched_sessions <= 0) { + server_log(DS4_LOG_DEFAULT, + "ds4-server: --max-active-requests requires --batched-session"); + exit(2); + } + if (c.max_active_requests > c.batched_sessions) { + server_log(DS4_LOG_DEFAULT, + "ds4-server: --max-active-requests cannot exceed --batched-session"); + exit(2); + } if (c.engine.directional_steering_file && !directional_steering_scale_set) { c.engine.directional_steering_ffn = 1.0f; } @@ -14271,6 +14309,8 @@ int main(int argc, char **argv) { s.tp_leader = tp_leader; s.ctx_size = cfg.ctx_size; s.slot_count = slot_count; + s.max_active_requests = cfg.max_active_requests > 0 + ? cfg.max_active_requests : slot_count; s.batched_mode = cfg.batched_sessions > 0; s.mixed_prefill_quantum = cfg.mixed_prefill_quantum; s.last_prefill_slot = slot_count - 1; @@ -14322,8 +14362,9 @@ int main(int argc, char **argv) { } if (s.batched_mode) { server_log(DS4_LOG_DEFAULT, - "ds4-server: batched mode enabled resident_sessions=%d prefill_quantum=%d mixed_prefill_quantum=%d decode_coalesce_us=%ld", + "ds4-server: batched mode enabled resident_sessions=%d max_active_requests=%d prefill_quantum=%d mixed_prefill_quantum=%d decode_coalesce_us=%ld", s.slot_count, + s.max_active_requests, server_prefill_quantum_for(&s, false), server_prefill_quantum_for(&s, true), server_decode_coalesce_us()); @@ -14519,6 +14560,100 @@ static void test_mixed_prefill_quantum_option(void) { TEST_ASSERT(server_prefill_quantum_for(&s, true) == 128); } +static void test_max_active_requests_option(void) { + char *default_argv[] = {"ds4-server"}; + server_config defaults = parse_options(1, default_argv); + TEST_ASSERT(defaults.batched_sessions == 0); + TEST_ASSERT(defaults.max_active_requests == 0); + + char *custom_argv[] = { + "ds4-server", "--batched-session", "10", + "--max-active-requests", "1" + }; + server_config custom = parse_options(5, custom_argv); + TEST_ASSERT(custom.batched_sessions == 10); + TEST_ASSERT(custom.max_active_requests == 1); +} + +static void test_dispatch_respects_active_request_limit(void) { + server s = {0}; + pthread_mutex_init(&s.mu, NULL); + pthread_mutex_init(&s.tool_mu, NULL); + pthread_cond_init(&s.cv, NULL); + s.batched_mode = true; + server_slot slots[3] = {0}; + s.slots = slots; + s.slot_count = 3; + TEST_ASSERT(active_request_limit(&s) == 3); + s.max_active_requests = 1; + + job first = {0}, second = {0}, third = {0}; + job *jobs[3] = {&first, &second, &third}; + const char *call_ids[3] = {"call-0", "call-1", "call-2"}; + for (int i = 0; i < 3; i++) { + slots[i].id = i; + slots[i].responses_live.valid = true; + id_list_push_unique(&slots[i].responses_live.call_ids, call_ids[i]); + jobs[i]->req.responses_requires_live_tool_state = true; + id_list_push_unique(&jobs[i]->req.responses_live_call_ids, + call_ids[i]); + } + first.next = &second; + second.next = &third; + s.head = &first; + s.tail = &third; + + pthread_mutex_lock(&s.mu); + dispatch_jobs_locked(&s); + TEST_ASSERT(slots[0].assigned == &first); + TEST_ASSERT(slots[0].busy); + TEST_ASSERT(slots[1].assigned == NULL); + TEST_ASSERT(slots[2].assigned == NULL); + TEST_ASSERT(s.head == &second); + TEST_ASSERT(s.tail == &third); + TEST_ASSERT(active_request_count_locked(&s) == 1); + + /* A worker clears assigned before running the job, but busy must retain + * ownership and keep queued requests from entering other slots. */ + slots[0].assigned = NULL; + dispatch_jobs_locked(&s); + TEST_ASSERT(slots[1].assigned == NULL); + TEST_ASSERT(slots[2].assigned == NULL); + TEST_ASSERT(s.head == &second); + + /* Completing the first request opens exactly one queue position. */ + slots[0].busy = false; + dispatch_jobs_locked(&s); + TEST_ASSERT(slots[1].assigned == &second); + TEST_ASSERT(slots[1].busy); + TEST_ASSERT(s.head == &third); + TEST_ASSERT(s.tail == &third); + TEST_ASSERT(active_request_count_locked(&s) == 1); + + /* Raising the cap to two admits exactly one more job while the second + * worker owns its claimed request. */ + slots[1].assigned = NULL; + s.max_active_requests = 2; + dispatch_jobs_locked(&s); + TEST_ASSERT(slots[2].assigned == &third); + TEST_ASSERT(slots[2].busy); + TEST_ASSERT(s.head == NULL); + TEST_ASSERT(s.tail == NULL); + TEST_ASSERT(active_request_count_locked(&s) == 2); + + slots[1].busy = false; + slots[2].assigned = NULL; + slots[2].busy = false; + pthread_mutex_unlock(&s.mu); + for (int i = 0; i < 3; i++) { + live_tool_state_free(&slots[i].responses_live); + request_free(&jobs[i]->req); + } + pthread_mutex_destroy(&s.mu); + pthread_mutex_destroy(&s.tool_mu); + pthread_cond_destroy(&s.cv); +} + static void test_multimodal_prefill_resume_frontier(void) { TEST_ASSERT(server_multimodal_resume_frontier(160, 160, 170, true) == 160); TEST_ASSERT(server_multimodal_resume_frontier(160, 159, 170, true) == 0); @@ -18132,20 +18267,34 @@ static void test_cancel_unlinks_queued_jobs(void) { static void test_cancel_detaches_assigned_job(void) { server s; server_slot slot = {0}; - job j; + job j, queued; test_cancel_server_init(&s); test_cancel_job_init(&j); + test_cancel_job_init(&queued); s.batched_mode = true; + s.max_active_requests = 1; test_server_bind_slot(&s, &slot); + slot.responses_live.valid = true; + id_list_push_unique(&slot.responses_live.call_ids, "call-next"); + queued.req.responses_requires_live_tool_state = true; + id_list_push_unique(&queued.req.responses_live_call_ids, "call-next"); slot.assigned = &j; slot.busy = true; + s.head = s.tail = &queued; server_cancel_job(&s, &j); TEST_ASSERT(job_cancelled(&j)); TEST_ASSERT(j.done); - TEST_ASSERT(slot.assigned == NULL); - TEST_ASSERT(!slot.busy); + TEST_ASSERT(slot.assigned == &queued); + TEST_ASSERT(slot.busy); + TEST_ASSERT(s.head == NULL); + TEST_ASSERT(s.tail == NULL); + slot.assigned = NULL; + slot.busy = false; + live_tool_state_free(&slot.responses_live); + request_free(&queued.req); + test_cancel_job_destroy(&queued); test_cancel_job_destroy(&j); test_cancel_server_destroy(&s); } @@ -19405,6 +19554,8 @@ static void test_responses_inline_image_content(void) { static void ds4_server_unit_tests_run(void) { test_batched_prefill_round_robin(); test_mixed_prefill_quantum_option(); + test_max_active_requests_option(); + test_dispatch_respects_active_request_limit(); test_multimodal_prefill_resume_frontier(); test_batched_live_continuation_slot_binding(); test_request_defaults_use_min_p_filtering();