Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 1 addition & 1 deletion docs/subsystems/gateway-routing.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ served by a plugin are in [Provider and plugin ownership](provider-plugins.md).

- A fallback happens only while none of the reply has been sent. An agent never gets half a reply from one upstream and the rest from another.
- A 2xx that can't be an API's answer is the 502 it stands for (`notAnAPIReply`, #1012): a web page (`text/html`, or a body that begins as one: a sign-in page, Cloudflare's challenge), a body with nothing in it, or one sent as JSON that doesn't begin as JSON. Every upstream request goes through it (`forwardOnce`, and the built-in clients' own: the ChatGPT backend's `/responses`, Kiro, Zed, Qoder, Devin, Cursor, Command Code), so it fails over, rests and is logged as a 502 does, its reason naming the page's title. Only the first bytes are peeked, which a stream waits for anyway, and a compressed body is left as it came.
- A streamed request is never left silent while its vendor says nothing. `holdWriter.watch` runs beside each try: past `keepHeldAfter` (15s) it sends a held try's agent the stream's 200 and SSE comments (`keepAlive`), none of the try, so another candidate may still answer in the same stream; once the reply passes, it puts a comment between two of its lines after `keepaliveEvery` of quiet. It stops after `keepaliveLongest` (5 minutes) with nothing from the vendor. A proxy in front such as Cloudflare (125s) otherwise ends the request (#947). Gemini streams get no comments (#934).
- A streamed request is never left silent while its vendor says nothing. `holdWriter.watch` runs beside each try: past `keepHeldAfter` (15s) it sends a held try's agent the stream's 200 and SSE comments (`keepAlive`), none of the try, so another candidate may still answer in the same stream; once the reply passes, it puts a comment between two of its lines after `keepaliveEvery` of quiet. It stops after `keepaliveLongest` (5 minutes) with nothing from the vendor. A proxy in front such as Cloudflare (125s) otherwise ends the request (#947). Gemini streams get no comments (#934). Those comments tell the agent the reply is a stream, so a vendor that answers whole `application/json` instead leaves that stream with a reply that is not one: it is released as the vendor not streaming one, the same reason a request nothing went to ahead of it is answered 502 for, and not as the stream's error carrying the answer, which a 200 is not.
- Claude's, GPT's and Gemini's reasoning is held (`refusesAfterThinking`, up to `holdThinking`, 4 minutes) so that a refusal after it goes to the next candidate unseen (#248). It is held only while one could: on the last candidate a refusal isn't asked again, so Claude's and GPT's reasoning streams as it comes there, and a refusal after it reaches the agent as the vendor's own. Held there, a lone relay's Claude or GPT showed Claude Code its thinking all at once with the text, and Codex none until then (iTianbao on X). Gemini's stays held on the last candidate too, for a reply that only reasoned to be asked again (#667). Any other model's reasoning is never held.
- The order of the steps a request's accounts go through, how much the accounts' own order (*Make first*) counts under each routing, and how often an allowance is read are described for users in [How magpie picks an account](../reference.md#how-magpie-picks-an-account). Smart compares renewals window by window, biggest first, each truncated to the hour (`weighRouted`), so the order decides only when every window falls in the same hour. Affinity keeps a conversation on an account at 90–97%: only `SpentShareOf` (98%, 100% In order) stops it (`affine`), while the 90% tier (`lowShare`) orders only conversations nobody has answered.
- A resting candidate is tried last, never dropped (`restLast`). A lone candidate is tried even while it rests, because there is no other.
Expand Down
12 changes: 10 additions & 2 deletions internal/gateway/fallback.go
Original file line number Diff line number Diff line change
Expand Up @@ -1428,9 +1428,17 @@ func (h *holdWriter) release() {
return
}
if h.alive != nil && h.alive.sent && !h.stream {
// an error status, once the agent has a stream: told as its error
// the stream's 200 and comments went ahead of a reply that is not
// one. An error status is that error as it stands. A 2xx the vendor
// did not stream, said as the same words a vendor that answers a
// request nothing went to ahead of it is told in — not as the
// stream's error carrying its answer, which a 200 is not
h.passing = true
streamError(h.w, h.alive.proto, h.status, provider.APIError(h.held.Bytes(), http.StatusText(h.status)))
said := http.StatusText(h.status)
if h.status < 300 {
said = "did not stream"
}
streamError(h.w, h.alive.proto, h.status, provider.APIError(h.held.Bytes(), said))
return
}
h.pass()
Expand Down
96 changes: 96 additions & 0 deletions internal/gateway/quiet_keepalive_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -309,3 +309,99 @@ func TestKeepAliveBeforeTheVendorsHeaders(t *testing.T) {
t.Fatalf("body: %q", rec.Body.String())
}
}

// quietWholeGap is how long quietWhole's vendor says nothing before it
// answers: long enough for watch to have had many a look at it past
// keepHeldAfter (40ms in quietFast).
const quietWholeGap = 300 * time.Millisecond

const (
wholeChatAnswer = `{"id":"chatcmpl-1","object":"chat.completion","choices":[{"index":0,"message":{"role":"assistant","content":"the whole answer"},"finish_reason":"stop"}],"usage":{"prompt_tokens":3,"completion_tokens":2}}`
wholeResponsesAnswer = `{"id":"resp-1","object":"response","output":[{"type":"message","role":"assistant","content":[{"type":"output_text","text":"the whole answer"}]}]}`
)

// quietWhole serves a provider on endpoint that says nothing at all for
// quietWholeGap and then answers whole, with no event stream in it.
func quietWhole(t *testing.T, id string, endpoint provider.Protocol, body string) {
t.Helper()
up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
io.ReadAll(r.Body)
time.Sleep(quietWholeGap)
w.Header().Set("Content-Type", "application/json")
io.WriteString(w, body)
}))
t.Cleanup(up.Close)
p := provider.Provider{ID: id, Name: "Fixture", Key: "k", Models: []string{"gpt-test"}}
if endpoint == provider.Responses {
p.Responses = up.URL + "/v1"
} else {
p.Chat = up.URL + "/v1"
}
if err := provider.Save(p); err != nil {
t.Fatal(err)
}
}

