Wrap your engine in an OpenAI-compatible streaming API that behaves well when things go wrong.
1. What you build#
An HTTP server exposing /v1/completions and /v1/chat/completions with server-sent-event
streaming, a bounded request queue, cancellation on client disconnect, health endpoints, and
Prometheus metrics. One request runs at a time — batching arrives in Project 07.
From this project on, switch the engine to PyTorch (transformers model, your own generation
loop with past_key_values) so later projects can use a GPU. Keep the loop yours; do not call
model.generate().
Diagram — Architecture#
flowchart LR C["Clients"] -->|"HTTP"| EL["Async event loop<br/>request handlers"] EL -->|"enqueue, or 429 if full"| Q["Bounded queue"] Q --> ET["Engine thread<br/>owns the model"] ET -->|"token pieces"| OQ["Per-request output queue"] OQ --> EL EL -->|"SSE"| C EL -.->|"disconnect sets cancel flag"| ET ET --> M["Metrics endpoint"] class EL,OQ io class Q queue class ET compute class C,M neutral
2. Why it matters#
A model is not a service. The gap between them — queues, timeouts, streaming, backpressure, cancellation, metrics — is where most production incidents live, and none of it is ML.
3. Read first#
- V.14 — Streaming
- VIII.01 — Anatomy of a serving system
- VIII.02 — APIs
- VIII.03 — Queues and admission
- VIII.05 — Backpressure, timeouts, retries
4. Spec#
POST /v1/completions {model, prompt, max_tokens, temperature, top_p, stop, seed, stream}
POST /v1/chat/completions {model, messages, ...same}
GET /v1/models
GET /health/live process is up
GET /health/ready model is loaded AND queue is below limit
GET /metrics Prometheus text format
Behaviour
- stream=true → SSE chunks, terminated by `data: [DONE]`
- queue full → 429 with Retry-After, immediately
- client disconnect → generation stops within one token
- per-request deadline → 504 (or truncated stream with finish_reason)
- usage{prompt_tokens, completion_tokens} in every responseArchitecture: an asyncio HTTP layer, a bounded queue, and one dedicated engine thread that owns the model. The event loop must never run a forward pass.
5. Milestones#
- Non-streaming endpoint that works with the official
openaiPython client pointed at your base URL. - Streaming. Token-by-token SSE. Handle the detokenization problem: a token may be half of a UTF-8 character or a word piece — emit only complete text (V.14).
- Engine thread + queue.
asyncio.Queuein, per-requestasyncio.Queueout, bridged withloop.call_soon_threadsafe. - Admission. Bounded queue; 429 when full. Reject prompts longer than the context window with a 400 before they queue.
- Cancellation. Detect disconnect; set a flag the decode loop checks every token.
- Metrics. TTFT, inter-token latency, queue wait, end-to-end latency (histograms); queue depth, in-flight (gauges); tokens in/out, requests by status (counters).
- Load test. Write an open-loop Poisson load generator — you reuse it in 07, 08, 12.
6. Starter skeleton#
// Job is one request travelling from the HTTP handler to the engine and back.
type Job struct {
PromptIDs []int
Params Params
Out chan string // pieces of text; closed by the engine when done
Ctx context.Context // cancelled when the client disconnects
}
var workQ = make(chan *Job, 32) // bounded: this IS the admission control
// engineLoop owns the model. Exactly one goroutine ever touches it.
func engineLoop(m Model) {
for job := range workQ {
generateStream(job.Ctx, m, job.PromptIDs, job.Params, job.Out) // stops early if Ctx is cancelled
close(job.Out) // sentinel: no more tokens
}
}
func completions(w http.ResponseWriter, r *http.Request) {
job, err := jobFromRequest(r) // sets job.Ctx = r.Context()
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
select {
case workQ <- job:
default: // queue full: refuse immediately instead of piling up
w.Header().Set("Retry-After", "1")
http.Error(w, "server overloaded", http.StatusTooManyRequests)
return
}
w.Header().Set("Content-Type", "text/event-stream")
flusher := w.(http.Flusher)
for {
select {
case <-r.Context().Done(): // client disconnected; the engine sees the same context
return
case piece, ok := <-job.Out:
if !ok {
fmt.Fprint(w, "data: [DONE]\n\n")
flusher.Flush()
return
}
fmt.Fprintf(w, "data: %s\n\n", chunkJSON(job, piece))
flusher.Flush()
}
}
}
func main() {
go engineLoop(loadModel())
http.HandleFunc("POST /v1/completions", completions)
log.Fatal(http.ListenAndServe(":8000", nil))
}7. What to measure#
| Measurement | Expectation to write down first |
|---|---|
| TTFT and ITL at 1 concurrent client | Your baseline |
| TTFT p50/p95 as arrival rate rises to and past capacity | Hockey stick — at what utilization? |
| Throughput (tok/s) vs offered load | Flat: one request at a time |
| Server overhead: end-to-end minus pure engine time | Should be < 1-2 ms/token |
| Time from disconnect to engine idle | ≤ one token |
| Closed-loop vs open-loop load at the same “concurrency” | Closed-loop hides queueing |
The hockey stick you see here is Little’s Law and M/M/1 in practice. Keep the plot; Project 07 and 08 are measured against it.
8. Done when#
- The
openaiclient works against your server, streaming and non-streaming. - Overload produces fast 429s, not growing latency.
- Killing a client mid-stream stops generation; you verified it in the metrics.
-
/metricsexposes TTFT and ITL histograms, and you have a Grafana panel or a plot. -
readygoes false while loading and when saturated;livestays true.
9. Common pitfalls#
Running the model on the event loop. Every other connection freezes, including health checks — and the orchestrator kills your pod under load.
Unbounded queue. Latency grows without limit and every request eventually times out anyway.
Emitting raw token pieces. Broken multi-byte characters in the stream.
Buffering proxies. Nginx and friends buffer SSE unless told otherwise; TTFT looks terrible for reasons that are not yours.
Averaging latencies. Report percentiles from histograms.
Stop strings spanning tokens. You must hold back text that could be the start of a stop sequence.
10. Stretch goals#
- Add
logprobsandn > 1. - Structured output: constrain decoding to a JSON schema with a logit mask.
- Graceful shutdown: on SIGTERM stop admitting, drain in-flight, then exit (VIII.08).
- OpenTelemetry spans per request with GenAI semantic conventions (XI.08).
11. Interview questions this project answers#
- Why must the forward pass not run on the async event loop?
- Why bound the queue, and what should the server return when it is full?
- What is the difference between liveness and readiness for a model server?
- How do you stop wasting GPU on a client that went away?
- Why does an open-loop load test tell you more than a closed-loop one?