Pidoku

AsyncLLM and the Way Back: Tokens to Text

Basic 50 min Difficulty 3/5 Lesson 03 of 04

Prerequisites The Frontend

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 io

How 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:

Python
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#

Python
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 out

q 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:

Python
except (asyncio.CancelledError, GeneratorExit):
    if q is not None:
        await self.abort(q.request_id, internal=True)
    raise

The 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:

ClassUsed whenHow
FastIncrementalDetokenizerThe tokenizer is a Hugging Face “fast” (Rust) tokenizer — almost alwaysThe tokenizers library’s DecodeStream: feed one token, get back finished text or nothing
SlowIncrementalDetokenizerPython-only tokenizersRe-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:

Python
# 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 = 0

With 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:

KindEach output containsUsed by
DELTAOnly the new text and tokens since the last outputStreaming HTTP responses
CUMULATIVEEverything so farLegacy callers
FINAL_ONLYNothing until the request finishes, then everythingNon-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:

Python
if not (
    finished
    or self.sent_tokens_offset == 0
    or self.detokenizer.num_output_tokens() - self.sent_tokens_offset
    >= self.stream_interval
):
    return None

The 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.

Go
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 - 1 characters.
  • --stream-interval reduces per-token overhead without changing time to first token.

Try it#

  1. 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?
  2. Change main to call update(p, true) for the final piece of a run that never hits a stop string. Confirm that the held-back characters are flushed.
  3. 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#

  1. Why does process_outputs exist as a single loop, and what would a second loop cost?
  2. Why can a stop string not be detected in the engine core?
  3. 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.

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom