Skip to content

tl-serve: OpenAI-compatible HTTP + SSE on tokio/hyper

ModuleL10.5 · build · Rust · Pass 7 · 16 to 24 h
You buildrust/crates/tl-engine/src/engine.rs: Engine, the step loop over runner, scheduler, block manager, and per-request samplers · rust/crates/tl-serve/src/server.rs: the v1 server on tokio and hyper (runtime.toml, the engine thread, bounded admission, abort on disconnect, drain on SIGTERM, health, readiness, /metrics) · openai.rs (validation and the response and chunk documents) · sse.rs (SSE framing, incremental UTF-8, stop strings) · template.rs (the chat-template Jinja subset) · the module lines of your L10.0 crate root tl-serve/src/lib.rs · your entry point’s --config form
ContractHTTP: openapi/openai-subset.v1.yaml (engine tier) · the engine role, config form: spec/cli-roles.md · [engine]: config/runtime.schema.json · chat templates: formats/generation-config.schema.json · metric names: otel/metrics.yaml
Testscourse/tests/rust/l10_5.rs, 19 tests (what they check: section 4); conformance ol conform openapi:v1 --target engine from MS-L10, including the official openai Python client
NeedsL10.1 runner and sampler (chapter) · L10.2 scheduler (chapter) · L10.3 chunked prefill (chapter) · L10.4 block manager (chapter) · L10.8 speculative decoding (chapter) · L1.5 tl-tok: the byte tokenizer and tokenizer.json BPE (chapter) · reading: lang.09 async Rust and tokio (primer), lang.05 HTTP and SSE, L10.0 the tracer server, L8.2 generate and the detokenizer · or --ref-deps
Used byL10.6 (disaggregation), L10.7 (metrics and tracing), and L10.9 (tool calls) build on this server, and gw.04, ag.01, load.01, and L12.3 reach it over HTTP
MilestoneMS-L10 (ol conform openapi:v1 --target engine 100%)
Optional depthOpenAI API reference: chat completions and streaming (free); hyper 1.x guide (free); Tokio tutorial: channels, select (free); Jinja template designer docs (free); vLLM’s OpenAI server (free)
  • The model runs on one OS thread and HTTP on tokio tasks; they meet in channels: a bounded admission queue in, one event channel per request out (concurrent_streams_equal_serial).
  • A chat request is rendered through the model’s own chat template, with the generation prompt, before tokenization, and usage.prompt_tokens counts that rendered prompt (hand_example_chat_completion, usage_counts_the_templated_prompt).
  • Streaming is data: <json>\n\n events: the role first, text as it becomes valid UTF-8 and safe from a stop string, finish_reason on the last content chunk, the usage chunk when asked, then [DONE] (stream_framing_role_first_then_done, text_stream_holds_back_stop_prefixes).
  • Every error has the v1 shape, and the status says what kind: 400 bad value, 422 unsupported value, 404 unknown model, 429 full queue with Retry-After (errors_use_the_openai_shape, full_queue_answers_429_with_retry_after).
  • A client that leaves frees its KV blocks within one step, and SIGTERM drains: /readyz turns 503, in-flight streams finish, then exit 0 (disconnect_frees_blocks_within_a_step, drain_finishes_in_flight_and_refuses_new).
Terminal window
ol start L10.5 # stubs engine.rs, server.rs, openai.rs, sse.rs, template.rs
ol tests L10.5
ol check L10.5
ol conform openapi:v1 --target engine # once your entry point serves --config (MS-L10 runs it)

Add to your tl-serve/Cargo.toml what the course manifest lists: tl-engine and tl-tok by path, tokio (full), hyper (server, http1), hyper-util (tokio), http-body-util, bytes, serde_json, toml (contracts/allowed-deps.toml). Add pub mod engine; to tl-engine/src/lib.rs and pub mod openai; pub mod server; pub mod sse; pub mod template; to tl-serve/src/lib.rs. Your http.rs (the tracer) stays as it is: --model-dir --port --health-port keeps serving API v0, and --config <runtime.toml> now starts server::run:

runtime.toml
[engine]
model_dir = "artifacts/models/smol-135m/v3"
http_listen = ":8000"
health_listen = ":9464"
prefix_cache = "radix"
prefill_chunk = 512
Terminal window
rust/target/release/tl-serve --config runtime.toml
curl -N localhost:8000/v1/chat/completions -d '{"model":"smol-135m","messages":[{"role":"user","content":"Hi"}],"stream":true}'

L10.1 to L10.4 built an engine that batches requests, chunks prompts, and shares prefixes, and nothing outside a Rust test can use it: the tracer server (L10.0) still runs one bigram request per thread. Every client in the rest of the course (the gateway, the load generator, the agent SDK, the official OpenAI client) speaks the v1 API: chat messages, streaming with roles, stop sequences, error codes. This module puts the engine behind that API, on tokio and hyper, with the operational behaviors a server needs to live in a cluster: bounded admission, disconnect handling, readiness, and draining.

HTTP connections spend almost all their time waiting; tokio runs thousands of them as tasks on a few threads (lang.09). The model is the opposite: each step is milliseconds of pure CPU, and an .await inside it would block a runtime thread. So the engine runs on its own OS thread and the two sides talk through channels:

handler task --try_send(Generate)--> bounded admission queue --> engine thread
loop:
take commands (add_request)
drop requests whose channel closed
Engine::step (schedule, plan, forward, sample, on_step)
handler / SSE task <--Ev::Token------- one unbounded channel per request --

Engine::step is the whole of continuous batching: schedule (L10.2), plan (L10.3) over the block manager’s tables (L10.4), ModelRunner::forward (L10.1), one sample per sampled row with the request’s own parameters and generator stream(seed, sample), then on_step. A token in stop_ids (EOS) finishes the request with stop.

An unbounded queue turns overload into unbounded latency: every new request waits behind all the others and times out at the client anyway. The admission channel and the scheduler’s waiting queue have capacity queue_capacity; a request that does not fit is answered at once with 429 rate_limit_exceeded and Retry-After: 1, which load balancers and the gateway (gw.04) understand. A request admitted but kept waiting for KV blocks past queue_deadline_ms gets 503 no_capacity.

The model was trained on conversations written in one exact text format. The model directory carries it as a template (generation_config.json’s chat_template, else tokenizer_config.json’s); template.rs renders messages through it with add_generation_prompt, the opening of an assistant turn, so the model continues as the assistant. SmolLM2’s template is

{% for message in messages %}{% if loop.first and messages[0]['role'] != 'system' %}{{ '<|im_start|>system\nYou are a helpful AI assistant named SmolLM, trained by Hugging Face<|im_end|>\n' }}{% endif %}{{'<|im_start|>' + message['role'] + '\n' + message['content'] + '<|im_end|>' + '\n'}}{% endfor %}{% if add_generation_prompt %}{{ '<|im_start|>assistant\n' }}{% endif %}

The interpreter is a small compiler: split the source into text, {{ expression }}, and {% statement %} segments (with - trimming whitespace on that side), parse statements into a tree (for, if / elif / else), parse expressions by recursive descent (or < and < not < ==, !=, is defined < +, ~ < .name, [key], | filter < literals and names), and evaluate over JSON values. Syntax errors surface when the model loads, never on a request. A model with no template (the tracer bigram, the tiny test models) gets a plain default, <|role|>\ncontent\n per message and <|assistant|>\n.

The rendered prompt is tokenized with the model’s tokenizer (tl-tok: bytes, or tokenizer.json BPE with its special tokens matched first), and usage.prompt_tokens counts those tokens: what the model reads and what the gateway’s ledger bills.

SymbolMeaningType
b1,b2,…b_1, b_2, \dotsbytes of the generated tokens, in orderu8
PPbytes received but not yet a complete UTF-8 sequencebytes
TTdecoded text not yet sentstring
Σ\Sigmathe stop strings of the request (at most 4)strings

TextStream::push appends a token’s bytes to PP, moves the longest valid UTF-8 prefix of PP into TT (an invalid sequence becomes U+FFFD; an incomplete tail waits), then: if some σ∈Σ\sigma \in \Sigma occurs in TT, it sends TT up to the first occurrence and stops; otherwise it holds back the longest suffix of TT that is a proper prefix of some σ\sigma and sends the rest. So a stop string never appears in the output, even split across tokens, and the concatenation of everything sent equals the non-streamed text.

An SSE stream is the bytes data: <ChatCompletionChunk json>\n\n per chunk: the first with delta.role: "assistant", then delta.content pieces (possibly empty), the last content chunk with finish_reason (length, or stop for EOS and stop strings), the usage chunk (choices: []) when stream_options.include_usage is set, and data: [DONE]\n\n. After 15 s without a token the server sends the comment : ping\n\n. hyper frames a body of unknown length with chunked transfer encoding; every event is written as soon as it exists.

Every request body is checked against the v1 schema before it reaches the engine: a value out of range (temperature outside 0 to 2, top_p not in (0, 1], max_tokens 0, a stop list of 5) is 400 invalid_request_error with param naming the field; a supported field with a value this engine does not serve (n > 1, echo, content as an array of parts, and until L10.9 tools, tool_choice, response_format) is 422 unsupported_parameter; an unknown model is 404 model_not_found; a prompt plus max_tokens past the context is 400 context_length_exceeded. Unknown fields are ignored. Every error body is {"error": {"message", "type", "param", "code"}} with application/json, and every response carries X-Request-Id (echoed, or generated).

When a client goes away, hyper drops the response body. For a stream, the body is the receiving end of a channel fed by the request’s SSE task; that task’s next send fails, it returns, and it drops the request’s event channel. The engine thread checks is_closed() on every live request before each step and aborts the closed ones: the scheduler releases their blocks within one step. A stop string ends a request the same way: the handler stops reading and drops its channel.

On SIGTERM (Kubernetes’ stop signal) the server closes its API listener, answers /readyz with 503 so no new traffic is routed, answers new requests on open connections with 503, waits (up to 30 s) until every in-flight request has finished, stops the engine thread, and exits 0. /healthz stays 200 throughout: the process is alive.

/v1/embeddings runs each input once (ModelRunner::run_once, between steps on the engine thread), takes the final-normed hidden states h1,…,hT∈Rdh_1, \dots, h_T \in \mathbb{R}^d, and returns e=hˉ/∥hˉ∥2e = \bar h / \lVert \bar h \rVert_2 with hˉ=1T∑tht\bar h = \frac{1}{T} \sum_t h_t. The gateway’s usage-policy classifier (D33) is a linear head over exactly this vector.

The succ bigram (after byte ii, byte i+1i + 1) in a directory succ, with the template {% for m in messages %}{{ m.content }}{% endfor %}.

Request. POST /v1/chat/completions with {"model": "succ", "messages": [{"role": "user", "content": "a"}], "max_tokens": 3, "temperature": 0}.

  1. Validation: model is the served id; temperature 0 is in range; max_tokens 3 plus the prompt fits the context.
  2. Template: one message, so the loop emits its content: the prompt is a.
  3. Tokenize: the byte tokenizer gives [97]; prompt_tokens = 1.
  4. Engine: step 1 prefills position 0 and samples greedily from row 97: b (98); steps 2 and 3 decode c, d. The third token reaches max_tokens: finish length.
  5. Response: {"object": "chat.completion", "choices": [{"index": 0, "message": {"role": "assistant", "content": "bcd"}, "finish_reason": "length"}], "usage": {"prompt_tokens": 1, "completion_tokens": 3, "total_tokens": 4}, ...}. This is hand_example_chat_completion.

With "stream": true, "stream_options": {"include_usage": true} the body is six events: role (content: ""), b, c, d with finish_reason: "length", the usage chunk, [DONE].

A stop string across tokens. With "stop": ["de"] and max_tokens 10, the tokens arrive as b, c, d, e, … After d the text not yet sent is d, a prefix of de: held back. After e it is de: the stop string, so nothing more is sent and the answer is bc with finish_reason: "stop" (stop_strings_end_generation_across_tokens).

The template, by hand (chat_template_hand_example). SmolLM2’s template with [{"role": "user", "content": "Hi"}]: loop.first is true and messages[0]['role'] is user, not system, so the default system turn is emitted first; then <|im_start|>user\nHi<|im_end|>\n; then, because add_generation_prompt is true, <|im_start|>assistant\n.

rust/crates/tl-engine/src/engine.rs
pub struct GenRequest { pub prompt: Vec<u32>, pub params: SamplingParams, pub seed: u64, pub max_tokens: usize, pub stop_ids: Vec<u32>, pub priority: i32 }
pub struct TokenOut { pub id: RequestId, pub token: u32, pub logprob: f64, pub finish: Option<FinishReason> }
impl Engine {
pub fn new(runner: ModelRunner, cfg: &EngineConfig, max_waiting: usize) -> Engine;
pub fn add_request(&mut self, r: GenRequest) -> Result<RequestId, AdmitError>;
pub fn abort(&mut self, id: RequestId); pub fn abort_all(&mut self);
pub fn step(&mut self) -> anyhow::Result<Vec<TokenOut>>;
pub fn has_work(&self) -> bool; pub fn stats(&self) -> EngineStats; pub fn max_model_len(&self) -> usize;
pub fn runner_mut(&mut self) -> &mut ModelRunner;
}
// rust/crates/tl-serve/src/{server,openai,sse,template}.rs
pub struct ServeConfig { pub model_dir: PathBuf, pub model_id: String, pub http_listen: String, pub health_listen: String,
pub engine: EngineConfig, pub queue_capacity: usize, pub queue_deadline_ms: u64,
pub handle_signals: bool, pub step_delay: Duration }
impl ServeConfig { pub fn for_model(dir: &Path) -> ServeConfig;
pub fn from_toml(text: &str, env: &[(String, String)]) -> Result<ServeConfig, String>;
pub fn load(path: &Path) -> Result<ServeConfig, String>; }
pub fn spawn(cfg: ServeConfig) -> Result<ServerHandle, String>; // binds, then serves on its own threads
impl ServerHandle { pub fn drain(&self); pub fn join(self) -> Result<(), String>; } // pub http_addr, health_addr
pub fn run(cfg: ServeConfig) -> Result<(), String>; // the binary: serve until SIGTERM, drain, return
pub fn pool_embedding(hidden: &[f32], d: usize) -> Vec<f32>; pub fn metrics_text(s: &EngineStats) -> String;
pub struct ApiError { pub status: u16, pub kind: &'static str, pub message: String, pub param: Option<String>, pub code: Option<&'static str> }
pub fn parse_chat(body: &[u8]) -> Result<GenSpec, ApiError>; pub fn parse_completion(body: &[u8]) -> Result<GenSpec, ApiError>;
pub fn parse_embeddings(body: &[u8]) -> Result<(String, Vec<String>), ApiError>; pub fn parse_tokenize(body: &[u8]) -> Result<(String, String, bool), ApiError>;
pub fn event(payload: &str) -> String; pub const DONE: &str; pub const PING: &str;
pub struct TextStream { /* pending bytes, held text, stops */ }
impl TextStream { pub fn new(stops: Vec<String>) -> Self; pub fn push(&mut self, bytes: &[u8]) -> String; pub fn finish(&mut self) -> String; pub fn stopped(&self) -> bool; }
pub struct Template { /* parsed nodes */ }
impl Template { pub fn parse(src: &str) -> Result<Template, String>; pub fn render(&self, ctx: &serde_json::Value) -> Result<String, String>;
pub fn render_chat(&self, messages: &Value, add_generation_prompt: bool, bos: &str, eos: &str, tools: Option<&Value>) -> Result<String, String>;
pub fn render_messages(&self, messages: &[(&str, &str)], add_generation_prompt: bool, bos: &str, eos: &str) -> Result<String, String>; }
pub const DEFAULT_TEMPLATE: &str;

Routes: API port POST /v1/chat/completions, POST /v1/completions, POST /v1/embeddings, POST /v1/tokenize, GET /v1/models, GET /v1/models/{id}, GET /healthz; health port GET /healthz, GET /readyz, GET /metrics (tl_engine_kv_blocks{state}, tl_engine_queue_depth, tl_engine_active_sequences, tl_engine_batch_tokens, tl_engine_prefix_cache_hit_ratio, tl_engine_preemptions_total, tl_engine_kv_evictions_total; L10.7 adds the histograms and traces). X-TL-Priority (an integer in -100 to 100) sets the request’s priority; X-TL-KV-Handle answers 404 kv_handle_not_found until L10.6.

Your main.rs: --config <path> loads ServeConfig::load (with TL_ENGINE__<KEY> overrides) and calls server::run; exit 0 after the drain.

TestKINDChecksWhy it matters downstream
hand_example_chat_completionunitsection 3: bcd, length, usage 1 + 3, every field of chat.completionthe worked example over the wire
stream_framing_role_first_then_doneconformancecontent type, one id, role first, finish on the last content chunk, usage chunk, [DONE]every SSE client parses this
stream_equals_nonstream_at_temperature_0differentialdeltas concatenate to the plain content on the tiny Llamacase chat.stream.equals_nonstream
seed_reproduces_samplingpropertythe same seed repeats, another seed differscase chat.seed
stop_strings_end_generation_across_tokensunitde split across tokens stops at bc, plain and streamedcase chat.stop
errors_use_the_openai_shapeconformancethirteen bad requests: status, shape, param, codecases err.400, err.422, route errors
completions_endpoint_and_usageconformance/v1/completions: default 16 tokens, usagethe text endpoint of v1
usage_counts_the_templated_promptconformanceprompt_tokens = tokens of the rendered default templatecase chat.usage, the gateway’s billing
full_queue_answers_429_with_retry_afterfaultqueue of one: the third request gets 429 with Retry-After at once; the queued one runs laterbounded admission
disconnect_frees_blocks_within_a_stepfault/metrics shows 0 active sequences and 0 used blocks within 2 s of a disconnectcase cancel.disconnect
drain_finishes_in_flight_and_refuses_newfaultdrain: /readyz 503, new connections refused, the open stream reaches [DONE]rolling deploys on Kubernetes
embeddings_are_mean_pooled_and_normalizeduniteach vector is the normalized mean of the hidden states, in input orderthe policy classifier’s input (D33)
tokenize_models_health_and_request_idconformance/v1/tokenize, /v1/models, health on both ports, X-Request-Idthe gateway’s TPM cost and routing
concurrent_streams_equal_serialdifferentialeight parallel streams equal the serial answerscase concurrency.16
chat_template_hand_exampleunitSmolLM2’s real template, filters, elif, loop.last, whitespace control, the default templatethe agent model (D20) gets the prompt it was trained on
template_errors_are_caught_at_parseboundaryunclosed blocks, unknown statements and filters refused at parsea broken template fails at load
text_stream_holds_back_stop_prefixesunitstop-prefix hold-back and incremental UTF-8, byte by bytestreams never leak a stop string
runtime_toml_and_env_overridesunit[engine] keys, defaults, TL_ENGINE__ overrides, refused keys and valuesthe config form of spec/cli-roles.md
speculative_runtime_config_drives_the_serving_pathintegration[engine].speculative selects the runner-backed L10.8 path; greedy output matches ordinary servingprompt lookup can be enabled without changing greedy answers
PitfallSymptomCaught by
Reporting a length cut as stopclients think the model finished its answerhand_example_chat_completion (mutant s01)
Rendering without the generation promptthe model continues the user’s turn instead of answeringusage_counts_the_templated_prompt (mutant s02)
Ending an SSE event with one newlineevery client merges all chunks into one eventstream_framing_role_first_then_done (mutant s03)
No role in the first deltathe OpenAI client cannot assemble the messagestream_framing_role_first_then_done (mutant s04)
Sending text that may begin a stop stringthe stream shows d of de before it stopstext_stream_holds_back_stop_prefixes (mutant s05)
Returning the stop string itselfanswers end with the stop sequencestop_strings_end_generation_across_tokens (mutant s06)
Accepting out-of-range valuestemperature: -1 reaches the samplererrors_use_the_openai_shape (mutant s07)
Ignoring na client asking for 2 choices silently gets 1errors_use_the_openai_shape (mutant s08)
A full queue answered with 503 or a waitclients and the gateway cannot tell overload from failure; latency grows unboundedfull_queue_answers_429_with_retry_after (mutant s09)
abort that does not reach the schedulera gone client keeps generating and holding blocksfull_queue_answers_429_with_retry_after (mutant s10)
Exiting on SIGTERM without waitingevery rolling deploy cuts streams mid-answerdrain_finishes_in_flight_and_refuses_new (mutant s11)
Returning the mean without normalizingcosine similarities and the policy head are offembeddings_are_mean_pooled_and_normalized (mutant s12)
loop.first off by onethe default system turn appears in the wrong placechat_template_hand_example (mutant s13)
Ignoring TL_ENGINE__ overridesa ConfigMap cannot be tuned per podruntime_toml_and_env_overrides (mutant s14)
Dropping X-Request-Idone request cannot be followed through gateway and engine logstokenize_models_health_and_request_id (mutant s15)
Running a step inside an async taskone long prefill stalls every connection on that runtime threadconcurrent_streams_equal_serial