// askWhole posts body to the gateway and reads the whole of its answer: a
// reply the vendor did not stream is not read as one.
func askWhole(t *testing.T, path, body string) (int, string, string) {
t.Helper()
gw := httptest.NewServer(New().Handler())
t.Cleanup(gw.Close)
resp, err := http.Post(gw.URL+path, "application/json", strings.NewReader(body))
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
b, err := io.ReadAll(resp.Body)
if err != nil {
t.Fatal(err)
}
return resp.StatusCode, resp.Header.Get("Content-Type"), string(b)
}

// A vendor quiet for longer than keepHeldAfter, with nothing heard from it
// at all, and then a whole, non-streamed application/json 200. watch keeps
// the agent of a stream alive while the vendor has yet to answer, which
// tells the agent the reply will be a stream, so the whole reply is held
// for a stream that never was: it is told the vendor did not stream, as the
// same vendor answering a request nothing went to ahead of it is, and not
// as that stream's error with its answer inside it.
func TestQuietThenWholeJSONSaysTheVendorDidNotStream(t *testing.T) {
fresh(t)
quietFast(t)
quietWhole(t, "fixture", provider.Chat, wholeChatAnswer)
code, _, got := askWhole(t, "/v1/chat/completions", `{"model":"fixture/gpt-test","stream":true,"messages":[{"role":"user","content":"hi"}]}`)
if !strings.Contains(got, "did not stream") {
t.Fatalf("%d %q", code, got)
}
if strings.Contains(got, "OK:") {
t.Fatalf("a 200 was told as the stream's error: %q", got)
}
}

// The same vendor, an agent that asked for no stream: no watch, no
// keepalive, and the whole reply goes as it is.
func TestQuietThenWholeJSONNoStream(t *testing.T) {
fresh(t)
quietFast(t)
quietWhole(t, "fixture", provider.Chat, wholeChatAnswer)
code, _, got := askWhole(t, "/v1/chat/completions", `{"model":"fixture/gpt-test","messages":[{"role":"user","content":"hi"}]}`)
if code != http.StatusOK || !strings.Contains(got, `"content":"the whole answer"`) || strings.Contains(got, `"error"`) {
t.Fatalf("%d %q", code, got)
}
}

// A provider with no Chat endpoint at all, quiet and then a whole JSON
// 200: the agent that asked for a stream has already been given it, so the
// provider not streaming one is told as that stream's error, with the
// vendor's own reason and not a 200's.
func TestQuietThenWholeJSONFromResponsesOnly(t *testing.T) {
fresh(t)
quietFast(t)
quietWhole(t, "fixture", provider.Responses, wholeResponsesAnswer)
_, _, got := askWhole(t, "/v1/chat/completions", `{"model":"fixture/gpt-test","stream":true,"messages":[{"role":"user","content":"hi"}]}`)
if !strings.Contains(got, "did not stream") || strings.Contains(got, "OK:") {
t.Fatalf("the stream's agent was told: %q", got)
}
}
Loading