The idea in one minute#
The engine core returns integers: for each request that made progress this step, a short list
of new token IDs. Turning those into the text a client sees is harder than it sounds. One
token is not one character — an emoji or a Chinese character can be split across several
tokens — so text can only be released when it is complete. Stop strings must be found even
when they straddle two tokens, and text that might be the start of a stop string must be
held back. And all of it has to be done for hundreds of requests on one event loop without
stalling any of them. The class that connects an HTTP handler to this machinery is
AsyncLLM; the class that does the conversion is OutputProcessor.
A picture#
flowchart LR
E[":vllm: <b>Engine core</b><br/><small>EngineCoreOutputs per step</small>"] -->|"PULL socket"| OT[":i-activity: <b>Socket task</b><br/><small>decode msgpack</small>"]
OT --> OQ[("outputs_queue")]
OQ --> OH[":i-recycle: <b>output_handler</b><br/><small>one task for all requests</small>"]
OH --> OP[":i-message-square: <b>OutputProcessor</b><br/><small>detokenise, stop, logprobs</small>"]
OP --> C1[("collector A")]
OP --> C2[("collector B")]
OP --> C3[("collector C")]
C1 --> G1[":i-globe: <b>generate() A</b><br/><small>SSE to client A</small>"]
C2 --> G2[":i-globe: <b>generate() B</b>"]
C3 --> G3[":i-globe: <b>generate() C</b>"]
OP -.->|"stop string hit: abort"| E
class E io
class OT,OH queue
class OQ,C1,C2,C3 memory
class OP compute
class G1,G2,G3 ioHow it really works#
One reader, many writers#
Every HTTP request is its own asyncio task running AsyncLLM.generate(). There could be a
thousand of them. They do not each read from the engine socket. Exactly one background
task, the output handler, reads the engine’s output and distributes it:
async def output_handler():
while True:
# 1) Pull EngineCoreOutputs from the EngineCore.
outputs = await engine_core.get_output_async()
...
for start in range(0, num_outputs, chunk_size):
outputs_slice = engine_core_outputs[start:end]
# 2) Process EngineCoreOutputs.
processed_outputs = output_processor.process_outputs(
outputs_slice, outputs.timestamp, iteration_stats
)
# Allow other asyncio tasks to run between chunks
if end < num_outputs:
await asyncio.sleep(0)
# 3) Abort any reqs that finished due to stop strings.
if processed_outputs.reqs_to_abort:
await engine_core.abort_requests_async(processed_outputs.reqs_to_abort)
# 4) Logging.
logger_ref[0].record(...)Three design points are visible in those few lines.
Chunks of 128. One engine step can return outputs for hundreds of requests. Detokenising
all of them without a pause would block the event loop — no new connection accepted, no byte
written to any client — for the whole batch. So the handler processes
VLLM_V1_OUTPUT_PROC_CHUNK_SIZE (default 128) outputs, yields with asyncio.sleep(0) so
other tasks can run, and continues.
One loop over the batch. The docstring of process_outputs carries a rule for
contributors: this “is the only function that should loop over EngineCoreOutputs. If you need
to touch every element of the batch, do it from within the loop below.” At 500 requests a
step, a second Python loop over the batch is measurable.
Metrics are recorded here, in the API server, from timestamps the engine attached. The engine core never touches Prometheus.
generate(): the consumer side#
q = await self.add_request(request_id, prompt, sampling_params, ...)
finished = False
while not finished:
# Note: drain queue without await if possible (avoids
# task switching under load which helps performance).
out = q.get_nowait() or await q.get()
finished = out.finished
if out is not STREAM_FINISHED:
yield outq is a RequestOutputCollector, a one-slot mailbox rather than a queue. If the output
handler deposits a second output before the client task has taken the first, the two are
merged: in streaming (“delta”) mode the new token IDs and text are appended to the waiting
output. A slow client therefore gets larger, less frequent chunks instead of an ever-growing
backlog, and the server’s memory use does not depend on how slowly clients read.
get_nowait() or await get() is a small optimisation with a large effect under load: if an
output is already waiting, the task takes it without suspending, which avoids a round trip
through the event loop’s scheduler.
Cancellation is automatic#
When a client disconnects, the web server cancels that request’s task. generate() catches
the cancellation and sends an abort to the engine:
except (asyncio.CancelledError, GeneratorExit):
if q is not None:
await self.abort(q.request_id, internal=True)
raiseThe abort travels as an ABORT message, lands on the engine core’s fast aborts_queue, and
the scheduler frees that request’s KV blocks at the end of the step in progress. Closing the
connection is the correct way for a client to stop generation; no separate cancel call is
needed.
Incremental detokenisation#
tokenizer.decode(all_ids) on every step would cost time proportional to the answer’s length,
on every token: quadratic overall. It would also produce wrong output mid-stream. Some
characters span several tokens, and decoding only the first of them yields the Unicode
replacement character �.
So each request owns an incremental detokenizer that holds state between steps and emits only text that is known to be complete:
| Class | Used when | How |
|---|---|---|
FastIncrementalDetokenizer | The tokenizer is a Hugging Face “fast” (Rust) tokenizer — almost always | The tokenizers library’s DecodeStream: feed one token, get back finished text or nothing |
SlowIncrementalDetokenizer | Python-only tokenizers | Re-decodes a short trailing window of tokens each step and emits the difference |
“Nothing” is a normal result. When a token ends in the middle of a character, the stream
returns no text and the piece appears with the next token. Clients see this as an occasional
event with an empty content.
Stop strings, and why text is held back#
A client can pass stop: ["\n\nUser:"]. There are two different kinds of stop:
- Stop token IDs (including end-of-sequence) are checked in the engine core, right after sampling. They cost nothing extra.
- Stop strings can only be checked on text, and the engine core has no tokenizer. They
are checked in the API server, inside
detokenizer.update().
That split has two consequences.
The engine finds out late. When the API server finds a stop string, the engine core is
still generating. So process_outputs adds the request to reqs_to_abort, and the output
handler sends an abort. The engine may compute one or two tokens past the stop before the
abort arrives; they are discarded. This is the dotted arrow in the picture.
Text must be withheld. Suppose the stop string is "\n\nUser:" and the text so far ends
with "\n\nUs". That might be the start of the stop string, or the start of “Usually”. If it
were streamed and turned out to be the stop string, the client would already have received
text it should never have seen. The detokenizer therefore holds back the last few characters:
# Number of chars to hold back when stop strings are to be excluded
# from streamed output.
if self.stop and not self.include_stop_str_in_output:
self.stop_buffer_length = max(len(s) for s in self.stop) - 1
else:
self.stop_buffer_length = 0With stop strings set, the stream always trails generation by longest_stop - 1 characters,
and the held-back text is flushed when the request finishes. A long stop string therefore
makes streaming look slightly laggy; that is the cost of correctness.
min_tokens interacts here too: while fewer than min_tokens tokens have been produced, the
position from which stop strings are searched keeps moving forward, so a stop string cannot
fire early.
Three shapes of output#
SamplingParams.output_kind decides what make_request_output returns each step:
| Kind | Each output contains | Used by |
|---|---|---|
DELTA | Only the new text and tokens since the last output | Streaming HTTP responses |
CUMULATIVE | Everything so far | Legacy callers |
FINAL_ONLY | Nothing until the request finishes, then everything | Non-streaming HTTP responses, offline LLM |
FINAL_ONLY matters for throughput: a non-streaming request creates no intermediate objects,
so a server handling only non-streaming traffic does far less Python work per token.
Trading smoothness for throughput: --stream-interval#
By default every token becomes one output and one server-sent event. For small models that generate thousands of tokens per second across all requests, the per-event cost — building objects, serialising JSON, a system call per write — becomes a real fraction of total CPU.
--stream-interval N (default 1) releases output only when the request has finished, or this
is its first token, or N tokens have accumulated since the last release:
if not (
finished
or self.sent_tokens_offset == 0
or self.detokenizer.num_output_tokens() - self.sent_tokens_offset
>= self.stream_interval
):
return NoneThe first token is always sent immediately, so time to first token is unaffected. Only the cadence afterwards changes.
When the engine dies#
If the socket task receives the ENGINE_CORE_DEAD marker or raises, the exception is put on
the outputs queue, the output handler catches it, and output_processor.propagate_error(e)
delivers it to every waiting collector. All in-flight requests fail at once with
EngineDeadError, and the server’s health check starts failing. No request is left hanging.
Code#
An incremental detokenizer with a stop-string hold-back buffer. The “tokens” are short strings so the logic is visible; the algorithm is the one above.
package main
import (
"fmt"
"strings"
)
type detok struct {
stops []string
holdBack int // longest stop string minus one
text string // everything decoded so far
sent int // how much of text has been released to the client
stoppedBy string
}
func newDetok(stops ...string) *detok {
d := &detok{stops: stops}
for _, s := range stops {
if len(s)-1 > d.holdBack {
d.holdBack = len(s) - 1
}
}
return d
}
// update appends one decoded piece and returns the text that is safe to stream.
func (d *detok) update(piece string, finished bool) string {
searchFrom := len(d.text) - d.holdBack // a stop string may straddle the boundary
if searchFrom < 0 {
searchFrom = 0
}
d.text += piece
for _, s := range d.stops {
if i := strings.Index(d.text[searchFrom:], s); i >= 0 {
d.text = d.text[:searchFrom+i] // truncate at the stop string
d.stoppedBy = s
finished = true
}
}
limit := len(d.text)
if !finished {
limit -= d.holdBack // might still turn out to be the start of a stop string
}
if limit <= d.sent {
return ""
}
out := d.text[d.sent:limit]
d.sent = limit
return out
}
func main() {
pieces := []string{"The", " answer", " is", " 42", ".", "\n", "\n", "User", ":", " thanks"}
d := newDetok("\n\nUser:")
fmt.Printf("hold back %d characters\n\n", d.holdBack)
for i, p := range pieces {
out := d.update(p, false)
fmt.Printf("token %-2d %-10q -> stream %q\n", i+1, p, out)
if d.stoppedBy != "" {
fmt.Printf("\nstop string %q found: abort sent to the engine core\n", d.stoppedBy)
fmt.Printf("tokens generated after this one are discarded\n")
break
}
}
fmt.Printf("\nclient received: %q\n", d.text[:d.sent])
}Notice that the stream stays six characters behind generation the whole way, that the
stop string is found on the token that completes it (":"), and that the client never sees
any part of it.
Remember this#
- One output-handler task reads the engine and feeds every request’s collector.
- It works in chunks of 128 outputs and yields between them so the event loop stays responsive.
- A collector merges outputs for a slow client instead of queueing them.
- A client disconnect cancels the task, which aborts the request in the engine and frees its blocks.
- Detokenisation is incremental and may emit nothing for a token.
- Stop token IDs are checked in the engine core; stop strings in the API server, followed by an abort.
- With stop strings set, streamed text trails by
longest_stop - 1characters. --stream-intervalreduces per-token overhead without changing time to first token.
Try it#
- In the program, add a second stop string
".". What is the hold-back now, and when does generation stop? Which stop string determines the hold-back, the first to match or the longest? - Change
mainto callupdate(p, true)for the final piece of a run that never hits a stop string. Confirm that the held-back characters are flushed. - Against a real server, send the same streaming request with and without
"stop": ["zzzzzzzzzzzzzzzzzzzz"](a 20-character string that will never occur). Compare how the chunks arrive.
Check yourself#
- Why does
process_outputsexist as a single loop, and what would a second loop cost? - Why can a stop string not be detected in the engine core?
- What happens to tokens the engine generates between a stop string appearing and the abort arriving?
Sources#
Checked on 5 October 2026 against main at commit 0c16eee.
vllm/v1/engine/async_llm.py—generate,_run_output_handlervllm/v1/engine/output_processor.py—process_outputs,make_request_output,RequestOutputCollectorvllm/v1/engine/detokenizer.pyvllm/config/scheduler.py—stream_interval