| Forward | L12.3 | Registered module relationship. |

DirectionModuleHow it uses this
BackL10.1the runner (also for embeddings) and the sampler with per-request stream(seed, sample)
BackL10.2the scheduler: add, abort, schedule, on_step; QueueFull becomes 429
BackL10.3plan builds every step; prefill_chunk turns chunking on
BackL10.4the block manager in the mode prefix_cache names; hit ratio and block counts in /metrics
BackL10.8runner-backed speculative generation when [engine].speculative is configured
BackL1.5tl-tok: byte tokenizer and tokenizer.json BPE for prompts and token bytes
ForwardL10.6disaggregated prefill and decode add EngineControl and the X-TL-KV-Handle resume to this server
ForwardL10.7histograms, OpenTelemetry spans, and trace propagation join /metrics
ForwardL10.9tools, tool_choice, and response_format replace today’s 422
Forwardgw.04, load.01, ag.01the gateway proxies this API; the load generator and the agent SDK are its clients

If you skip this module, MS-L10’s conformance run has no v1 engine to test.

Your pieceProduction equivalentWhat it addsWhere to look
one engine thread and channelsvLLM’s AsyncLLM with a separate engine core processthe engine in another process (ZeroMQ), so Python’s GIL never blocks itvllm/v1/engine/async_llm.py
a Jinja subsetminijinja, Hugging Face apply_chat_templatethe full template language: macros, raise_exception, namespacestransformers apply_chat_template
bounded queue + 429admission control with token buckets and fairnessper-tenant limits, priority classesgw.03, gw.04
hand-written validationOpenAPI-generated servers (utoipa, progenitor)schema and code from one sourcecraft.14’s v2 migration
/metrics gaugesOpenTelemetry metrics and traceshistograms of TTFT and TPOT, spans per requestL10.7