fix: harden provider stream lifecycle - #125
Merged
Merged
Conversation
…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.
Patel230
added a commit
that referenced
this pull request
Oct 2, 2026
* docs: add OSS competitive roadmap
* fix: harden stream lifecycle
* docs(agents): explain the rho go.local.mod overlay instead of a parent 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.
* fix(router): report the deployment that actually served a request
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.
* fix(router): count only deployment-health failures against circuit breakers
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.
* fix(stream): guarantee one terminal event with route and error info
- 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.
* docs(changelog): describe the stream lifecycle changes accurately
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.
* fix(engine): end a stream cancelled before its first event with a terminal
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.
* fix(stream): let only the caller's context signal cancellation
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.
* fix(engine): map stream ErrorInfo kinds to engine error codes
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.
* fix(stream): emit the cancelled terminal when cancelled mid-delivery
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.
* fix(stream): release the provider request before awaiting the terminal
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.
* fix(router): keep forwarding deployment events after output until done
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.
* fix(router): treat warning-marked diagnostics as non-fatal
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.
* perf(stream): skip re-coordinating an already coordinated stream
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.
* test(stream): pin the remaining stream lifecycle guarantees
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.
* fix(sse): bound the size of a single SSE event
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.
* fix(cache): key cached responses on the whole request
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.
* docs(research): re-snapshot Bifrost releases in the gateway landscape
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.
* docs(research): analyze OpenRouter, Vercel AI Gateway and Models.dev 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.
* docs(engine): document stream terminals and cancellation for hosts
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.
* fix(observability): record a streamed interaction before its terminal
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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, acancelledterminal (only when the caller's context ends), or an error whoseErrorInfokind maps to a typed engine code. Upstream timeouts fail over instead of masquerading as user cancellation, routed streams no longer truncate after the firstusageevent, and cancelling releases the provider connection immediately.Commits, in order:
b3e36eddocs(agents)go.local.modoverlay instead of a parentgo.workfde0340fix(router)DeploymentID,Attempts)a4a1d39fix(router)d799d7afix(stream)ErrorInfo;EventCancelled; wrappers coordinate every stream path2943450docs(changelog)c24743efix(engine)Err() == nil666c9dffix(stream)17b0c9cfix(engine)ErrorInfokinds to engine error codes withRetryable(F038)af91953fix(stream)d4cb0c9fix(stream)60f2a83fix(router)done(new, pre-existing on main)be7b5cafix(router)b4d0844perf(stream)0ed463ftest(stream)b86bfeefix(sse)1c08f5bfix(cache)47139c5docs(research)e4e6601docs(research)5be33b6docs(engine)b559b27fix(observability)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
inferStreamErrorKindmatched bare "cancelled"/any timeout text, and router + engine treated those kinds as the caller's cancellation, so an upstreamClient.Timeout(which matchescontext.DeadlineExceeded) stopped failover, skipped breaker accounting and reached rho asErrorCancelled. Now only the consumer's own context decides; upstream timeout/cancel events fail over before output, count against the breaker, and surface as retryableErrorProviderUnavailable.Chatalso stops failing over after the caller's deadline.TransformStreamResultdropped the cancelled terminal when cancelled while parked delivering an event; all cancellation paths now share one helper.cancelledterminal documented for hosts (CHANGELOG,docs/architecture/HOST-ENGINE-BOUNDARY.md); rho mapping is a follow-up (below).ErrorInfo.Kind→ErrorRateLimited/ErrorAuthentication/ErrorContextExceeded/ErrorInvalidRequest(incl. content filtering) /ErrorProviderUnavailable, withRetryable, provider and model.StreamResult,EventStreamer,engine.StreamandTransformStreamResultdocument that Close (or draining) is required.done, and router error terminals carryingErrorInfo+ route.DeploymentRouterignores warning-marked diagnosticerrorevents as terminals, like core and the engine.CoordinateStreamResultreturns 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.ParseSSEStreamaccumulated an event'sdata:lines without limit; the Concentrate reader bounded neither lines nor events. Both stop with a stream error pastcore.SSEMaxEventBytes(16 MiB). Before the fix a 32 MiB event was accepted whole.ChatOptionsand message field under a versioned prefix, and unencodable requests bypass the cache.grep -ci openrouter = 0was 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.Found during the work and fixed:
main):streamWithDeploymentended a successful stream at the first non-output event after output, so Anthropic's post-contentusage(andttft/provider_block) cut off the rest of the answer and surfaced a truncation error under deployment routing.engine.Streamcancelled before its first event could close with no events and a nilErr().-shuffle=onrun ofTestRecorderStreamRecord): the coordinated wrapper ends the caller's stream atdone, butrecordStreamappended 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
fromEngineEventhas no case forEventCancelledorEventProviderBlock, so each cancellation logsgateway: forwarding unrecognized engine event typeand forwardscancelledwithoutError/ErrorInfo, and the trailingErr()still yieldserror: context canceled.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'sfmtuses whatever gofumpt is inGOPATH/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.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.GOWORK=off go build -modfile=go.local.mod ./...ok;go test -modfile=go.local.mod ./internal/engine/... ./internal/provider/gateway/...ok.Follow-ups
internal/provider/gateway/engine_client.gomapfluxengine.EventCancelled(typecancelled, copyError/ErrorInfo/Route) andEventProviderBlock; skip the trailingErr()error event when a cancelled terminal was forwarded andIsCode(err, ErrorCancelled); handlecancelledas a clean stop ininternal/engine/stream.go; add table tests. Do not bump flux in rho before the release.order/only/sorthint within approved routing policy, and the planned Models.dev catalog generator with per-field provenance.fmttarget somake cimatches CI.🤖 Generated with Claude Code