Skip to content

fix: harden provider stream lifecycle - #125

Open
Patel230 wants to merge 22 commits into
mainfrom
fix/stream-lifecycle
Open

Patel230 wants to merge 22 commits into
mainfrom
fix/stream-lifecycle

Conversation

@Patel230

@Patel230 Patel230 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Hardens the provider stream lifecycle end to end (provider/core → router → engine), lands the owner's in-progress stream work, fixes the audit findings against it, and refreshes the OSS research notes that ship on this branch.

Every stream now ends exactly once: done, a cancelled terminal (only when the caller's context ends), or an error whose ErrorInfo kind maps to a typed engine code. Upstream timeouts fail over instead of masquerading as user cancellation, routed streams no longer truncate after the first usage event, and cancelling releases the provider connection immediately.

Commits, in order:

Commit What
b3e36ed docs(agents) owner: explain rho's go.local.mod overlay instead of a parent go.work
fde0340 fix(router) owner: report the deployment that actually served a request (DeploymentID, Attempts)
a4a1d39 fix(router) owner: count only deployment-health failures against circuit breakers (+ table test)
d799d7a fix(stream) owner: one terminal event with route and ErrorInfo; EventCancelled; wrappers coordinate every stream path
2943450 docs(changelog) accurate CHANGELOG for the above (F036)
c24743e fix(engine) a stream cancelled before its first event could close with no events and Err() == nil
666c9df fix(stream) only the caller's context signals cancellation (F034)
17b0c9c fix(engine) map ErrorInfo kinds to engine error codes with Retryable (F038)
af91953 fix(stream) cancelled terminal when cancelled mid-delivery (F035)
d4cb0c9 fix(stream) release the provider request before awaiting the terminal; document Close (F039)
60f2a83 fix(router) keep forwarding deployment events after output until done (new, pre-existing on main)
be7b5ca fix(router) warning-marked diagnostics are non-fatal (F044)
b4d0844 perf(stream) skip re-coordinating an already coordinated stream (F045)
0ed463f test(stream) remaining lifecycle guarantees (F040)
b86bfee fix(sse) bound a single SSE event to 16 MiB in core and the Concentrate reader (F256)
1c08f5b fix(cache) response-cache key covers the whole request (F257)
47139c5 docs(research) Bifrost re-snapshot (F301)
e4e6601 docs(research) OpenRouter / Vercel AI Gateway / Models.dev as design donors (F293)
5be33b6 docs(engine) stream terminal and cancellation contract for hosts (F037/F199 flux side)
b559b27 fix(observability) record a streamed interaction before its terminal (regression from wrapping the recorder, caught by shuffled CI)

The owner's uncommitted edits were committed as-is except: restored doc comments the edits had dropped, a gofumpt v0.10.0 fix CI would have rejected, a breaker-classification test, and the CHANGELOG entry (which claimed provider/core "compiles again" after a never-defined symbol; no ref ever failed to build).

Findings addressed

  • F034 (high): inferStreamErrorKind matched bare "cancelled"/any timeout text, and router + engine treated those kinds as the caller's cancellation, so an upstream Client.Timeout (which matches context.DeadlineExceeded) stopped failover, skipped breaker accounting and reached rho as ErrorCancelled. Now only the consumer's own context decides; upstream timeout/cancel events fail over before output, count against the breaker, and surface as retryable ErrorProviderUnavailable. Chat also stops failing over after the caller's deadline.
  • F035: TransformStreamResult dropped the cancelled terminal when cancelled while parked delivering an event; all cancellation paths now share one helper.
  • F036: CHANGELOG rewritten to describe the new behavior under Added/Fixed.
  • F037 (flux side): cancelled terminal documented for hosts (CHANGELOG, docs/architecture/HOST-ENGINE-BOUNDARY.md); rho mapping is a follow-up (below).
  • F038: stream ErrorInfo.Kind → ErrorRateLimited / ErrorAuthentication / ErrorContextExceeded / ErrorInvalidRequest (incl. content filtering) / ErrorProviderUnavailable, with Retryable, provider and model.
  • F039: the source is released (cancel + Close) before the terminal is handed over, so a host that cancels without draining strands one goroutine, not the HTTP body; StreamResult, EventStreamer, engine.Stream and TransformStreamResult document that Close (or draining) is required.
  • F040: tests for cancel-mid-delivery, idempotent Close (core + engine), upstream-timeout failover, caller-cancel single terminal, no events after done, and router error terminals carrying ErrorInfo + route.
  • F044: DeploymentRouter ignores warning-marked diagnostic error events as terminals, like core and the engine.
  • F045: CoordinateStreamResult returns a stream unchanged when it is already coordinated under a context with the same Done channel and the same request ID (Router/ProtocolRouter re-wraps no longer add a goroutine + buffer). Handlers and differently scoped contexts still wrap.
  • F256: ParseSSEStream accumulated an event's data: lines without limit; the Concentrate reader bounded neither lines nor events. Both stop with a stream error past core.SSEMaxEventBytes (16 MiB). Before the fix a 32 MiB event was accepted whole.
  • F257: the cache key hashed only model/system/temperature/message text+tools; it now hashes every ChatOptions and message field under a versioned prefix, and unencodable requests bypass the cache.
  • F293: partially reproduced — the audit's grep -ci openrouter = 0 was wrong (OpenRouter and Vercel appear as exclusions), but neither they nor Models.dev were analyzed as routing/catalog donors. Added a dated addendum with facts verified on 2026-09-27 and Borrow/Adopt rows in the roadmap table.
  • F301: Bifrost re-snapshot (HTTP transport v2.2.3 on 2026-09-24, Core v1.10.4 on 2026-09-25), original snapshot kept.

Found during the work and fixed:

  • Router truncation after output (pre-existing on main): streamWithDeployment ended a successful stream at the first non-output event after output, so Anthropic's post-content usage (and ttft/provider_block) cut off the rest of the answer and surfaced a truncation error under deployment routing.
  • Silent early cancel: engine.Stream cancelled before its first event could close with no events and a nil Err().
  • Double error after output: the router forwarded a failing event and then emitted its own error event (masked by the wrapper; now reported once).
  • Recorder race (introduced by the owner's recorder wrap, caught by CI's -shuffle=on run of TestRecorderStreamRecord): the coordinated wrapper ends the caller's stream at done, but recordStream appended the interaction only after the provider closed its channel. It now records before forwarding the terminal; a deterministic regression test failed before the fix.

Not reproduced / deferred

  • F037 / F199 (rho side): deferred per the campaign plan — rho maps the new event only after a flux release. Verified against this branch: rho builds and its engine and gateway tests pass; fromEngineEvent has no case for EventCancelled or EventProviderBlock, so each cancellation logs gateway: forwarding unrecognized engine event type and forwards cancelled without Error/ErrorInfo, and the trailing Err() still yields error: context canceled.
  • Auto router / Fusion claims from the research refresh were not found in OpenRouter's May–September 2026 announcements and were left out of the research addendum.

Verification

All run from this checkout on darwin/arm64, Go 1.26.6, GOWORK=off:

  • make ci GOFUMPT=<gofumpt v0.10.0> → tidy, verify, fmt, vet, golangci-lint v2.1.0 (0 issues), boundary + layering guards, check-replace, go test ./... -race, govulncheck (no vulnerabilities) — All CI checks passed; tree unchanged afterwards. (The Makefile's fmt uses whatever gofumpt is in GOPATH/bin; overridden to CI's pinned v0.10.0.)
  • go run mvdan.cc/gofumpt@v0.10.0 -l . → no files.
  • go test -race -count=3 ./engine/ ./provider/core/ ./router/ → ok.
  • CI's exact test command go test ./... -race -count=1 -shuffle=on -timeout=300s → ok twice after the recorder fix.
  • deadcode -test ./... → same 224 pre-existing reports as the base branch, none new.
  • markdownlint-cli2 '**/*.md' → 0 issues.
  • go test -fuzz=FuzzBuildCacheKey -fuzztime=15s ./provider → pass.
  • jscpd on touched packages: 4.58% → 4.52% duplicated lines; new clones are pre-existing Chat/StreamChat loop headers.
  • Each fix's regression test was run against the pre-fix commit and failed there (F034 router 4/5 + engine, F035, F039, router truncation, F044, F257; F256 shown with a 32 MiB event accepted before and rejected after).
  • rho compatibility (not committed in rho): GOWORK=off go build -modfile=go.local.mod ./... ok; go test -modfile=go.local.mod ./internal/engine/... ./internal/provider/gateway/... ok.

Follow-ups

  • rho, after a flux release with this PR: in internal/provider/gateway/engine_client.go map fluxengine.EventCancelled (type cancelled, copy Error/ErrorInfo/Route) and EventProviderBlock; skip the trailing Err() error event when a cancelled terminal was forwarded and IsCode(err, ErrorCancelled); handle cancelled as a clean stop in internal/engine/stream.go; add table tests. Do not bump flux in rho before the release.
  • Tag flux v0.0.2 after merge so rho pins a release instead of v0.0.1.
  • Research addendum inferences: OpenRouter provider-preference passthrough and opt-in attribution headers, a normalized order/only/sort hint within approved routing policy, and the planned Models.dev catalog generator with per-field provenance.
  • Consider pinning gofumpt v0.10.0 in the Makefile fmt target so make ci matches CI.

🤖 Generated with Claude Code

…t go.work

A parent go.work breaks every sibling module it does not list, and an
exported GOWORK leaks into child go processes that rho's tests spawn.
Point contributors at rho's gitignored -modfile overlay and note that
without it rho builds against the published flux module.
DeploymentRouter.Chat now attaches a ResolvedRoute with the serving
deployment ID and the number of attempts made, and the engine keeps a
route the provider reported instead of overwriting it with the planned
route. Hosts can attribute usage and price to the backend that answered
after a failover.
…eakers

Caller cancellation, rate limits and 4xx request errors say nothing
about a deployment's health, yet every failure tripped its breaker and
could take a healthy deployment out of rotation. Record a breaker
failure only for 5xx/529 statuses and transport-level errors.
- provider/core: TransformStreamResult emits a terminal "cancelled"
  event when the caller's context ends, Close is idempotent and unblocks
  a parked forwarder, and every terminal error/cancelled event carries a
  StreamErrorInfo inferred from the message when the producer set none.
- engine: add EventCancelled; Stream emits it (with ErrorInfo and route)
  before Err() reports ErrorCancelled, forwards Error/Warning/Route and
  ProviderBlock, and deep-copies mutable event fields.
- router: DeploymentRouter emits route_changed per attempt, annotates
  events with the serving deployment, and stops on caller cancellation;
  Router, ProtocolRouter and the recorder wrap their streams so every
  path shares the same lifecycle.
The drafted entry claimed provider/core "compiles again" after a
never-defined ensureStreamErrorInfo, but no committed ref ever lacked the
symbol or failed to build; the ErrorInfo guarantee is new behavior. List
the cancelled terminal event, the ErrorInfo guarantee and route reporting
under Added, and the breaker classification under Fixed.
…minal

forward raced the route_selected emit against the context; when the
context was already cancelled it could pick ctx.Done, return, and close
the stream with no events and a nil Err(), breaking the one-terminal
guarantee. Emit the cancelled terminal on that path too.
inferStreamErrorKind labels any timeout text and any message containing
"cancelled" as timeout/canceled, and the router and engine treated
those kinds as the caller's cancellation. An upstream read timeout
(net/http Client.Timeout matches context.DeadlineExceeded) or provider
prose such as "subscription cancelled" therefore stopped deployment
failover, skipped breaker accounting and reached hosts as
ErrorCancelled.

- core: only Go's "context canceled" text infers the canceled kind, and
  cancelled events are never retryable.
- router: a failure is cancellation only when the router's context is
  done; otherwise cancelled/timeout events become provider errors that
  fail over before output, are forwarded once after output, count
  against the breaker, and keep their ErrorInfo on the router's own
  terminal event. Chat stops failing over after the caller's deadline.
- engine: map a failure to EventCancelled/ErrorCancelled only when the
  engine's context is done, using the context's own error as the cause;
  otherwise report a retryable ErrorProviderUnavailable.
normalizeEvent turned every fatal stream error into a non-retryable
ErrorProviderUnavailable and ignored ErrorInfo, so hosts could not tell a
401 from a 429 or an outage on the stream path even though provider/core
now guarantees ErrorInfo. Map rate_limited, auth, context_exceeded,
invalid_request and content_filtered to their engine codes, carry
Retryable, and fill the route's provider and model.
When the caller cancelled while TransformStreamResult was parked
delivering an event to a slow consumer, sendLifecycleEvent returned false
and the forwarder closed the channel without the cancelled terminal that
every other cancellation path emits. Route all cancellation paths
through one helper, including that one.
Guaranteed terminal delivery made engine.Stream.forward and every
TransformStreamResult goroutine wait for the consumer after
cancellation, and they closed the source only after that wait. A host
that cancelled the context but neither drained nor closed the stream -
previously enough to tear down - kept the provider HTTP body open
indefinitely.

Release the source (cancel + Close) before handing over the cancelled
terminal, and document on StreamResult, EventStreamer, engine.Stream and
TransformStreamResult that callers must drain or Close; Close stays
idempotent and unblocks the parked goroutine.
streamWithDeployment treated any non-output event after output as the end
of a successful stream. Real providers emit usage (Anthropic's
message_delta), ttft and provider_block events between content events,
so a routed stream stopped after the first of them, dropped the rest of
the answer and its done event, and surfaced as a truncation error.
Forward non-terminal events after output and end only on done.
provider/core emits {Type: error, Warning: ...} diagnostics that precede
the real terminal, and core and the engine forward them without ending
the stream. DeploymentRouter ignored the Warning marker, so under
deployment routing a diagnostic ended the stream with a fatal error (or
triggered a failover before output) and the following done/usage was
lost. Buffer or forward them like other non-terminal events.
Adapters, Router, ProtocolRouter, the recorder and several wrappers each
called CoordinateStreamResult, so one request stacked 3-5 identical
lifecycle wrappers, each with its own goroutine, buffer and ErrorInfo
pass. Record the context and request ID coordinating each live stream
and return the source unchanged when a caller would wrap it again under
a context with the same Done channel: the extra wrapper could not behave
differently. Handlers and differently scoped contexts still wrap.
Cover idempotent engine Close (source closed and request cancelled
exactly once), no events after done through the engine's public Stream
API, and router error terminals that keep the adapter's ErrorInfo kind
and the failing deployment's route, alongside the cancellation, failover
and mid-delivery tests added with their fixes.
ParseSSEStream capped each line at 2 MiB but appended every data: line of
an event to a strings.Builder until a blank line, so a hostile or broken
endpoint could stream an endless run of data lines and grow client
memory without bound. The Concentrate Responses reader used
ReadString('\n'), bounding neither lines nor events. Stop with a stream
error once an event's accumulated fields exceed core.SSEMaxEventBytes
(16 MiB) in both parsers.
buildCacheKey hashed only the model, system prompt, temperature and each
message's role, content and tool data. With caching enabled, a reply
generated for one tool set, max_tokens, stop sequence, top_p/top_k,
thinking or response-format setting, image or caller identity was
served to a request that differed only in those fields, and a caller
could pre-seed replies for others sending the same messages.

Hash the JSON of every ChatOptions and message field under a versioned
prefix; a request that cannot be encoded returns an empty key and
bypasses the cache.
The 2026-09-24 snapshot listed helm-chart-v2.1.43 and HTTP transport
2.1.0; the releases API now shows HTTP transport v2.2.3 (2026-09-24,
pinned provider keys per routing fallback) and Core v1.10.4
(2026-09-25). Keep the original snapshot and add the dated re-check.
…as donors

The landscape excluded the hosted meta-gateways from the OSS peer set but
never analyzed them as routing or catalog donors, although flux ships an
openrouter adapter that passes no provider preferences and plans a
Models.dev-generated catalog. Add a dated addendum with the verified
OpenRouter provider-routing fields and 2026 announcements, Vercel's
order/only/sort, caching and BYOK options, and Models.dev's data layout,
plus Borrow/Adopt rows in the roadmap table. Unverified claims from the
research refresh (Auto router, Fusion, in-region hostnames) are recorded
as gaps, not facts.
Hosts need to know that EventCancelled is a terminal of its own, that
Err() then reports ErrorCancelled (so no second terminal should be
emitted), how provider failures map to error codes, and that cancelling
the context does not replace Close.
The recorder's output is now wrapped by CoordinateStreamResult, which
ends the caller's stream at the terminal event. recordStream appended
the interaction only after the provider closed its channel, so a caller
that drained the stream and saved the cassette could miss it (seen as a
shuffled-CI failure of TestRecorderStreamRecord). Save the interaction
before forwarding the terminal event; an early close by the caller still
records nothing.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant