How ds4-server Handles Batched Sessions for Multi-User Inference
ds4-server uses a batched mode where concurrent client requests are collected by a dedicated decode-worker thread, coalesced into a single array, and submitted together to GPU or Metal backends via ds4_sessions_eval_batch, dramatically improving throughput while maintaining thread-safety and atomic failure semantics.
The ds4-server component of the antirez/ds4 repository implements a sophisticated batching system designed to maximize hardware utilization when serving multiple simultaneous users. Instead of processing each inference request individually—which would leave GPU/Metal cores underutilized—the server aggregates pending evaluations and executes them as a single batch operation.
The Three-Stage Batching Pipeline
Multi-user inference in ds4-server flows through three coordinated stages: request intake, batch assembly, and unified backend execution.
Stage 1: Per-Request Entry Point (server_eval_token)
Every client request enters through server_eval_token in ds4_server.c (lines 10801–10809). This function acts as a router that behaves differently depending on the batched_mode configuration:
- Non-batched path: The function acquires
inference_mu, callsds4_session_evaldirectly for single-session evaluation, and returns immediately. - Batched path: The request is recorded in its assigned slot—storing the
decode_token, settingdecode_pending = true, and incrementing the global pending counter. The function then signals the decode-worker viapthread_cond_broadcast(&s->model_cv)and yields control.
This design allows the client thread to return to handling I/O while the heavy inference work is deferred to a specialized worker.
Stage 2: Decode-Worker Thread (decode_worker_main)
The decode_worker_main function (lines 10855–10931 of ds4_server.c) implements the core batching logic. Its operation follows a precise sequence:
-
Wait for work: The worker spins on
while (s->decode_pending == 0 && …), sleeping onpthread_cond_timedwaituntil at least one slot hasdecode_pendingset. -
Optional coalescing: To maximize batch size, the worker can wait an additional configurable window—controlled by
DS4_SERVER_DECODE_COALESCE_US—allowing straggler requests to arrive before proceeding. -
Batch construction: The worker iterates through all slots, collecting those with
decode_pending == true. For each pending slot, it:- Clears
decode_pendingand setsdecode_in_flight = true - Populates a
ds4_decode_itemstruct with the session pointer and token - Stores the slot pointer in a
members[]array for result distribution
- Clears
-
Batch execution: With the
items[count]array assembled, the worker acquiresinference_muand callsds4_sessions_eval_batch(items, count, …). -
Result distribution: After the backend returns, the worker copies
decode_rcanddecode_errinto each slot, setsdecode_done = true, and broadcasts to wake waiting client threads.
Stage 3: Core Library Batch Evaluation (ds4_sessions_eval_batch)
The ds4_sessions_eval_batch function in ds4.c (lines 61333–61372) serves as the gateway to hardware acceleration. Its responsibilities include:
-
Validation: Verifies that all sessions belong to the same
ds4_engine, tokens are within vocabulary bounds, no session appears twice, and no session has exceeded its context window. -
Single-item optimization: If
count == 1, delegates tods4_session_evalto avoid batch overhead. -
Backend dispatch: For multi-item batches, routes to:
ds4_sessions_eval_batch_cudafor NVIDIA GPU executionds4_sessions_eval_batch_metalfor Apple Silicon GPU execution
-
Atomic failure handling: If the backend returns an error, the function invalidates all sessions in the batch, ensuring no partial progress contaminates the system state.
Synchronization Architecture
The batching system employs a two-lock strategy to balance concurrency and correctness:
| Phase | Mechanism | Purpose |
|---|---|---|
| Slot state & pending count | model_mu + pthread_cond_broadcast |
Protects request enqueue/dequeue and worker coordination |
| Inference engine state | inference_mu |
Serializes access to GPU/Metal context and model weights |
| Worker sleep | pthread_cond_timedwait |
Allows microsecond-precision coalescing without busy-waiting |
This separation enables the decode-worker to assemble batches under model_mu while only briefly holding inference_mu during actual execution—minimizing contention on the critical path.
Practical Code Flow
Client thread enqueueing a request:
/* Client thread – enqueue a token for batched evaluation */
int rc = server_eval_token(srv, slot, token, err_buf, sizeof(err_buf));
if (rc != 0) { /* handle error */ }
Worker thread building and executing the batch:
/* Inside decode_worker_main – build the batch */
for (int i = 0; i < s->slot_count; i++) {
server_slot *slot = &s->slots[i];
if (!slot->decode_pending) continue;
slot->decode_pending = false;
slot->decode_in_flight = true;
s->decode_pending--;
members[count] = slot;
items[count].session = slot->session;
items[count].token = slot->decode_token;
count++;
}
/* Execute the batch */
char batch_err[160] = {0};
int rc = ds4_sessions_eval_batch(items, count, batch_err, sizeof(batch_err));
Core library validation and dispatch:
/* ds4_sessions_eval_batch – validation & dispatch */
if (count == 1) return ds4_session_eval(items[0].session, items[0].token, err, errlen);
...
if (e->backend == DS4_BACKEND_CUDA)
return ds4_sessions_eval_batch_cuda(items, count, err, errlen);
...
Key Source Files
| File | Role in Batched Multi-User Inference |
|---|---|
ds4_server.c |
Implements server_eval_token, per-slot state machine, and decode_worker_main batching thread |
ds4.c |
Core library containing ds4_sessions_eval_batch with CUDA/Metal backend dispatch |
ds4_session.c |
Defines ds4_session structure and non-batched ds4_session_eval fallback |
ds4_tp.c |
Thread-pool utilities for slot worker management |
Summary
ds4-serverachieves high-throughput multi-user inference through explicit request batching rather than per-request GPU dispatch.- The decode-worker thread (
decode_worker_main) coalesces pending evaluations with microsecond-precision timing (DS4_SERVER_DECODE_COALESCE_US). - Two-lock synchronization (
model_mufor coordination,inference_mufor engine access) prevents contention while ensuring correctness. ds4_sessions_eval_batchvalidates inputs rigorously, optimizes single-item cases, and guarantees all-or-nothing failure semantics.- Client threads experience blocking semantics identical to non-batched mode, hiding the complexity of asynchronous batch assembly.
Frequently Asked Questions
What is the purpose of DS4_SERVER_DECODE_COALESCE_US?
DS4_SERVER_DECODE_COALESCE_US configures a microsecond delay before batch execution, allowing the decode-worker to accumulate additional pending requests. This increases average batch size and GPU utilization at the cost of minor latency inflation. The trade-off is tunable based on workload characteristics and latency requirements.
How does ds4-server prevent partial failures in batched sessions?
The ds4_sessions_eval_batch function validates all sessions before backend submission and implements all-or-nothing failure handling. If the CUDA or Metal backend returns an error, the function invalidates every session in the batch rather than allowing partial success. This prevents state corruption where some sessions advance their context while others fail.
Can batched and non-batched modes coexist in the same server instance?
No—batching is controlled by a global batched_mode flag at server initialization. When disabled, server_eval_token acquires inference_mu and calls ds4_session_eval synchronously for each request. When enabled, all requests flow through the decode-worker thread. The mode is selected based on expected concurrency and latency requirements.
What happens to client threads while their requests are being batched?
Client threads block on a per-slot condition variable after enqueueing their request. The thread yields after signaling model_cv, then sleeps until decode_worker_main completes the batch and broadcasts model_cv again. From the client's perspective, the latency appears as standard request processing time; the batching optimization is completely transparent.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →