Skip to content

SSE streaming proxy (upgrades gw.00)

Modulegw.04 · build · Go · Pass 7 · 4 h
You buildgo/gateway/proxy/proxy.go, which you wrote in gw.00 and now extend: StreamProxy, StreamObserver, ExchangeObserver, Observers, Forward, Outbound, NewOutbound, UpstreamError, and Handler (the chain’s proxy stage), keeping NewProxy and everything gw.00 proved
ContractStreaming, Disconnect, and Errors rules of course/contracts/openapi/openai-subset.v1.yaml (gateway tier); NewProxy still serves course/contracts/openapi/openai-subset.v0.yaml
Testscourse/tests/go/gw_04/ (what they check: section 4), with the gw.00 suite as smoke regression · the bench BenchmarkFirstByteOverhead (ol bench gw.04 --assert)
Needsgw.01 server skeleton (server.Usage, the Exchange) · reading: gw.00 tracer gateway, L10.5 the engine’s v1 streams, SSE in ai-platform-engineering/03
Used bygw.05 forwards every routing attempt through Forward and relies on its commitment boundary · inherited from gw.00: dep.00 builds it into the gateway image, obs.00 traces it
MilestoneMS-gateway
Optional depthWHATWG HTML, server-sent events (free); RFC 9110 section 9.3 on request semantics and retries (free)
  • The proxy still forwards bytes: every read from the engine is written and flushed before the next read, so the client’s stream is byte-identical to the engine’s whatever the read sizes (TestRecordedStreamsByteExact, TestStreamsEventByEvent).
  • While the bytes pass, a splitter cuts them into SSE events at blank lines for an observer: the first byte (time to first token), every event, and the usage from the last usage chunk, delivered once at a clean end (TestHandExample).
  • Forward draws the commitment boundary: before the engine answers, a failure writes nothing and is reported uncommitted, so routing may try another worker; after the first byte nothing is retried and a broken stream ends with an SSE error event (TestForwardUncommittedOnRefused, TestChaosResetMidStream).
  • The client’s context is the upstream request’s context: when the client leaves, the engine’s request is cancelled (TestClientDisconnectCancelsUpstream).
  • Handler is the chain’s proxy stage: no static key (authn already ran), the gateway’s X-Request-Id and no duplicate from the engine, the first byte and usage recorded on the Exchange for the limiter and the meter (TestHandlerInChain).
Terminal window
ol start gw.04 # proxy.go exists (gw.00): prints the contract diff instead of overwriting
ol tests gw.04 # read the test catalog first
ol check gw.04 # gw.00's smoke tests, then gw.04's tests
ol bench gw.04 --assert # local only: first-byte overhead within 5 ms
ol diff gw.04 # after passing: your code against the reference

gw.04 takes over a file you own (upgrades): edit your own proxy.go in place. Your gw.00 verdict is frozen as passed, and its tests now run as gw.04’s smoke regression, so the tracer behavior must survive the upgrade. In your composition root, the Proxy slot becomes proxy.Handler(proxy.Config{Upstream: ...}).


The gw.00 proxy copies bytes and knows nothing about them. Pass 7 needs to know three things about every stream as it passes: when its first byte reached the client (time to first token, the SLO of obs.03), how many tokens it used (the limiter settles with it in gw.03, the ledger records it in gw.07), and whether it may still be retried elsewhere (failover in gw.05). Parsing and re-emitting events would be the easy way to learn those things and the wrong one: it delays tokens, it can change bytes (whitespace, key order, a character split across chunks), and it couples the gateway to every field of the engine’s JSON. This module watches the stream without touching it.

An engine’s streamed answer (openai-subset.v1.yaml) is Content-Type: text/event-stream and a sequence of events, each data: <json>\n\n: the first chunk carries delta.role, content chunks follow (some with empty content while a multi-byte character is incomplete), the last content chunk carries finish_reason, an optional chunk with choices: [] carries usage, and the stream ends with data: [DONE]\n\n. Lines starting with : are comments, such as : ping\n\n every 15 s while idle. A non-streamed answer is one JSON body with a usage object.

The proxy relays to the client side and reads from the upstream side. The bytes it reads arrive in arbitrary pieces: TCP and the engine’s buffers decide the boundaries, not the event structure.

StreamProxy(w, r, up, o) does, in order:

  1. Copy the end-to-end headers of up to w (hop-by-hop headers dropped, as in gw.00), keeping the gateway’s own X-Request-Id: Header.Add would append the engine’s copy as a second value, and the client would see two ids.
  2. Write the status.
  3. Loop: read up to 32 KiB; if it is the first body byte, call o.OnFirstByte() before writing it (so the timestamp is not inflated by the write); write the bytes; flush; then feed the same bytes to the splitter.

The splitter appends each piece to a buffer and cuts complete events at \n\n. For each event it joins the payloads of its data: lines (one space after the colon removed), skips events with no data: line (comments), and calls o.OnChunk(data). Because it keeps the remainder between reads, an event split across ten 1-byte reads is still delivered once and whole. If an event’s JSON has a non-null usage, it becomes the stream’s usage (the last one wins). For a JSON body the proxy keeps a copy (up to 4 MiB) and reads usage from it at the end. When the engine omits total_tokens, it is prompt_tokens + completion_tokens.

  1. At a clean end (EOF), call o.OnDone(usage) once.

ExchangeObserver(ex) is the observer the chain uses by default: it records MarkFirstByte and SetUsage on the request’s Exchange, which is how gw.03 settles and gw.07 meters without importing this package. Observers{a, b} fans out to several (gw.06 adds a recorder for stream replay).

SymbolMeaning
tet_ewhen the engine writes its first byte
tct_cwhen the client receives it through the gateway
Δ=tc−te\Delta = t_c - t_ethe gateway’s added time to first byte

With a flush per read, Δ\Delta is one socket read, one write, and one flush: well under a millisecond on one host. Without the flush (section 5 of gw.00), Δ\Delta is the time for 4 KiB of events to accumulate, about a second at 20 tokens per second. The contract budget is Δ≤5\Delta \le 5 ms; BenchmarkFirstByteOverhead measures it as the difference between direct and proxied time to first byte, and ol bench gw.04 --assert checks it. A latency bound is a bench, never part of ol check, because shared CI machines make timing noisy (DESIGN 4.0).

HTTP sends the status line before the body, so a proxy’s options depend on what it has sent:

WhenWhat happenedForward returnsClient sees
before any answerrefused, reset, DNS, timeoutUpstreamError{Committed: false}, nothing writtenwhatever the caller decides: Handler answers 503 no_capacity, gw.05 tries another worker
the engine answered 503 and Retry503 is setdraining or queue full: no work was doneUpstreamError{Committed: false, StatusCode: 503}, nothing writtenthe same
any other answer200, 400, 404 …relayed; status and body pass throughthe engine’s answer
after the first bytethe stream brokeUpstreamError{Committed: true} after one data: {"error": ...} eventthe bytes so far, then an error event, never [DONE]
any timethe client leftthe upstream request is cancelled (it runs on the client’s context)nothing

Retrying after the first byte would repeat or splice tokens: the client already printed “The cat”, and a second engine starts again from “The”. So a committed failure is final. Outbound is the request as it will be sent (method, path, query, the gw.00 header rules, the body as bytes), built once by NewOutbound so gw.05 can send the same request to several workers in turn.

The engine answers with five events, 163, 147, 148, 156, and 14 bytes long (628 in total), and the gateway happens to read 40 bytes at a time (16 reads):

data: {"id":"c1",...,"delta":{"role":"assistant","content":""},...}\n\n
data: {"id":"c1",...,"delta":{"content":"Hel"},...}\n\n
data: {"id":"c1",...,"delta":{"content":"lo"},"finish_reason":"stop"}]}\n\n
data: {"id":"c1",...,"choices":[],"usage":{"prompt_tokens":5,"completion_tokens":2,"total_tokens":7}}\n\n
data: [DONE]\n\n
ReadBytesWhat happens
10 to 39OnFirstByte() (the client has 0 bytes), write, flush; the splitter holds 40 bytes, no \n\n yet
5160 to 199write, flush; the splitter finds the first \n\n (ending at byte 162): OnChunk with the role chunk; 37 bytes stay buffered
8280 to 319the "Hel" event completes (ends at byte 309): OnChunk
12440 to 479the "lo" event completes (ends at byte 457): OnChunk
16600 to 627the usage event (ends at 613) and [DONE] complete: two OnChunk calls; the usage event sets usage 5/2/7
EOFOnDone({5, 2, 7})

The client received exactly the 628 bytes, OnFirstByte was called once, OnChunk five times (the last with [DONE]), and OnDone once. This is TestHandExample.

package proxy // import "tinyllm/gateway/proxy"
// gw.00, unchanged
func NewProxy(cfg Config) http.Handler
func ParseTraceparent(h string) (traceID, parentID, flags string, ok bool)
// gw.04
type Config struct { Upstream, APIKey string; Client *http.Client; OnSpan func(Span)
Observer func(r *http.Request) StreamObserver } // new: per-request observer
func Handler(cfg Config) http.Handler // the chain's proxy stage: no key check
type StreamObserver interface { OnFirstByte(); OnChunk(data []byte); OnDone(u server.Usage) }
func ExchangeObserver(ex *server.Exchange) StreamObserver
type Observers []StreamObserver
func StreamProxy(w http.ResponseWriter, r *http.Request, up *http.Response, o StreamObserver) error
type Outbound struct { Method, Path, RawQuery string; Header http.Header; Body []byte; Retry503 bool }
func NewOutbound(r *http.Request, body []byte) (*Outbound, Span)
type UpstreamError struct { Committed bool; StatusCode int; Err error }
func Forward(ctx context.Context, w http.ResponseWriter, base string, out *Outbound, client *http.Client, o StreamObserver) error

Use only the standard library and tinyllm/gateway/server.

TestKINDChecksWhy it matters downstream
TestHandExampleunitsection 3: bytes, one OnFirstByte before any write, five events, usage 5/2/7, status and content typeyou and the tests agree on what is observed
TestRecordedStreamsByteExactconformancefive fixture streams (chat, UTF-8 split, ping without usage, tool-call fragments, text completions) at read sizes 1, 7, 64, 65536: bytes equal, event counts, usageconformance chat.stream.framing through the gateway
TestPingCommentsPassThroughunit: ping reaches the client and is not an eventidle streams stay open without fake events
TestNonStreamUsageunita JSON body’s usage, total_tokens computed, no eventsthe limiter and the ledger for non-streamed calls
TestForwardUncommittedOnRefusedfaulta dead engine: uncommitted error, nothing writtengw.05 can fail over
TestForwardRetry503unitwith Retry503 a 503 is handed back unwritten; without it, relayeddraining engines are skipped, not shown to clients
TestNewOutboundDropsKeyunitno Authorization or hop-by-hop headers; request id, path, query, trace keptthe engine tier stays key-free
TestStreamsEventByEventunitthrough real sockets, event 1 arrives while the engine holds event 2tokens stream
TestHandlerInChainunitin the chain: no key needed, one X-Request-Id (the gateway’s), first byte, usage, and status on the Exchangegw.03 settles and gw.07 meters from it
TestHandlerUpstreamDown503boundarya dead engine before the first byte: 503 no_capacitythe gw.00 answer, kept
TestChaosResetMidStreamfaultthe testkit chaos proxy resets the engine connection after the first event: the client gets an SSE error event and no [DONE]; no OnDoneno silent truncation; no metering of a broken stream
TestClientDisconnectCancelsUpstreamfaulta client that leaves cancels the engine’s request (contract: 100 ms; checked here within 1 s)the engine stops generating for nobody
PitfallSymptomCaught by
1. no flush after each writetokens arrive in bursts; with a fake engine that waits, nothing arrivesTestStreamsEventByEvent (mutant s01)
2. OnFirstByte after the write, or on every readTTFT includes the write; the metric fires dozens of times per streamTestHandExample (mutant s02)
3. usage from the wrong place: never from the SSE usage chunk, never from a JSON body, total_tokens left 0the limiter refunds everything; the ledger records zero tokensTestHandExample, TestNonStreamUsage (mutants s03, s11, s13)
4. cutting events per read, or reporting comments as eventsevents split across reads vanish; pings count as tokensTestRecordedStreamsByteExact, TestPingCommentsPassThrough (mutants s04, s05)
5. a broken stream that ends quietly and reports donethe client shows a short answer as complete; the ledger meters itTestChaosResetMidStream (mutant s08)
6. Forward that answers the client itself on failure, or ignores Retry503failover is impossible: the 503 is already sentTestForwardUncommittedOnRefused, TestForwardRetry503 (mutants s06, s07)
7. the upstream request on context.Background()a client that hangs up leaves the engine generatingTestClientDisconnectCancelsUpstream (mutant s09)
8. Header.Add for every upstream headerthe client gets two X-Request-Id valuesTestHandlerInChain (mutant s10)
9. the chain’s handler still checking the static keyevery request is a 401 once real keys replace TL_API_KEYTestHandlerInChain (mutant s12)
DirectionModuleHow it uses this
Backgw.01server.Usage, the Exchange the default observer records on, the Proxy slot
Backgw.00the request path you upgrade; its suite is now your smoke regression
Forwardgw.05sends each routing attempt with Forward and retries only uncommitted failures
Forwarddep.00the gateway image now serves Handler in the chain (inherited call site)
Forwardobs.00the gateway.proxy span is unchanged (inherited call site)
Your pieceProduction equivalentWhat it addsWhere to look
Forward and the commitment boundarysaige’s retry providerbuffer until the first content delta, then never retry (drainUntilContentOrError)saige provider/retry
StreamProxyhttputil.ReverseProxy with FlushInterval: -1the standard library’s streaming proxy, hooks for rewritinghttputil.ReverseProxy (free)
the observerEnvoy AI Gateway token usage extractionusage read from provider streams for metering and limitsEnvoy AI Gateway (free)
disconnect handlingvLLM’s abort on client disconnectthe engine side of the same contractvLLM (free)