Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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 as stop_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:

requestconstraint
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

  1. ★ Add request deadlines: an X-Request-Timeout header after which the server aborts the request and returns what it has with finish_reason: "length". Where: header handling in engine/serve/api.py, with timed cancellation via AsyncLLM.abort in engine/serve/async_engine.py.
  2. ★★ Rate-limit per API key with a token bucket on tokens, not requests (prompt + max_tokens), and return 429 with a Retry-After header. Where: API-key admission in engine/serve/api.py; add token-bucket state there or in engine/serve/limits.py.
  3. ★★ 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_app in engine/serve/api.py, with draining in engine/serve/async_engine.py.
  4. ★★★ 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 in engine/serve/async_engine.py.

Check your understanding

  1. Why does running the engine core in its own process improve the decode latency of every request?
  2. Which failures should fail one request, and which should fail the whole engine? Why?
  3. Why are stop strings handled by the frontend and not the engine core?
  4. Trace what happens, step by step, when a client closes a streaming connection.
  5. Why must the server refuse an overloaded request before it starts streaming the response?
  6. 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.py and serving_chat.py (routes and streaming), vllm/v1/engine/async_llm.py, core_client.py (the ZeroMQ link to the engine process) and output_processor.py (detokenization and stop strings in the frontend).
  • SGLang: python/sglang/srt/managers/tokenizer_manager.py, scheduler.py and detokenizer_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.