Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
22 commits
Select commit Hold shift + click to select a range
cb95acf
docs: add OSS competitive roadmap
Patel230 Sep 24, 2026
2832455
fix: harden stream lifecycle
Patel230 Sep 24, 2026
b3e36ed
docs(agents): explain the rho go.local.mod overlay instead of a paren…
Patel230 Sep 26, 2026
fde0340
fix(router): report the deployment that actually served a request
Patel230 Sep 26, 2026
a4a1d39
fix(router): count only deployment-health failures against circuit br…
Patel230 Sep 26, 2026
d799d7a
fix(stream): guarantee one terminal event with route and error info
Patel230 Sep 26, 2026
2943450
docs(changelog): describe the stream lifecycle changes accurately
Patel230 Sep 26, 2026
c24743e
fix(engine): end a stream cancelled before its first event with a ter…
Patel230 Sep 26, 2026
666c9df
fix(stream): let only the caller's context signal cancellation
Patel230 Sep 26, 2026
17b0c9c
fix(engine): map stream ErrorInfo kinds to engine error codes
Patel230 Sep 26, 2026
af91953
fix(stream): emit the cancelled terminal when cancelled mid-delivery
Patel230 Sep 26, 2026
d4cb0c9
fix(stream): release the provider request before awaiting the terminal
Patel230 Sep 26, 2026
60f2a83
fix(router): keep forwarding deployment events after output until done
Patel230 Sep 26, 2026
be7b5ca
fix(router): treat warning-marked diagnostics as non-fatal
Patel230 Sep 26, 2026
b4d0844
perf(stream): skip re-coordinating an already coordinated stream
Patel230 Sep 26, 2026
0ed463f
test(stream): pin the remaining stream lifecycle guarantees
Patel230 Sep 26, 2026
b86bfee
fix(sse): bound the size of a single SSE event
Patel230 Sep 26, 2026
1c08f5b
fix(cache): key cached responses on the whole request
Patel230 Sep 26, 2026
47139c5
docs(research): re-snapshot Bifrost releases in the gateway landscape
Patel230 Sep 26, 2026
e4e6601
docs(research): analyze OpenRouter, Vercel AI Gateway and Models.dev …
Patel230 Sep 26, 2026
5be33b6
docs(engine): document stream terminals and cancellation for hosts
Patel230 Sep 26, 2026
b559b27
fix(observability): record a streamed interaction before its terminal
Patel230 Sep 26, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 11 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,10 +58,19 @@ make ci # Full CI suite
- `provider/core.Provider` is the lower-level provider contract; keep its
method set stable and use it across feature packages
- Streaming tests need careful goroutine management
- `go.work` here should stay minimal; the parent `graycode-eco/go.work`
connects this independent `flux` checkout beside Rho for local development.
- `go.work` here should stay minimal; the `graycode-eco` workspace connects
this independent `flux` checkout beside Rho for local development.
Do not add extra local `replace` directives here without coordinating with
the parent workspace.
- There is usually **no** `go.work` in the parent folder, and that is expected.
A parent `go.work` breaks every sibling it does not list, and exporting
`GOWORK` is worse: it is inherited by child `go` processes, so rho's own
tests that shell out to `go test` in temporary projects fail. To exercise
rho against this checkout, use a gitignored module overlay in rho instead —
`go test -modfile=go.local.mod ./...`. See rho/AGENTS.md
"Workspace workflow" for the setup. Without that overlay, rho compiles
against the *published* flux from the module cache, so a green rho suite
does not exercise local flux changes at all.

## Naming Conventions

