36. An OpenAI-compatible server
In this chapter
- The architecture of a serving process: an engine core that only schedules and runs the model, and an asyncio frontend that owns text, HTTP and streaming.
- Why production servers run the engine core in its own process, and how the two halves talk.
- The OpenAI API in detail: completions and chat completions, streaming with server-sent events, tools, JSON output, logprobs,
n, errors and usage. - The behaviors that separate a demo from a service: cancellation when clients leave, backpressure, error isolation, metrics and health.
You will build
core_loop, AsyncLLM._dispatch and AsyncLLM.generate (engine/serve/async_engine.py), and Server.sampling_params and Server.chat_prompt (engine/serve/api.py). Then python serve.py --model-dir ... serves a checkpoint to any OpenAI client.
Time: 6-8 hours. GPU: not needed.
One request’s journey
A chat request to the finished server passes through these stages:
client ──HTTP──> FastAPI route ──> Server.chat
│ validate, apply defaults, render the chat template,
│ tokenize, build SamplingParams (+ grammar for tools / JSON)
▼
AsyncLLM.generate ──("add", rid, ids, params)──> engine core (own process)
▲ │ Scheduler, ModelRunner,
│ per-request asyncio.Queue │ Sampler: Chapters 31-34
│ ▼
AsyncLLM._dispatch <──("outputs", [RequestOutput...])──
│ detokenize, stop strings, metrics
▼
OutputParser (reasoning, tool calls) ──SSE chunks──> client
The rest of Part VIII’s machinery sits inside the engine-core box. This chapter builds everything around it.
Two halves
The engine core is a tight loop: plan, run, sample, update, repeat. Anything that delays it delays every request’s next token. The frontend does work that’s unpredictable in size: parsing JSON bodies, rendering templates, tokenizing a 100,000-token document, detokenizing, serializing streamed chunks for hundreds of connections. In one Python process, the two compete for a single interpreter. CPython runs one thread’s bytecode at a time (the GIL), and Chapter 33 measured a millisecond of host work per step at 256 requests, so a frontend that tokenizes a long prompt can easily stall the decode loop for several steps.
So production servers split them. vLLM V1 runs an API-server process and an EngineCore process connected by ZeroMQ sockets; SGLang runs a tokenizer manager, a scheduler and a detokenizer manager as separate processes. This book’s AsyncLLM does the same with two queues, and can run the core either in a thread (simple to debug, and what the tests use) or in its own process (mode="process", the default for serve.py):
def core_loop(factory, inputs, outputs, stats_every=8):
"""Run an EngineCore until told to stop. (Your engine: Chapter 36)
Block for input only when idle; otherwise take whatever messages are waiting and step.
An invalid request fails alone ("error"); an exception inside step() fails every request
("fatal"), because the engine's state can no longer be trusted.
"""
try:
core = factory()
except Exception as error: # noqa: BLE001
outputs.put(("fatal", f"engine failed to start: {error!r}"))
return
outputs.put(("ready", {"max_model_len": core.max_model_len, "num_blocks": core.blocks.num_blocks}))
busy = False
while True:
if busy and not core.has_unfinished: # just went idle (finished or aborted): report it
outputs.put(("stats", core.stats()))
busy = core.has_unfinished
block = not busy
while True:
try:
kind, payload = inputs.get(block=block)
except queue.Empty:
break
block = False
if kind == "shutdown":
return
if kind == "abort":
core.abort_request(payload)
elif kind == "add":
rid, prompt, params, priority, options = payload
try:
core.add_request(rid, prompt, params, priority, **options)
except ValueError as error:
outputs.put(("error", (rid, str(error))))
if core.has_unfinished:
try:
step_outputs = core.step()
except Exception as error: # noqa: BLE001
outputs.put(("fatal", f"engine step failed: {error!r}"))
return
if step_outputs:
outputs.put(("outputs", step_outputs))
if core.steps % stats_every == 0:
outputs.put(("stats", core.stats()))
busy = True
The loop blocks only when idle. While requests are running it takes whatever messages are waiting and steps, so a burst of new requests is admitted at the next step without waiting for the batch to drain.
It also distinguishes two kinds of failure. A request the engine rejects (too long, a guided-decoding pattern it can’t compile) fails alone: the core reports ("error", rid, message) and keeps serving everyone else. An exception inside step() is different: the engine’s state, its block tables and its num_computed_tokens, may now be inconsistent, so the core reports "fatal", every in-flight request fails, and /health turns red so that an orchestrator (Kubernetes, systemd) restarts the process. Trying to limp on after a failed step risks serving one request’s cache to another.
In process mode, the factory that builds the engine is picklable (functools.partial of a module-level function with a checkpoint path), and the model is loaded in the child. Weights are never copied between processes, and the child process is started with spawn, because a forked process inherits CUDA state that it can’t use.
The frontend
AsyncLLM runs in the server’s event loop. Each request gets an asyncio.Queue; a background thread receives the core’s messages and hands them to the loop with call_soon_threadsafe, where _dispatch turns token IDs into text:
class AsyncLLM:
def __init__(self, factory, tokenizer, mode="thread", max_pending=1024, metrics=None,
request_timeout=120, max_buffered_updates=256):
if mode not in ("thread", "process") or max_pending < 1 or request_timeout <= 0 or max_buffered_updates < 1:
raise ValueError("Need thread|process mode and positive admission/deadline/buffer limits")
self.tokenizer, self.max_pending, self.mode = tokenizer, max_pending, mode
context = mp.get_context("spawn")
self.inputs = context.Queue() if mode == "process" else queue.Queue()
self.outputs = context.Queue() if mode == "process" else queue.Queue()
target = (context.Process if mode == "process" else threading.Thread)
self.worker = target(target=run_core, args=(factory, self.inputs, self.outputs), daemon=True)
self.requests, self.core_stats, self.info = {}, {}, None
self.metrics = metrics
self.request_timeout, self.max_buffered_updates = request_timeout, max_buffered_updates
self.closing = False
self.dead = None
self.ids = itertools.count()
async def start(self):
"""Start the core, and a daemon thread that hands the core's messages to the event loop;
wait until the model is loaded. (A daemon thread, not a loop executor: a thread blocked
on queue.get() must never keep the process from exiting.)"""
self.loop = asyncio.get_running_loop()
self.ready = self.loop.create_future()
self.worker.start()
threading.Thread(target=self._pump, daemon=True).start()
kind, payload = await self.ready
if kind != "ready":
self.dead = payload
raise EngineDead(payload)
self.info = payload
async def stop(self, drain_timeout=30):
"""Reject new work, give existing streams a deadline, then abort and join the worker."""
self.closing = True
deadline = time.monotonic() + drain_timeout
while self.requests and time.monotonic() < deadline:
await asyncio.sleep(0.01)
for rid, state in list(self.requests.items()):
self.abort(rid)
state.queue.put_nowait(EngineDead("Server drain deadline exceeded"))
self.inputs.put(("shutdown", None))
await asyncio.to_thread(self.worker.join, 2)
if self.mode == "process" and self.worker.is_alive():
self.worker.terminate()
await asyncio.to_thread(self.worker.join, 2)
self.outputs.put(("closed", None)) # ends the pump
await asyncio.sleep(0)
def _pump(self):
while True:
kind, payload = self.outputs.get()
if kind == "closed":
return
self.loop.call_soon_threadsafe(self._handle, kind, payload)
if kind == "fatal":
return
def _handle(self, kind, payload):
if not self.ready.done():
self.ready.set_result((kind, payload))
if kind == "ready":
return
if kind == "outputs":
for out in payload:
self._dispatch(out)
elif kind == "stats":
self.core_stats = payload
elif kind == "error":
rid, message = payload
state = self.requests.get(rid)
if state is not None:
state.queue.put_nowait(ValueError(message))
elif kind == "fatal":
self.dead = payload
for state in self.requests.values():
state.queue.put_nowait(EngineDead(payload))
def _dispatch(self, out: RequestOutput):
"""Token ids -> text for one request: detokenize, check stop strings, time it. (Your engine: Chapter 36)"""
state = self.requests.get(out.request_id)
if state is None or state.finished:
return # aborted, or already stopped by a stop string
if state.queue.qsize() >= self.max_buffered_updates:
self.abort(out.request_id)
state.finished = True
while not state.queue.empty():
state.queue.get_nowait()
state.queue.put_nowait(Overloaded("Client is not consuming streamed updates"))
return
now = time.monotonic()
if out.new_token_ids:
if state.first_token is None:
state.first_token = now
if self.metrics:
self.metrics.observe("time_to_first_token_seconds", now - state.arrival)
elif self.metrics:
self.metrics.observe("inter_token_latency_seconds", (now - state.last_token) / len(out.new_token_ids))
state.last_token = now
text = state.detok.add(out.new_token_ids)
if out.finished:
text += state.detok.flush()
visible = state.stop.add(text)
finished, reason, stop_reason = out.finished, out.finish_reason, out.stop_reason
if state.stop.stopped is not None:
if not out.finished:
self.abort(out.request_id) # the core doesn't know about stop strings
finished, reason, stop_reason = True, "stop", state.stop.stopped
elif finished:
visible += state.stop.finish()
state.finished = finished
if finished and self.metrics:
self.metrics.inc("request_success_total", labels={"finished_reason": reason})
self.metrics.observe("e2e_request_latency_seconds", now - state.arrival)
self.metrics.inc("prompt_tokens_total", out.num_prompt_tokens)
self.metrics.inc("generation_tokens_total", out.num_output_tokens)
state.queue.put_nowait(TextOutput(out.request_id, visible, out.new_token_ids, finished, reason, stop_reason,
out.logprobs, out.prompt_logprobs, out.num_prompt_tokens,
out.num_output_tokens, out.num_cached_tokens))
async def generate(self, prompt_ids, params, request_id=None, priority=0, lora=None, features=None):
"""Async iterator of TextOutputs for one request (n must be 1; the API layer fans out). (Your engine: Chapter 36)
Leaving the loop early (a client disconnect cancels the task) aborts the request in the
core, so its KV blocks are freed at once.
"""
if self.dead:
raise EngineDead(self.dead)
if self.closing:
raise Overloaded("Server is draining")
if len(self.requests) >= self.max_pending:
raise Overloaded(f"{len(self.requests)} requests in flight")
rid = request_id or f"req-{next(self.ids)}"
if rid in self.requests:
raise ValueError(f"Duplicate request id {rid!r}")
state = RequestState(asyncio.Queue(), IncrementalDetokenizer(self.tokenizer, params.skip_special_tokens),
StopChecker(params.stop, params.include_stop_str_in_output))
self.requests[rid] = state
self.inputs.put(("add", (rid, list(prompt_ids), params, priority, {"lora": lora, "features": features})))
deadline = time.monotonic() + self.request_timeout
try:
while True:
item = await asyncio.wait_for(state.queue.get(), max(0, deadline - time.monotonic()))
if isinstance(item, Exception):
raise item
yield item
if item.finished:
return
finally:
if not state.finished:
state.finished = True
self.abort(rid)
self.requests.pop(rid, None)
def abort(self, request_id):
self.inputs.put(("abort", request_id))
@property
def num_in_flight(self):
return len(self.requests)
Three responsibilities live here because they need text:
- Detokenization with Chapter 35’s incremental UTF-8 decoder, one per request.
- Stop strings, which the core can’t see. When the frontend finds one, it truncates the text, marks the request finished with
finish_reason="stop"and the matched string asstop_reason, and sends("abort", rid)to the core so that its blocks are freed at once. - Latency metrics: time to first token (arrival to first output), inter-token latency, end-to-end latency, measured where the client would measure them.
The reader is a daemon thread, not a task that calls queue.get() through the event loop’s executor. That’s a detail that bit this chapter’s first draft: a test failed an assertion before stopping the engine, and the interpreter then hung forever, because asyncio.run waits for its executor’s threads at exit, and one of them was blocked in queue.get(). A server whose shutdown can hang can’t be restarted cleanly; a daemon thread never holds the process.
The API
OpenAI’s API is the de facto standard: every client library, agent framework, evaluation harness and benchmark tool speaks it. Matching it closely is worth more than any feature you could add.
Requests and defaults
The request schemas list OpenAI’s fields plus vLLM’s extensions (top_k, min_p, repetition_penalty, min_tokens, ignore_eos, guided_json, guided_regex, priority). Unknown fields are ignored, because clients send fields like user and parallel_tool_calls that a server may not use:
class SamplingFields(BaseModel):
"""Fields shared by both endpoints. None means "the server's default" (generation_config.json)."""
model_config = ConfigDict(extra="ignore")
model: str | None = None
max_tokens: int | None = None
temperature: float | None = None
top_p: float | None = None
top_k: int | None = None
min_p: float | None = None
repetition_penalty: float | None = None
presence_penalty: float = 0.0
frequency_penalty: float = 0.0
n: int = 1
seed: int | None = None
stop: str | list[str] | None = None
stop_token_ids: list[int] | None = None
include_stop_str_in_output: bool = False
ignore_eos: bool = False
min_tokens: int = 0
logit_bias: dict[str, float] | None = None
skip_special_tokens: bool = True
stream: bool = False
stream_options: dict | None = None
guided_json: dict | str | None = None
guided_regex: str | None = None
priority: int = 0
lora: str | None = None
class CompletionRequest(SamplingFields):
prompt: str | list[str] | list[int] | list[list[int]]
logprobs: int | None = None
prompt_logprobs: int | None = None
echo: bool = False
class ChatRequest(SamplingFields):
messages: list[dict]
tools: list[dict] | None = None
tool_choice: str | dict | None = None
response_format: dict | None = None
logprobs: bool = False
top_logprobs: int | None = None
max_completion_tokens: int | None = None
chat_template_kwargs: dict | None = None
Sampling fields default to None, meaning “the server’s default”, and the server’s defaults come from the checkpoint’s generation_config.json. Qwen3’s instruct models, for example, ship temperature=0.6, top_p=0.95, top_k=20, and its authors recommend against greedy decoding for thinking mode. A request that sets nothing should get what the model’s authors intended, not OpenAI’s default of temperature=1.0.
def sampling_params(self, req, prompt_len, default_max, logprobs=None, prompt_logprobs=None,
guided_regex=None, guided_json=None):
"""Request fields -> SamplingParams, with server defaults and validation. (Your engine: Chapter 36)"""
max_tokens = req.max_tokens if req.max_tokens is not None else default_max
if max_tokens is None:
max_tokens = self.max_model_len - prompt_len
if prompt_len + max_tokens > self.max_model_len:
raise APIError(400, f"This model's maximum context length is {self.max_model_len} tokens; the request has "
f"{prompt_len} prompt tokens and asks for {max_tokens} more.", param="max_tokens")
pick = lambda name: getattr(req, name) if getattr(req, name) is not None else self.defaults[name] # noqa: E731
try:
return SamplingParams(
max_tokens=max_tokens, temperature=pick("temperature"), top_p=pick("top_p"), top_k=pick("top_k"),
min_p=pick("min_p"), repetition_penalty=pick("repetition_penalty"), seed=req.seed, n=req.n,
presence_penalty=req.presence_penalty, frequency_penalty=req.frequency_penalty,
stop=req.stop or (), include_stop_str_in_output=req.include_stop_str_in_output,
stop_token_ids=tuple(req.stop_token_ids or ()) + self.default_stop_token_ids,
ignore_eos=req.ignore_eos, min_tokens=req.min_tokens,
logit_bias={int(k): v for k, v in (req.logit_bias or {}).items()},
logprobs=logprobs, prompt_logprobs=prompt_logprobs, skip_special_tokens=req.skip_special_tokens,
guided_regex=guided_regex or req.guided_regex, guided_json=guided_json or req.guided_json)
except ValueError as error:
raise APIError(400, str(error)) from None
max_tokens defaults to 16 for completions (OpenAI’s legacy default) and to the rest of the context for chat. A request that can’t fit is rejected with a 400 that says why, in OpenAI’s error format: {"error": {"message", "type", "param", "code"}}. Clients parse that format and some retry on it, so every error path uses it, including validation errors from SamplingParams (a negative temperature is the client’s mistake, not a server crash).
The generation config’s eos_token_id list is added to every request’s stop tokens. Qwen3 lists both <|im_end|> (end of turn) and <|endoftext|>; a server that stops only on the tokenizer’s EOS lets the model ramble past the end of its answer.
Completions
async def completions(self, req, http_request):
self.check(http_request, req.model)
self.check_lora(req.lora)
prompt = req.prompt
if isinstance(prompt, str) or (prompt and isinstance(prompt[0], int)):
prompt = [prompt]
prompts = [self.tokenizer.encode(p) if isinstance(p, str) else list(p) for p in prompt]
if not prompts or any(not p for p in prompts):
raise APIError(400, "The prompt is empty", param="prompt")
params = [self.sampling_params(req, len(p), 16, req.logprobs, req.prompt_logprobs) for p in prompts]
rid, created = f"cmpl-{uuid.uuid4().hex}", int(time.time())
self.admit(len(prompts) * params[0].n)
streams = []
for i, (ids, p) in enumerate(zip(prompts, params)):
streams += [self.llm.generate(ids, p.child(j) if p.n > 1 else p, f"{rid}-{i}-{j}", req.priority,
req.lora) for j in range(p.n)]
n = params[0].n
base = {"id": rid, "object": "text_completion", "created": created, "model": self.model_name}
if req.stream:
return self.sse(self.completion_chunks(streams, base, n, prompts, req))
texts, finish, logprobs, usage = [""] * len(streams), [None] * len(streams), [[] for _ in streams], [0, 0, 0]
async for k, out in merge(streams):
texts[k] += out.text
logprobs[k] += out.logprobs or []
if out.finished:
finish[k] = (out.finish_reason, out.stop_reason)
usage[1] += out.num_output_tokens
usage[2] += out.num_cached_tokens
usage[0] = sum(len(p) for p in prompts)
choices = [{"index": k, "text": (req.prompt if req.echo and isinstance(req.prompt, str) else "") + texts[k],
"logprobs": self.completion_logprobs(logprobs[k]) if req.logprobs is not None else None,
"finish_reason": finish[k][0], "stop_reason": finish[k][1]} for k in range(len(streams))]
return {**base, "choices": choices, "usage": self.usage(*usage)}
async def completion_chunks(self, streams, base, n, prompts, req):
completion_tokens = cached = 0
async for k, out in merge(streams):
lp = self.completion_logprobs(out.logprobs) if out.logprobs and req.logprobs is not None else None
choice = {"index": k, "text": out.text, "logprobs": lp,
"finish_reason": out.finish_reason if out.finished else None}
yield {**base, "choices": [choice]}
if out.finished:
completion_tokens += out.num_output_tokens
cached += out.num_cached_tokens
if (req.stream_options or {}).get("include_usage"):
yield {**base, "choices": [], "usage": self.usage(sum(len(p) for p in prompts), completion_tokens, cached)}
prompt can be text, token IDs, or a list of either: a batch. With n samples per prompt, one HTTP request becomes len(prompts) × n engine requests, with choice index prompt × n + sample. merge interleaves their streams as items arrive, so a streaming response sends each choice’s chunks as soon as they exist rather than one choice after another.
Chat completions
def chat_prompt(self, req):
"""Messages -> token ids, plus the constraints the request implies. (Your engine: Chapter 36)"""
messages = []
for m in req.messages:
content = m.get("content")
if isinstance(content, list): # [{"type": "text", "text": ...}, ...]
parts = []
for part in content:
if part.get("type") == "text":
parts.append(part.get("text", ""))
elif part.get("type") == "image_url" and self.multimodal_processor is not None:
parts.append(self.multimodal_processor.marker)
else:
raise APIError(400, "This server does not support that content part", param="messages")
content = "".join(parts)
messages.append({**m, "content": content if content is not None else ""})
tools = req.tools if req.tools and req.tool_choice != "none" else None
guided_regex = guided_json = None
if tools and req.tool_choice == "required":
guided_regex = tool_call_pattern(tools)
elif tools and isinstance(req.tool_choice, dict):
guided_regex = tool_call_pattern(tools, req.tool_choice["function"]["name"])
fmt = req.response_format or {}
if fmt.get("type") == "json_schema":
guided_json = fmt["json_schema"].get("schema", {})
elif fmt.get("type") == "json_object":
guided_regex = json_object_regex()
try:
text = render_chat(self.chat_template, messages, tools=tools, add_generation_prompt=True,
**(req.chat_template_kwargs or {}))
except Exception as error: # noqa: BLE001 (template errors are the client's)
raise APIError(400, f"Chat template error: {error}") from None
return self.tokenizer.encode(text), tools is not None, guided_regex, guided_json
async def chat(self, req, http_request):
self.check(http_request, req.model)
self.check_lora(req.lora)
ids, use_tools, guided_regex, guided_json = self.chat_prompt(req)
images = [part for m in req.messages if isinstance(m.get("content"), list)
for part in m["content"] if part.get("type") == "image_url"]
features = None
if images:
ids, features = await asyncio.to_thread(self.multimodal_processor, ids, images)
if req.max_completion_tokens is not None:
req.max_tokens = req.max_completion_tokens
logprobs = (req.top_logprobs or 0) if req.logprobs else None
params = self.sampling_params(req, len(ids), None, logprobs, None, guided_regex, guided_json)
rid, created = f"chatcmpl-{uuid.uuid4().hex}", int(time.time())
streams = self.streams([ids], params, rid, req.priority, req.lora, features)
parsers = [OutputParser(self.reasoning, use_tools) for _ in streams]
base = {"id": rid, "created": created, "model": self.model_name}
if req.stream:
return self.sse(self.chat_chunks(streams, parsers, base, len(ids), req))
messages = [{"role": "assistant", "content": ""} for _ in streams]
finish, entries, completion, cached = [None] * len(streams), [[] for _ in streams], 0, 0
async for k, out in merge(streams):
for delta in parsers[k].feed(out.text, final=out.finished):
self.accumulate(messages[k], delta)
entries[k] += out.logprobs or []
if out.finished:
finish[k] = "tool_calls" if messages[k].get("tool_calls") else out.finish_reason
completion += out.num_output_tokens
cached += out.num_cached_tokens
for m in messages:
if m.get("tool_calls") and not m["content"]:
m["content"] = None
choices = [{"index": k, "message": messages[k], "finish_reason": finish[k],
"logprobs": self.chat_logprobs(entries[k]) if req.logprobs else None} for k in range(len(streams))]
return {**base, "object": "chat.completion", "choices": choices,
"usage": self.usage(len(ids), completion, cached)}
async def chat_chunks(self, streams, parsers, base, prompt_len, req):
base = {**base, "object": "chat.completion.chunk"}
for k in range(len(streams)):
yield {**base, "choices": [{"index": k, "delta": {"role": "assistant", "content": ""}, "finish_reason": None}]}
completion, cached, called = 0, 0, [False] * len(streams)
async for k, out in merge(streams):
lp = self.chat_logprobs(out.logprobs) if out.logprobs and req.logprobs else None
deltas = parsers[k].feed(out.text, final=out.finished) or ([{}] if lp else [])
for delta in deltas:
called[k] |= "tool_calls" in delta
yield {**base, "choices": [{"index": k, "delta": delta, "logprobs": lp, "finish_reason": None}]}
lp = None
if out.finished:
reason = "tool_calls" if called[k] else out.finish_reason
yield {**base, "choices": [{"index": k, "delta": {}, "finish_reason": reason}]}
completion += out.num_output_tokens
cached += out.num_cached_tokens
if (req.stream_options or {}).get("include_usage"):
yield {**base, "choices": [], "usage": self.usage(prompt_len, completion, cached)}
chat_prompt normalizes messages (OpenAI allows content to be a list of typed parts; Chapter 43 adds images), renders the chat template with the tools, and turns two request fields into grammars with Chapter 34’s machinery:
| request | constraint |
|---|---|
tool_choice: "required" | the output is one of the tools’ call formats, arguments matching each tool’s schema |
tool_choice: {"function": {"name": "get_weather"}} | that tool’s call, with valid arguments |
response_format: {"type": "json_schema", ...} | the schema’s regex |
response_format: {"type": "json_object"} | any JSON object, nested up to three levels |
Without a constraint (tool_choice: "auto"), the model decides whether to call a tool, and the OutputParser of Chapter 35 separates calls, reasoning and content as the text streams in. A choice that contains a tool call finishes with finish_reason: "tool_calls", and its content is null if there was no other text, as OpenAI clients expect.
Streaming
A streamed response is a sequence of server-sent events: data: {json} lines separated by blank lines, ending with data: [DONE]. For chat, the first chunk of each choice carries {"role": "assistant", "content": ""}, later chunks carry content, reasoning_content or tool_calls deltas, a final chunk carries the finish_reason, and if the client asked for stream_options: {"include_usage": true}, one more chunk with empty choices carries the token counts, including prompt_tokens_details.cached_tokens, the prefix-cache hits that some providers bill at a discount.
An error after streaming has started can’t change the HTTP status, which was 200 and is long gone. It’s sent as an error event before [DONE]. That’s why the server does its admission check before returning the streaming response.
Behaviors that make it a service
Cancellation
Users close tabs; agents time out; load balancers drop connections. Each abandoned request still holds KV blocks and still costs a slot in every batch until it reaches max_tokens. A busy server can spend a large fraction of its GPU on output nobody will read.
When a client disconnects during a streamed response, the server’s response task is cancelled. Cancellation propagates into merge, which cancels its pumps, which raises CancelledError inside each AsyncLLM.generate, whose finally block sends ("abort", rid). The core frees the request’s blocks at its next loop. The test checks the whole chain: start a 500-token request, read one chunk, close the stream, and within a few steps the core reports zero running requests and an empty KV pool.
Backpressure
Accepting every request isn’t kind to clients: past some load, queueing delay grows without bound and every request times out. max_pending caps the requests in flight; beyond it the server answers 503, and well-behaved clients back off and retry, or a load balancer sends the request to another replica (Chapter 41). The right cap comes from measurement (Chapter 44): the load at which time to first token breaks its target.
Observability
/metrics serves Prometheus’s text format, with names that mirror vLLM’s (izh: instead of vllm:), so existing dashboards translate directly:
"""Prometheus metrics for the server (Chapter 36), in the text exposition format, without a client
library. Names follow vLLM's (vllm:...) with an izh: prefix, so existing dashboards translate."""
import bisect
import threading
LATENCY_BUCKETS = (0.001, 0.005, 0.01, 0.02, 0.04, 0.06, 0.08, 0.1, 0.25, 0.5, 0.75, 1.0, 2.5, 5.0, 7.5, 10.0,
20.0, 40.0, 80.0)
HELP = {
"time_to_first_token_seconds": "Time from arrival to the first output token.",
"inter_token_latency_seconds": "Time between output tokens of one request.",
"e2e_request_latency_seconds": "Time from arrival to completion.",
"request_success_total": "Finished requests by finish reason.",
"prompt_tokens_total": "Prompt tokens of finished requests.",
"generation_tokens_total": "Output tokens of finished requests.",
"num_requests_running": "Requests in the running batch.",
"num_requests_waiting": "Requests waiting for admission.",
"kv_cache_usage_perc": "Fraction of KV blocks in use.",
"prefix_cache_hit_rate": "Fraction of prompt tokens served by the prefix cache.",
"requests_in_flight": "Requests the server is streaming.",
}
class Metrics:
def __init__(self, prefix="izh:"):
self.prefix, self.lock = prefix, threading.Lock()
self.counters, self.histograms = {}, {}
def inc(self, name, value=1, labels=None):
key = (name, tuple(sorted((labels or {}).items())))
with self.lock:
self.counters[key] = self.counters.get(key, 0) + value
def observe(self, name, value):
with self.lock:
counts, total = self.histograms.setdefault(name, ([0] * (len(LATENCY_BUCKETS) + 1), [0.0, 0]))
counts[bisect.bisect_left(LATENCY_BUCKETS, value)] += 1
total[0] += value
total[1] += 1
def render(self, gauges):
"""The /metrics page: counters, histograms (cumulative buckets) and current gauges."""
lines = []
def header(name, kind):
lines.append(f"# HELP {self.prefix}{name} {HELP.get(name, name)}")
lines.append(f"# TYPE {self.prefix}{name} {kind}")
with self.lock:
for name in sorted({n for n, _ in self.counters}):
header(name, "counter")
for (n, labels), value in sorted(self.counters.items()):
if n == name:
label = ",".join(f'{k}="{v}"' for k, v in labels)
lines.append(f"{self.prefix}{name}{{{label}}} {value}" if label else f"{self.prefix}{name} {value}")
for name, (counts, (total, count)) in sorted(self.histograms.items()):
header(name, "histogram")
running = 0
for bound, c in zip(list(LATENCY_BUCKETS) + ["+Inf"], counts):
running += c
lines.append(f'{self.prefix}{name}_bucket{{le="{bound}"}} {running}')
lines.append(f"{self.prefix}{name}_sum {total}")
lines.append(f"{self.prefix}{name}_count {count}")
for name, value in sorted(gauges.items()):
header(name, "gauge")
lines.append(f"{self.prefix}{name} {value}")
return "\n".join(lines) + "\n"
The ones to alert on: time_to_first_token_seconds p99 (users waiting), num_requests_waiting (load beyond capacity), kv_cache_usage_perc near 1 together with preemptions (memory-bound; Chapter 31), and request_success_total{finished_reason="abort"} (clients giving up).
Security
The API key check is a minimum. The other protections are already in place from earlier chapters: chat templates render in a sandbox (Chapter 35), prefix-cache keys are cryptographic and can carry a per-tenant salt (Chapter 31), and every request is bounded by max_model_len and the server’s limits. Chapter 44 adds request deadlines, raw-body limits, bounded stream buffers and drain/readiness; per-key rate limiting remains deployment work.
Serving a checkpoint
uv pip install -r serve-requirements.txt
python serve.py --model-dir models/Qwen3-0.6B --attention-backend triton --cuda-graphs auto --reasoning
serve.py --help lists every flag; each maps to a field of Chapter 31’s EngineConfig or to this chapter’s server settings. Any OpenAI client works:
from openai import OpenAI
client = OpenAI(base_url="http://localhost:8000/v1", api_key="unused")
reply = client.chat.completions.create(model="Qwen3-0.6B", messages=[{"role": "user", "content": "Hi!"}])
Run it
python run.py serve --requests 32 --slots 32 --new-tokens 32
Without a checkpoint, the demo trains a 1,000-token tokenizer, builds a random 2-layer model with the same vocabulary, starts the server with the engine core in its own process, and acts as a client:
{"models": [{"id": "izh-tiny", "object": "model", "owned_by": "izh", "max_model_len": 2048}]}
{"finish_reason": "tool_calls", "tool_calls": [{"id": "call_35b136318e8944f082ba71c1", "type": "function", "function": {"name": "get_weather", "arguments": "{\"city\": \" depreerဵs\"}"}}], "usage": {"prompt_tokens": 281, "completion_tokens": 68, "total_tokens": 349, "prompt_tokens_details": {"cached_tokens": 0}}}
{"sse_events": 13, "first": "data: {\"id\": \"cmpl-0c6f8af0647d4499850aa6a2112b09e1\", \"object\": \"text_completion\", \"create...", "last": "data: [DONE]"}
{"concurrent_requests": 32, "output_tokens": 1024, "seconds": 0.67, "tok_s": 1529.8, "ttft_p50_ms": 120.4, "ttft_max_ms": 124.1}
{"metrics": ["izh:generation_tokens_total 1104", "izh:request_success_total{finished_reason=\"length\"} 33", "izh:request_success_total{finished_reason=\"stop\"} 1", "izh:time_to_first_token_seconds_count 34"]}
tool_choice="required" forced a well-formed call to get_weather from a model with random weights (its choice of city is as random as its weights). A streamed completion of 12 tokens arrived as 13 events plus [DONE]. Thirty-two concurrent streaming requests through HTTP, the frontend, the queues and the engine process ran at about 1,500 tokens per second on a laptop CPU, every request receiving its first token within 125 ms of the others: they were admitted in the same step.
Build it
Engine milestone 36: the server. Implement core_loop, AsyncLLM._dispatch and AsyncLLM.generate in engine/serve/async_engine.py, and Server.sampling_params and Server.chat_prompt in engine/serve/api.py (the route handlers, streaming, merging, metrics and the launcher are provided).
pytest tests/test_ch36_server.py
python run.py serve --impl engine
python serve.py --model-dir <a Qwen3 checkpoint> # then point any OpenAI client at it
The tests check that greedy completions equal the offline engine’s output, streamed and not, with usage chunks; batch prompts with n=2 and seeded reproducibility; stop strings and logprobs, streamed and not; a forced tool call in chat, streamed and not; JSON-schema response format with chat logprobs; 401, 400 (too long, invalid temperature) and 404 errors in OpenAI’s format and the metrics page; cancellation freeing the core’s blocks and backpressure refusing excess requests; and the engine core in its own process producing the same answer.
Stretch exercises
- ★ Add request deadlines: an
X-Request-Timeoutheader after which the server aborts the request and returns what it has withfinish_reason: "length". Where: header handling inengine/serve/api.py, with timed cancellation viaAsyncLLM.abortinengine/serve/async_engine.py. - ★★ Rate-limit per API key with a token bucket on tokens, not requests (prompt +
max_tokens), and return 429 with aRetry-Afterheader. Where: API-key admission inengine/serve/api.py; add token-bucket state there or inengine/serve/limits.py. - ★★ Graceful shutdown: on SIGTERM, stop accepting requests (503), finish the running ones, then exit. Test it with a request in flight. Where: application lifespan/shutdown in
build_appinengine/serve/api.py, with draining inengine/serve/async_engine.py. - ★★★ Replace the queues with ZeroMQ sockets and msgpack serialization, and measure the frontend-to-core overhead per step at 256 concurrent streams against
multiprocessing.Queue. Where: queue creation and send/receive paths inengine/serve/async_engine.py.
Check your understanding
- Why does running the engine core in its own process improve the decode latency of every request?
- Which failures should fail one request, and which should fail the whole engine? Why?
- Why are stop strings handled by the frontend and not the engine core?
- Trace what happens, step by step, when a client closes a streaming connection.
- Why must the server refuse an overloaded request before it starts streaming the response?
- Why should a server’s default temperature come from the checkpoint rather than from OpenAI’s API defaults?
Going deeper
- vLLM:
vllm/entrypoints/openai/api_server.pyandserving_chat.py(routes and streaming),vllm/v1/engine/async_llm.py,core_client.py(the ZeroMQ link to the engine process) andoutput_processor.py(detokenization and stop strings in the frontend). - SGLang:
python/sglang/srt/managers/tokenizer_manager.py,scheduler.pyanddetokenizer_manager.py, a three-process design. - OpenAI’s API reference for Completions, Chat Completions and streaming; the WHATWG HTML standard’s section on server-sent events; Prometheus’s Exposition formats documentation.
- Beyer et al., Site Reliability Engineering (O’Reilly, 2016), chapters on handling overload and cascading failures, for backpressure and load shedding.