diff --git a/go-api/internal/gateway/failover.go b/go-api/internal/gateway/failover.go index 0f4ad9a..e85dfff 100644 --- a/go-api/internal/gateway/failover.go +++ b/go-api/internal/gateway/failover.go @@ -39,6 +39,33 @@ func NewFailover(primary Gateway, rest ...Gateway) Gateway { return &failover{providers: append([]Gateway{primary}, rest...)} } +// Standby is a gateway that has somewhere else to go. +// +// The runtime needs this and must NOT learn what a provider is. A mid-run +// failure cannot be moved by this package — the conversation is half built and +// its tool calls belong to whoever issued them (see canFailOver) — so the only +// thing that can rescue it is starting the run again somewhere else, and only +// the loop can do that. This is the whole of what the loop is told: "there is +// another one, here it is", with no vendor, credential or model id crossing the +// boundary. +type Standby interface { + // Standby returns a gateway beginning at the NEXT provider, and whether + // there was one. The receiver is unchanged. + Standby() (Gateway, bool) +} + +// Standby drops the provider that just failed and returns the rest. +// +// The remainder keeps its own fallbacks, so a second failure on a three +// provider deployment still has somewhere to go. With one provider left there +// is no wrapper at all, which is NewFailover's own rule. +func (f *failover) Standby() (Gateway, bool) { + if len(f.providers) < 2 { + return nil, false + } + return NewFailover(f.providers[1], f.providers[2:]...), true +} + func (f *failover) Complete(ctx context.Context, req Request) (*Response, error) { var last error for i, p := range f.providers { diff --git a/go-api/internal/runtime/loop.go b/go-api/internal/runtime/loop.go index 98a3e94..418101a 100644 --- a/go-api/internal/runtime/loop.go +++ b/go-api/internal/runtime/loop.go @@ -127,7 +127,67 @@ func (m *ModelExecutor) ExecuteAgent(ctx context.Context, agent *Agent, input Ex func (m *ModelExecutor) executeWithLimits( ctx context.Context, agent *Agent, input ExecutionInput, limits Limits, ) (*ExecutionResult, error) { - return m.executeRun(ctx, agent, input, limits, delegation{}) + res, err := m.executeRun(ctx, agent, input, limits, delegation{}) + if next, ok := m.standbyFor(res, input); ok { + return next.executeRun(ctx, agent, input, limits, delegation{}) + } + return res, err +} + +// standbyFor decides whether to run the whole turn again on another provider. +// +// WHY A RESTART AND NOT A HANDOVER. gateway.canFailOver will not move a +// conversation that has called a tool: the assistant turn echoing that call is +// the provider's own, and a vendor that signs its function calls rejects a +// follow-up carrying somebody else's. But a rate limit lands where the request +// is BIGGEST, which is the second or third call, once the catalogue, the +// retrieved block, the tool results and the whole prior conversation are being +// re-sent. Measured on this deployment: every GatewayFailure recorded had +// already called a tool, so in-place failover covered none of them — +// +// http 429 … Rate limit reached … tokens per minute (TPM): Limit 8000, Used 7183 +// +// Starting over sends no transcript, so nothing provider-specific travels and +// the question is simply asked again somewhere with budget left. It costs the +// work already done, charged to the budget that is NOT exhausted. +// +// THE RULE THAT MAKES IT SAFE: a run that carried a confirmation never +// restarts. Re-running re-runs its tools, and a read twice is two reads while a +// write twice is two shifts assigned. I4 is what makes the test this cheap — +// a write executes ONLY against a resolved token (Registry.gate), so a run with +// no confirmation cannot have written anything, and one with a confirmation is +// refused here without inspecting what it did. +// +// Once. Not a loop over every provider: a question worth asking twice is not +// worth asking five times, and each attempt spends a real budget. The second +// result is returned as it stands, whatever it says. +func (m *ModelExecutor) standbyFor(res *ExecutionResult, input ExecutionInput) (*ModelExecutor, bool) { + if res == nil || res.Termination != TerminationGatewayFailure { + return nil, false + } + // An approved write may already have happened. Nothing below is worth a + // double assignment. + if input.Confirmation != "" { + return nil, false + } + // Only a transient fault moves, on the same line gateway.canFailOver draws: + // a rejected credential or a model this deployment cannot use fails the + // same way everywhere, and asking twice only doubles the bill. + var gwErr *gateway.Error + if !errors.As(res.Error, &gwErr) || !gwErr.Retryable() { + return nil, false + } + sb, ok := m.gw.(gateway.Standby) + if !ok { + return nil, false + } + next, ok := sb.Standby() + if !ok { + return nil, false + } + clone := *m + clone.gw = next + return &clone, true } // executeRun is the loop. `del` is what a SUBAGENT inherits from its parent — diff --git a/go-api/internal/runtime/standby_test.go b/go-api/internal/runtime/standby_test.go new file mode 100644 index 0000000..45097b8 --- /dev/null +++ b/go-api/internal/runtime/standby_test.go @@ -0,0 +1,119 @@ +package runtime + +import ( + "context" + "testing" + + "github.com/krow/krow-backend/go-api/internal/gateway" +) + +// standbyGateway is a gateway with somewhere else to go: `first` answers until +// it is stood down, then `second` does. +type standbyGateway struct { + first gateway.Gateway + second gateway.Gateway +} + +func (s *standbyGateway) Complete(ctx context.Context, req gateway.Request) (*gateway.Response, error) { + return s.first.Complete(ctx, req) +} + +func (s *standbyGateway) Standby() (gateway.Gateway, bool) { + if s.second == nil { + return nil, false + } + return s.second, true +} + +func rateLimited() error { + return &gateway.Error{Code: gateway.CodeRateLimited, Status: 429, Message: "TPM limit 8000"} +} + +// The production case: a rate limit on a run that had already called a tool, +// which in-place failover will not move. +func TestARateLimitedRunIsRetriedOnTheStandbyProvider(t *testing.T) { + busy := &fakeGateway{err: rateLimited()} + spare := &fakeGateway{text: "15 open roles"} + exec := NewModelExecutor(&standbyGateway{first: busy, second: spare}, &MemorySink{}, nil) + + res, err := exec.ExecuteAgent(context.Background(), testAgent(), testInput("how many open positions?")) + if err != nil { + t.Fatalf("the standby should have answered: %v", err) + } + if res.Termination != TerminationCompleted { + t.Fatalf("Termination = %q, want Completed", res.Termination) + } + if res.Output != "15 open roles" { + t.Errorf("Output = %q, want the standby's answer", res.Output) + } + if spare.calls == 0 { + t.Error("the standby provider was never asked") + } +} + +// A run carrying a confirmation has performed an approved write. Re-running it +// re-runs its tools, and a write twice is two shifts assigned. +func TestAConfirmedRunIsNeverRestarted(t *testing.T) { + busy := &fakeGateway{err: rateLimited()} + spare := &fakeGateway{text: "should never be reached"} + exec := NewModelExecutor(&standbyGateway{first: busy, second: spare}, &MemorySink{}, nil) + + in := testInput("assign Maria to the Friday shift") + in.Confirmation = "a-token-a-person-approved" + + res, _ := exec.ExecuteAgent(context.Background(), testAgent(), in) + if res.Termination == TerminationCompleted { + t.Error("a confirmed run was restarted; an approved write could run twice") + } + if spare.calls != 0 { + t.Errorf("the standby was asked %d times; a confirmed run must not be replayed", spare.calls) + } +} + +// A credential or a model id fails the same way everywhere. Asking twice only +// doubles the bill and hides the fault. +func TestATerminalGatewayErrorIsNotRetriedElsewhere(t *testing.T) { + busy := &fakeGateway{err: &gateway.Error{ + Code: gateway.CodeUnauthorized, Status: 401, Message: "bad key", + }} + spare := &fakeGateway{text: "should never be reached"} + exec := NewModelExecutor(&standbyGateway{first: busy, second: spare}, &MemorySink{}, nil) + + res, _ := exec.ExecuteAgent(context.Background(), testAgent(), testInput("anything")) + if res.Termination == TerminationCompleted { + t.Error("a terminal error was retried on another provider") + } + if spare.calls != 0 { + t.Errorf("the standby was asked %d times on a 401", spare.calls) + } +} + +// With one provider there is no standby, and nothing about the single-provider +// path may change. +func TestWithNoStandbyTheFailureStands(t *testing.T) { + busy := &fakeGateway{err: rateLimited()} + exec := NewModelExecutor(busy, &MemorySink{}, nil) + + res, _ := exec.ExecuteAgent(context.Background(), testAgent(), testInput("anything")) + if res.Termination != TerminationGatewayFailure { + t.Errorf("Termination = %q, want GatewayFailure", res.Termination) + } + if busy.calls == 0 { + t.Error("the only provider was never asked") + } +} + +// Once, not until the providers run out. +func TestTheStandbyIsAskedOnlyOnce(t *testing.T) { + busy := &fakeGateway{err: rateLimited()} + alsoBusy := &fakeGateway{err: rateLimited()} + exec := NewModelExecutor(&standbyGateway{first: busy, second: alsoBusy}, &MemorySink{}, nil) + + res, _ := exec.ExecuteAgent(context.Background(), testAgent(), testInput("anything")) + if res.Termination != TerminationGatewayFailure { + t.Errorf("Termination = %q, want GatewayFailure", res.Termination) + } + if alsoBusy.calls == 0 { + t.Error("the standby was never tried") + } +}