Expand Down
61 changes: 61 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,73 @@ Format: [Keep a Changelog](https://keepachangelog.com/en/1.0.0/) · Versioning:
credential lookup is required for the replicated route.
- Client-owned OpenAI-compatible provider registration through
`FluxClient.RegisterCustomProvider`.
- Streams end with exactly one terminal event. When the caller's context is
cancelled or its deadline passes, `provider/core` stream wrappers emit a
terminal `cancelled` event and the engine emits `engine.EventCancelled`
(with `ErrorInfo` and the route) before `Err()` reports `ErrorCancelled`.
Hosts that switch on event types should handle `cancelled`. Cancelling
the context releases the provider request immediately, but the terminal is
delivered like any other event, so callers must still read the stream to
the end or call `Close` (as `StreamResult`, `EventStreamer` and
`engine.Stream` now document).
- Terminal `error`/`cancelled` stream events leaving `provider/core` always
carry a `StreamErrorInfo` (`Kind`/`Retryable`), inferred from the provider
message when the adapter set none. Existing `ErrorInfo` is never
overwritten, and a cancellation never reports as an internal fault.
- Responses and stream events report the route that actually served the
request: `ResolvedRoute.DeploymentID` and `Attempts` after a deployment
failover, a `route_changed` event per deployment attempt, and the route on
engine events.

### Fixed
- Circuit breakers now admit at most one concurrent half-open probe and do
not reserve probes during route filtering.
- `DeploymentRouter` records a circuit-breaker failure only for errors that
describe the deployment's health (5xx, 529, transport failures). Caller
cancellation, rate limits and 4xx request errors no longer take a healthy
deployment out of rotation.
- Upstream timeouts and provider messages that merely contain "cancelled" are
no longer mistaken for the caller's cancellation. Only the caller's own
context decides that: `DeploymentRouter` fails over from an upstream
timeout (including net/http `Client.Timeout`) and counts it against the
deployment, and the engine reports it as a retryable
`ErrorProviderUnavailable` instead of `ErrorCancelled`.
- Engine stream errors keep the provider's classification: a stream error's
`ErrorInfo.Kind` now maps to `ErrorRateLimited`, `ErrorAuthentication`,
`ErrorContextExceeded` or `ErrorInvalidRequest` (content filtering
included), carries `Retryable`, and names the route's provider and model,
instead of every stream failure becoming a non-retryable
`ErrorProviderUnavailable`.
- `DeploymentRouter` streams no longer end at the first non-output event
after output. A `usage`, `ttft` or `provider_block` event between content
events (Anthropic reports output usage before `message_stop`) used to end
the deployment stream, dropping the remaining content and `done` and
surfacing a truncation error.
- `DeploymentRouter` treats warning-marked diagnostic `error` events (for
example a reasoning-only response) as non-fatal, like `provider/core` and
the engine already did, instead of ending the stream or failing over.
- SSE parsing is bounded per event. `core.ParseSSEStream` capped each line at
2 MiB but accumulated an event's `data:` lines without limit, and the
Concentrate Responses reader bounded neither lines nor events, so a hostile
or broken endpoint could grow client memory indefinitely. Both now stop
with a stream error once one event exceeds `core.SSEMaxEventBytes`
(16 MiB).
- The response cache key now covers the whole request. It hashed only the
model, system prompt, temperature and message text/tool data, so with
caching enabled a reply produced for one tool set, `max_tokens`, stop
sequence, sampling or thinking setting, response format, image or caller
was served to a different request. Keys now hash every `ChatOptions` and
message field under a versioned prefix, and requests that cannot be
encoded bypass the cache.
- A `provider/core` stream cancelled while its consumer was behind (the
forwarder blocked delivering an event) now still ends with the terminal
`cancelled` event instead of closing silently.

### Changed
- `core.CoordinateStreamResult` returns a stream unchanged when it is already
coordinated under the same context and request ID, so layers that re-wrap
an adapter's stream with the caller's context (`Router`, `ProtocolRouter`)
no longer add a goroutine and buffer per layer.
- Removed process-global custom gateway and dynamic provider registration,
the ambient `OPENAI_API_BASE` auto-registration path, and no-op API-key
prefix inference. Custom gateways and endpoints now require explicit,
Expand Down
18 changes: 18 additions & 0 deletions docs/architecture/HOST-ENGINE-BOUNDARY.md
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,24 @@ additive and must be ignored safely. Flux emits tool requests; Rho authorizes
and executes tools, appends results to its history, and begins the next model
turn.

Every stream ends exactly once:

- `done` — success; `Err()` stays nil.
- `cancelled` — the request context was cancelled or its deadline passed. The
event carries `ErrorInfo` (`canceled` or `timeout`) and the route, and
`Err()` then reports `ErrorCancelled` wrapping the context's error. Hosts
should treat it as the user's cancellation, not as a failure, and must not
emit a second terminal for it.
- no terminal event — `Next` returns false and `Err()` reports the provider
failure with the code from its `ErrorInfo` (`rate_limited`,
`authentication_failed`, `context_exceeded`, `invalid_request`, or
`provider_unavailable`) and `Retryable`. An upstream timeout is such a
retryable failure, never `cancelled`.

Cancelling the context releases the provider request at once, but the
terminal event waits for the host to read it, so hosts must still read to the
end or call `Close`.

## Readiness

Preflight has two explicit modes:
Expand Down
2 changes: 1 addition & 1 deletion docs/guides/DYNAMIC-MODEL-DISCOVERY.md
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ the safe credential/gateway reports without reading either file directly.
| catalog cache missing/corrupt | report unavailable; do not claim bootstrap readiness |
| live list fails or selected model is absent | fail live preflight |
| custom gateway URL contains embedded data | reject configuration |
| stream caller exits | close/cancel the Engine stream |
| stream caller exits | close the Engine stream (cancelling its context alone leaves the terminal event waiting for a reader) |

Provider-specific friendly error formatting remains Flux-owned; Rho decides
where and how to display it.
Expand Down
3 changes: 3 additions & 0 deletions engine/classify.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ func classify(operation string, route Route, err error) error {
switch {
case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded):
code = ErrorCancelled
case errors.Is(err, core.ErrStreamTruncated):
code = ErrorProviderUnavailable
retryable = true
default:
var providerErr *core.FluxError
if errors.As(err, &providerErr) {
Expand Down
15 changes: 11 additions & 4 deletions engine/continuation.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ func streamWithContinuation(ctx context.Context, provider core.Provider, message
var segment strings.Builder
hadToolCall := false
var terminal core.FluxStreamEvent
segmentLoop:
for event := range current.Events {
switch event.Type {
case "content":
Expand All @@ -52,20 +53,26 @@ func streamWithContinuation(ctx context.Context, provider core.Provider, message
}
case "done":
terminal = event
continue
break segmentLoop
case "error":
if event.Warning == "" {
_ = emitEngineEvent(streamCtx, out, event)
current.Close()
return
}
}
if !emitEngineEvent(streamCtx, out, event) {
current.Close()
return
}
}
current.Close()
if terminal.Type == "" {
return
}

needsContinuation := terminal.StopReason == "max_tokens" || terminal.StopReason == "length"
if !needsContinuation || hadToolCall || totalOutput >= maxTotalTokens || attempt >= maxContinuations {
if terminal.Type == "" {
terminal = core.FluxStreamEvent{Type: "done", StopReason: terminal.StopReason, RequestID: requestID}
}
_ = emitEngineEvent(streamCtx, out, terminal)
return
}
Expand Down
37 changes: 37 additions & 0 deletions engine/contract_e2e_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,43 @@ func TestEngineContinuationPreservesEventsAndConversationShape(t *testing.T) {
}
}

type terminalErrorProvider struct{}

func (*terminalErrorProvider) Name() string { return "terminal-error" }
func (*terminalErrorProvider) Ping(context.Context) error { return nil }
func (*terminalErrorProvider) Chat(context.Context, []core.FluxMessage, core.ChatOptions) (*core.FluxResponse, error) {
return nil, nil
}

func (*terminalErrorProvider) StreamChat(context.Context, []core.FluxMessage, core.ChatOptions) (*core.StreamResult, error) {
events := make(chan core.FluxStreamEvent, 2)
events <- core.FluxStreamEvent{Type: "error", Error: "connection reset"}
events <- core.FluxStreamEvent{Type: "done", StopReason: "stop"}
close(events)
return llm.NewStreamResult(events, "request-error", nil), nil
}

func TestEngineContinuationDoesNotAppendDoneAfterFatalError(t *testing.T) {
source, err := streamWithContinuation(
context.Background(), &terminalErrorProvider{},
[]core.FluxMessage{{Role: "user", Content: "hello"}},
core.ChatOptions{Model: "model"},
Limits{MaxContinuations: 1, MaxTotalOutputTokens: 100},
)
if err != nil {
t.Fatal(err)
}
defer source.Close()

var events []core.FluxStreamEvent
for event := range source.Events {
events = append(events, event)
}
if len(events) != 1 || events[0].Type != "error" || events[0].Error != "connection reset" {
t.Fatalf("events = %+v, want one fatal error", events)
}
}

func firstCatalogModelID(cat catalog.Catalog) string {
for id := range cat.Models {
return id
Expand Down
36 changes: 32 additions & 4 deletions engine/convert.go
Original file line number Diff line number Diff line change
Expand Up @@ -85,14 +85,42 @@ func cloneStringMap(in map[string]string) map[string]string {
return out
}

func cloneRoute(route *Route) *Route {
if route == nil {
return nil
}
cloned := *route
return &cloned
}

// mergeRoute keeps the route the provider reported and fills any blank
// identity fields from the route the engine planned.
func mergeRoute(actual *Route, planned Route) *Route {
if actual == nil {
return cloneRoute(&planned)
}
merged := *actual
if merged.Provider == "" {
merged.Provider = planned.Provider
}
if merged.Model == "" {
merged.Model = planned.Model
}
if !merged.DeploymentRouting {
merged.DeploymentRouting = planned.DeploymentRouting
}
return &merged
}

// fromClientResponse attaches the resolved route to a client response. The
// engine and the client both speak the canonical contract response type, so
// this only sets the route the engine selected.
// engine and the client both speak the canonical contract response type. A
// route the provider already reported (for example the deployment that served
// the request after a failover) wins; the planned route only fills blanks.
func fromClientResponse(resp *core.FluxResponse, route Route) *GenerateResponse {
if resp == nil {
return &GenerateResponse{Route: &route}
return &GenerateResponse{Route: cloneRoute(&route)}
}
resp.Route = &route
resp.Route = mergeRoute(resp.Route, route)
return resp
}

Expand Down
17 changes: 17 additions & 0 deletions engine/convert_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,23 @@ func TestFromClientResponse_AttachesRoute(t *testing.T) {
}
}

func TestFromClientResponse_PreservesActualRoute(t *testing.T) {
resp := &core.FluxResponse{
Content: "hello",
Route: &Route{
Provider: "router", Model: "planned/model", DeploymentID: "deployment-2", Attempts: 2,
},
}
planned := Route{Provider: "planned", Model: "planned/model", DeploymentRouting: true}
out := fromClientResponse(resp, planned)
if out.Route == nil || out.Route.DeploymentID != "deployment-2" || out.Route.Attempts != 2 || !out.Route.DeploymentRouting {
t.Fatalf("route = %+v, want actual route preserved", out.Route)
}
if out.Route.Provider != "router" || out.Route.Model != "planned/model" {
t.Fatalf("route identity = %+v, want provider route preserved", out.Route)
}
}

func TestFromClientResponse_NilResponse(t *testing.T) {
route := Route{Provider: "test", Model: "test/model"}
out := fromClientResponse(nil, route)
Expand Down
71 changes: 71 additions & 0 deletions engine/engine_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -240,6 +240,77 @@ func TestStreamFatalErrorEventStillTerminal(t *testing.T) {
}
}

func TestStreamSilentSourceCloseIsTruncated(t *testing.T) {
sourceEvents := make(chan core.FluxStreamEvent)
close(sourceEvents)

ctx, cancel := context.WithCancel(context.Background())
stream := newStream(ctx, cancel, llm.NewStreamResult(sourceEvents, "", nil), Route{Provider: "mock", Model: "mock/model"})
defer stream.Close()

var events []Event
for stream.Next() {
events = append(events, stream.Event())
}
err := stream.Err()
if !IsCode(err, ErrorProviderUnavailable) {
t.Fatalf("error = %v, want provider_unavailable", err)
}
if !errors.Is(err, core.ErrStreamTruncated) {
t.Fatalf("error = %v, want stream truncation cause", err)
}
if len(events) != 1 || events[0].Type != EventRouteSelected {
t.Fatalf("events = %+v, want only route_selected", events)
}
}

func TestStreamContextCancellationEmitsTerminalAndErr(t *testing.T) {
sourceEvents := make(chan core.FluxStreamEvent)
ctx, cancel := context.WithCancel(context.Background())
stream := newStream(ctx, cancel, llm.NewStreamResult(sourceEvents, "request-cancel", nil), Route{Provider: "mock", Model: "mock/model"})

if !stream.Next() {
t.Fatal("expected route event")
}
cancel()

if !stream.Next() {
t.Fatal("expected cancellation event")
}
event := stream.Event()
if event.Type != EventCancelled || event.ErrorInfo == nil || event.ErrorInfo.Kind != llm.ErrKindCanceled {
t.Fatalf("event = %+v, want cancellation terminal", event)
}
if stream.Next() {
t.Fatal("unexpected event after cancellation terminal")
}
if err := stream.Err(); !IsCode(err, ErrorCancelled) {
t.Fatalf("error = %v, want cancelled", err)
}
_ = stream.Close()
}

func TestStreamCancelledBeforeFirstEventStillTerminates(t *testing.T) {
// forward races the route_selected emit against the already-cancelled
// context; every run must still end with the cancelled terminal.
for i := 0; i < 200; i++ {
ctx, cancel := context.WithCancel(context.Background())
cancel()
stream := newStream(ctx, cancel, llm.NewStreamResult(make(chan core.FluxStreamEvent), "request-early", nil), Route{Provider: "mock", Model: "mock/model"})
var last Event
for stream.Next() {
last = stream.Event()
}
if last.Type != EventCancelled {
t.Fatalf("run %d: last event = %+v, want cancelled terminal", i, last)
}
if err := stream.Err(); !IsCode(err, ErrorCancelled) {
t.Fatalf("run %d: error = %v, want cancelled", i, err)
}
_ = stream.Close()
}
}

func TestSnapshotPublishesCapabilities(t *testing.T) {
compiled := &catalog.CompiledCatalog{
ModelsByID: map[string]catalog.Model{
Expand Down
Loading
Loading