diff --git a/go-api/internal/httpserver/gatewaylog_test.go b/go-api/internal/httpserver/gatewaylog_test.go new file mode 100644 index 0000000..2471c13 --- /dev/null +++ b/go-api/internal/httpserver/gatewaylog_test.go @@ -0,0 +1,171 @@ +package httpserver + +// Unit tests for the operator's side of a GatewayFailure. +// +// The user-facing sentence tells the reader whether retrying can work. These +// assert the other half: that the deployment says WHICH fault it was, to the +// only audience that can act on it. Four of the five gateway faults need an +// administrator, and until this line existed a deployment failing every run +// emitted a stream of 200s and nothing else. +// +// Internal rather than httpserver_test because the function under test is +// unexported. Pure: a result goes in and a log record comes out — no server +// wiring, no database. + +import ( + "bytes" + "encoding/json" + "log/slog" + "strings" + "testing" + + "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/gateway" + "github.com/krow/krow-backend/go-api/internal/runtime" +) + +// logging builds a Server that logs into a buffer, and a reader for the records +// it wrote. +func logging(t *testing.T) (*Server, func() []map[string]any) { + t.Helper() + var buf bytes.Buffer + s := &Server{log: slog.New(slog.NewJSONHandler(&buf, nil))} + + return s, func() []map[string]any { + var out []map[string]any + for _, line := range strings.Split(strings.TrimSpace(buf.String()), "\n") { + if line == "" { + continue + } + var rec map[string]any + if err := json.Unmarshal([]byte(line), &rec); err != nil { + t.Fatalf("log line is not JSON: %v", err) + } + out = append(out, rec) + } + return out + } +} + +func gatewayResult(term runtime.Termination, cause error) *runtime.ExecutionResult { + return &runtime.ExecutionResult{ + RunID: "run_1", + AgentID: "activity-agent", + AgentVersion: 3, + Termination: term, + Error: &runtime.RuntimeError{ + Code: "runtime." + strings.ToLower(string(term)), + Message: "internal wording", + Cause: cause, + }, + } +} + +// The code is what distinguishes the retryable fault from the four that need an +// administrator, so it is the field that must survive into the log. +func TestGatewayFailureIsLoggedWithItsCode(t *testing.T) { + for _, tc := range []struct { + name string + cause error + wantCode string + wantStatus float64 + }{ + { + name: "no credential configured", + cause: &gateway.Error{Code: gateway.CodeNotConfigured, Message: "no model credentials"}, + wantCode: gateway.CodeNotConfigured, + }, + { + name: "credential rejected", + cause: &gateway.Error{Code: gateway.CodeUnauthorized, Message: "refused", Status: 401}, + wantCode: gateway.CodeUnauthorized, + wantStatus: 401, + }, + { + name: "rate limited", + cause: &gateway.Error{Code: gateway.CodeRateLimited, Message: "slow down", Status: 429}, + wantCode: gateway.CodeRateLimited, + wantStatus: 429, + }, + } { + t.Run(tc.name, func(t *testing.T) { + s, records := logging(t) + s.logGatewayFailure(authctx.Identity{OrgID: "org_1"}, + gatewayResult(runtime.TerminationGatewayFailure, tc.cause)) + + recs := records() + if len(recs) != 1 { + t.Fatalf("wrote %d log records, want 1: %v", len(recs), recs) + } + rec := recs[0] + + if rec["level"] != "ERROR" { + t.Errorf("level = %v, want ERROR — a deployment that cannot reach its model is an outage", rec["level"]) + } + if got := rec["gateway_code"]; got != tc.wantCode { + t.Errorf("gateway_code = %v, want %q", got, tc.wantCode) + } + if got := rec["gateway_status"]; got != tc.wantStatus { + t.Errorf("gateway_status = %v, want %v", got, tc.wantStatus) + } + + // §10: every line carries these four. + for _, field := range []string{"run_id", "tenant_id", "agent_key", "agent_version"} { + if rec[field] == nil { + t.Errorf("log record has no %s", field) + } + } + // §10 again: no model or document text in the log store. The + // gateway's message can quote the provider's body, so it stays out. + if strings.Contains(strings.ToLower(rec["msg"].(string)), "refused") { + t.Errorf("msg = %q, want no provider text", rec["msg"]) + } + for k, v := range rec { + if str, ok := v.(string); ok && strings.Contains(str, "slow down") { + t.Errorf("field %s leaked the provider message: %q", k, str) + } + } + }) + } +} + +// A cause that is not a gateway error leaves the code visibly empty rather than +// guessed. "Which fault was it" is the whole point of the line, and a wrong +// answer to it is worse than a gap. +func TestGatewayFailureWithoutACauseLogsAnEmptyCode(t *testing.T) { + s, records := logging(t) + s.logGatewayFailure(authctx.Identity{OrgID: "org_1"}, + gatewayResult(runtime.TerminationGatewayFailure, nil)) + + recs := records() + if len(recs) != 1 { + t.Fatalf("wrote %d log records, want 1", len(recs)) + } + if got := recs[0]["gateway_code"]; got != "" { + t.Errorf("gateway_code = %v, want empty", got) + } +} + +// Every other termination is silent here. A run that hit its budget or was +// refused is not a gateway outage, and logging it as one would make the signal +// useless exactly when it is being read. +func TestOnlyGatewayFailureIsLogged(t *testing.T) { + for _, term := range []runtime.Termination{ + runtime.TerminationCompleted, + runtime.TerminationBudgetExceeded, + runtime.TerminationDeadline, + runtime.TerminationConfirmationPending, + runtime.TerminationToolFailure, + runtime.TerminationRefused, + } { + t.Run(string(term), func(t *testing.T) { + s, records := logging(t) + s.logGatewayFailure(authctx.Identity{OrgID: "org_1"}, + gatewayResult(term, &gateway.Error{Code: gateway.CodeRateLimited})) + + if recs := records(); len(recs) != 0 { + t.Errorf("wrote %d log records for %s, want none: %v", len(recs), term, recs) + } + }) + } +} diff --git a/go-api/internal/httpserver/gatewaymessage_test.go b/go-api/internal/httpserver/gatewaymessage_test.go new file mode 100644 index 0000000..ab52e34 --- /dev/null +++ b/go-api/internal/httpserver/gatewaymessage_test.go @@ -0,0 +1,170 @@ +package httpserver + +// Unit tests for the wording a GatewayFailure produces. +// +// Internal rather than httpserver_test because the function under test is the +// mapping itself, and the mapping is unexported. Pure: no server, no database, +// no fixture — a cause goes in and a sentence comes out. +// +// What these assert is one property, and it is the one the old wording broke: +// a reader is told to retry EXACTLY when retrying can work. A rate limit clears +// on its own; a rejected credential, a model id the endpoint does not have, and +// an unconfigured deployment do not, and telling somebody to wait a minute for +// any of those is a loop with no exit. + +import ( + "errors" + "strings" + "testing" + + "github.com/krow/krow-backend/go-api/internal/gateway" + "github.com/krow/krow-backend/go-api/internal/runtime" +) + +// invitesRetry reports whether a sentence tells the reader to try again. +// +// Deliberately looser than an equality check on the whole string: what must +// hold is the ADVICE, not the copy, so rewording a sentence does not fail a +// test that was never about the words. +func invitesRetry(message string) bool { + m := strings.ToLower(message) + return strings.Contains(m, "ask again") || strings.Contains(m, "try again") +} + +func TestGatewayFailureMessageInvitesRetryOnlyWhenRetryingCanWork(t *testing.T) { + cases := []struct { + name string + cause error + retry bool + }{ + { + name: "a rate limit clears on its own", + cause: &gateway.Error{Code: gateway.CodeRateLimited, Status: 429}, + retry: true, + }, + { + name: "a provider 5xx is worth another attempt", + cause: &gateway.Error{Code: gateway.CodeUpstream, Status: 503}, + retry: true, + }, + { + name: "a rejected credential will be rejected again", + cause: &gateway.Error{Code: gateway.CodeUnauthorized, Status: 401}, + retry: false, + }, + { + name: "a model the endpoint does not have stays absent", + cause: &gateway.Error{Code: gateway.CodeInvalidRequest, Status: 404}, + retry: false, + }, + { + name: "an unconfigured deployment cannot answer at all", + cause: &gateway.Error{Code: gateway.CodeNotConfigured}, + retry: false, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + got := gatewayFailureMessage(tc.cause) + if got == "" { + t.Fatal("a failed run must say something") + } + if invitesRetry(got) != tc.retry { + t.Errorf("retry advice = %v, want %v\n message: %q", + invitesRetry(got), tc.retry, got) + } + // Nothing ran, so nothing can have been written. The reassurance is + // the whole reason this termination is not frightening. + if !strings.Contains(got, "Nothing was changed") { + t.Errorf("message must say nothing was changed: %q", got) + } + }) + } +} + +// The three that need a person are the three that used to be indistinguishable +// from load. Each must point at one, or the reader has no idea what to do next. +func TestGatewayFailureMessageNamesAnAdministratorWhenOneIsNeeded(t *testing.T) { + for _, code := range []string{ + gateway.CodeUnauthorized, + gateway.CodeInvalidRequest, + gateway.CodeNotConfigured, + } { + got := gatewayFailureMessage(&gateway.Error{Code: code}) + if !strings.Contains(strings.ToLower(got), "administrator") { + t.Errorf("%s: must send the reader to an administrator: %q", code, got) + } + } +} + +// Vendor names, model ids and HTTP statuses are for the trajectory, not for a +// venue manager — they cannot act on any of them. +func TestGatewayFailureMessageLeaksNoOperatorDetail(t *testing.T) { + for _, code := range []string{ + gateway.CodeRateLimited, + gateway.CodeUpstream, + gateway.CodeUnauthorized, + gateway.CodeInvalidRequest, + gateway.CodeNotConfigured, + } { + got := gatewayFailureMessage(&gateway.Error{ + Code: code, + Status: 429, + Message: "gemini-3.5-flash-lite quota exceeded for project 12345", + }) + for _, leak := range []string{"gemini", "429", "quota", "http", "12345"} { + if strings.Contains(strings.ToLower(got), leak) { + t.Errorf("%s: message carries operator detail %q: %q", code, leak, got) + } + } + } +} + +// A cause that is not a gateway error — lost, wrapped away, or a non-gateway +// failure that reached this termination — still has to produce a sentence. +func TestGatewayFailureMessageFallsBackWithoutAGatewayError(t *testing.T) { + for _, cause := range []error{nil, errors.New("something else entirely")} { + if got := gatewayFailureMessage(cause); got == "" { + t.Errorf("cause %v produced no message", cause) + } + } +} + +// The cause arrives wrapped in a RuntimeError, which is how the surface +// actually receives it. If Unwrap ever stopped reaching the gateway error, +// every failure would silently fall back to "usually it is busy" — the exact +// bug this change exists to fix, reintroduced without a compile error. +func TestGatewayFailureMessageReadsThroughARuntimeError(t *testing.T) { + wrapped := &runtime.RuntimeError{ + Code: "runtime.gatewayfailure", + Cause: &gateway.Error{Code: gateway.CodeUnauthorized, Status: 401}, + } + + got := gatewayFailureMessage(wrapped) + if invitesRetry(got) { + t.Errorf("a wrapped credential failure must not invite a retry: %q", got) + } +} + +// The other terminations are unchanged by the new parameter: they ignore the +// cause, so passing one must not alter a single word. +func TestTerminationMessageIgnoresTheCauseElsewhere(t *testing.T) { + cause := &gateway.Error{Code: gateway.CodeUnauthorized} + + for _, term := range []runtime.Termination{ + runtime.TerminationBudgetExceeded, + runtime.TerminationDeadline, + runtime.TerminationConfirmationPending, + runtime.TerminationToolFailure, + runtime.TerminationRefused, + } { + if terminationMessage(term, nil) != terminationMessage(term, cause) { + t.Errorf("%s: wording changed with the cause", term) + } + } + + if terminationMessage(runtime.TerminationCompleted, nil) != "" { + t.Error("a completed run has nothing to say") + } +} diff --git a/go-api/internal/httpserver/runs.go b/go-api/internal/httpserver/runs.go index 994fbc6..004a6f9 100644 --- a/go-api/internal/httpserver/runs.go +++ b/go-api/internal/httpserver/runs.go @@ -9,6 +9,7 @@ import ( "github.com/krow/krow-backend/go-api/internal/authctx" "github.com/krow/krow-backend/go-api/internal/domain" + "github.com/krow/krow-backend/go-api/internal/gateway" "github.com/krow/krow-backend/go-api/internal/runtime" "github.com/krow/krow-backend/go-api/internal/tools" ) @@ -163,6 +164,7 @@ func (s *Server) handleAgentRun(w http.ResponseWriter, r *http.Request) { return } s.logUnsaved(ident, res) + s.logGatewayFailure(ident, res) writeJSON(w, http.StatusOK, buildRunResponse(res)) } @@ -180,6 +182,43 @@ func (s *Server) logUnsaved(ident authctx.Identity, res *runtime.ExecutionResult } } +// logGatewayFailure is the operator's record of the model provider not +// answering, and it exists because the reader of the message cannot report what +// the message does not say. +// +// Four of the five gateway faults need an administrator and will fail +// identically on every retry — a rejected credential, a model id this +// deployment cannot use, no key at all, and a 4xx from the endpoint. The person +// in the chat panel is told, correctly, that retrying will not help; but +// nothing until now told the side that CAN fix it. A deployment failing every +// run produced a stream of 200s and no error line, so the only record of which +// fault it was lived in a trajectory somebody had to know to go and read. +// +// Logged at Error because that is what it is: on a rate limit it is a capacity +// decision worth seeing, and on the other four it is an outage. It carries the +// gateway's code and status and NOT the message — §10 keeps model and document +// text out of the log store, and the code is the part that is actionable +// anyway. +func (s *Server) logGatewayFailure(ident authctx.Identity, res *runtime.ExecutionResult) { + if res.Termination != runtime.TerminationGatewayFailure { + return + } + + // Empty rather than invented when the cause did not survive: "which fault + // was it" is the whole point of this line, and a guessed answer to it is + // worse than a visible gap. + code, status := "", 0 + var gwErr *gateway.Error + if errors.As(res.Error, &gwErr) { + code, status = gwErr.Code, gwErr.Status + } + + s.log.Error("gateway failure", + "run_id", res.RunID, "tenant_id", ident.OrgID, + "agent_key", res.AgentID, "agent_version", res.AgentVersion, + "gateway_code", code, "gateway_status", status) +} + // buildRunResponse turns a runtime result into the client's shape. // // Every termination answers 200. That looks wrong at first and is not: the @@ -206,7 +245,7 @@ func buildRunResponse(res *runtime.ExecutionResult) runResponse { }, } if res.Termination != runtime.TerminationCompleted { - out.Message = terminationMessage(res.Termination) + out.Message = terminationMessage(res.Termination, res.Error) } return out } @@ -222,7 +261,9 @@ func buildRunResponse(res *runtime.ExecutionResult) runResponse { // Every one of the seven is spelled out. A default that said "something went // wrong" would be the place where a Refused run and a ToolFailure became // indistinguishable to the person best placed to tell us which it was. -func terminationMessage(t runtime.Termination) string { +// `cause` is the run's own error, carried so GatewayFailure can say which +// gateway failure it was. Every other termination ignores it. +func terminationMessage(t runtime.Termination, cause error) string { switch t { case runtime.TerminationCompleted: return "" @@ -239,16 +280,70 @@ func terminationMessage(t runtime.Termination) string { case runtime.TerminationRefused: return "The agent declined to answer this one." case runtime.TerminationGatewayFailure: - // The one termination where "try again" is honest advice: the - // dominant cause is a rate limit that clears within a minute, and - // nothing about the question itself was the problem. - return "The model behind this agent did not answer — usually it is busy. " + - "Wait a minute and ask again. Nothing was changed." + return gatewayFailureMessage(cause) default: return "The agent did not finish." } } +// gatewayFailureMessage tells a GatewayFailure apart from the four others it +// used to be indistinguishable from. +// +// GatewayFailure is everything the gateway can raise except a refusal and a +// timeout, which have terminations of their own. That is five different faults, +// and the one sentence they all produced was "usually it is busy — wait a +// minute and ask again". +// +// For a rate limit that is true and useful. For the other three it is advice +// that CANNOT work: a credential the provider rejected, a model id that is not +// on the configured endpoint, and no key at all will each fail identically on +// every retry, forever. Telling somebody to wait a minute for a misconfigured +// deployment sends them round a loop with no exit, and hides an operator +// problem behind what reads as a transient one — the reader retries instead of +// reporting it, so nobody with access to fix it ever hears. +// +// So each says what it is, and only the two that clear on their own invite a +// retry. The wording stays free of vendor names and status codes: the person +// reading it cannot act on "429 from the model endpoint", and the code is in +// the trajectory for the person who can. +func gatewayFailureMessage(cause error) string { + var gwErr *gateway.Error + if !errors.As(cause, &gwErr) { + // No gateway error to read — either the cause was lost or something + // non-gateway reached this termination. The old sentence is still the + // best guess, so it is what an unknown falls back to. + return "The model behind this agent did not answer — usually it is busy. " + + "Wait a minute and ask again. Nothing was changed." + } + + switch gwErr.Code { + case gateway.CodeRateLimited: + return "The model behind this agent is busy right now. " + + "Wait a minute and ask again. Nothing was changed." + case gateway.CodeUpstream: + if gwErr.Status >= 500 { + return "The model behind this agent is having trouble. " + + "Try again in a few minutes. Nothing was changed." + } + return "The model behind this agent could not be reached, and retrying is " + + "unlikely to help. This needs an administrator. Nothing was changed." + case gateway.CodeUnauthorized: + return "This deployment's model credentials were rejected, so the agent " + + "cannot answer. Retrying will not help — it needs an administrator. " + + "Nothing was changed." + case gateway.CodeInvalidRequest: + return "The agent is pointed at a model this deployment cannot use, so it " + + "cannot answer. Retrying will not help — it needs an administrator. " + + "Nothing was changed." + case gateway.CodeNotConfigured: + return "No model is configured for this deployment, so the agent cannot " + + "answer. It needs an administrator. Nothing was changed." + default: + return "The model behind this agent did not answer — usually it is busy. " + + "Wait a minute and ask again. Nothing was changed." + } +} + // runLoadError maps a pre-run failure onto the API's error vocabulary. // // These are the errors from LoadExecutableAgent, raised before any run began — @@ -351,6 +446,7 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au return } s.logUnsaved(ident, res) + s.logGatewayFailure(ident, res) writeJSON(w, http.StatusOK, buildRunResponse(res)) return } @@ -399,6 +495,7 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au } s.logUnsaved(ident, res) + s.logGatewayFailure(ident, res) send(map[string]any{"run": buildRunResponse(res)}) fmt.Fprint(w, "data: [DONE]\n\n") flusher.Flush() diff --git a/go-api/internal/runtime/budget.go b/go-api/internal/runtime/budget.go index 410d705..f79d854 100644 --- a/go-api/internal/runtime/budget.go +++ b/go-api/internal/runtime/budget.go @@ -56,6 +56,15 @@ type Limits struct { MaxToolCalls int MaxTokens int64 Deadline time.Duration + + // MaxOutputTokens caps a SINGLE model call; MaxTokens caps the whole run. + // + // Without it the run budget was the only ceiling on any one response, so a + // balanced run could spend its 120k as eight 16k generations — and + // generation time is the wall clock a person waits through. Latency is why + // this exists; cost is a side effect. Zero means the gateway's configured + // default (`MODEL_MAX_OUTPUT_TOKENS`). + MaxOutputTokens int64 } // LimitsForTier is what a run gets when its spec declares no limits of its own. @@ -70,14 +79,23 @@ type Limits struct { // expensive one: a fast run gets a third of a deep run's steps and a sixth of // its deadline, so a misrouted spec shows up as a truncated answer rather than // as a bill. +// +// MaxOutputTokens follows the same shape. It is sized for the longest answer a +// tier should ever give in one turn, not for the run: a tool-call step spends a +// few hundred tokens on arguments, and a chat answer past ~3k tokens is one +// nobody reads. A cap the model is not told about truncates rather than winding +// down, so these are set above any legitimate answer and not near it. func LimitsForTier(tier string) Limits { switch tier { case "fast": - return Limits{MaxSteps: 3, MaxToolCalls: 4, MaxTokens: 40_000, Deadline: 20 * time.Second} + return Limits{MaxSteps: 3, MaxToolCalls: 4, MaxTokens: 40_000, Deadline: 20 * time.Second, + MaxOutputTokens: 1_500} case "deep": - return Limits{MaxSteps: 12, MaxToolCalls: 20, MaxTokens: 300_000, Deadline: 120 * time.Second} + return Limits{MaxSteps: 12, MaxToolCalls: 20, MaxTokens: 300_000, Deadline: 120 * time.Second, + MaxOutputTokens: 4_000} default: // balanced, and anything unrecognised — ParseTier has already normalised it - return Limits{MaxSteps: 8, MaxToolCalls: 12, MaxTokens: 120_000, Deadline: 60 * time.Second} + return Limits{MaxSteps: 8, MaxToolCalls: 12, MaxTokens: 120_000, Deadline: 60 * time.Second, + MaxOutputTokens: 3_000} } } diff --git a/go-api/internal/runtime/loop.go b/go-api/internal/runtime/loop.go index 7662ab7..ef90ba9 100644 --- a/go-api/internal/runtime/loop.go +++ b/go-api/internal/runtime/loop.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "strings" + "time" "github.com/krow/krow-backend/go-api/internal/gateway" "github.com/krow/krow-backend/go-api/internal/knowledge" @@ -24,6 +25,17 @@ import ( // // The loop runs until the model stops asking for tools, or until a bound is // reached. Every exit is one of the six terminations. + +// trajectoryPersistTimeout bounds the trajectory write that happens after a run +// has produced its answer. Generous, because losing the record of a run is +// worse than a slow one, but finite: the caller is still waiting on this, so an +// unreachable database must cost seconds and not the request's whole write +// timeout. See finish. +// +// A var rather than a const only so a test can assert the bound holds without +// spending the bound. Nothing outside this package sets it. +var trajectoryPersistTimeout = 5 * time.Second + type ModelExecutor struct { // subagents resolves a spec's `subagents:` into runnable agents. nil means // delegation is off and a spec that declares subagents runs alone — see @@ -180,21 +192,51 @@ func (m *ModelExecutor) executeRun( } rec.Message("user", question) - // The tools this agent may use. Unknown names are recorded and dropped - // rather than failing the run: §3 says an unknown tool fails validation at - // *publish*, so one reaching run time means a tool was withdrawn under a - // live spec — degrading is better than an outage, provided someone is told. - toolDefs, unknown := m.toolsFor(agent) - for _, name := range unknown { - rec.Error("runtime.unknown_tool", fmt.Sprintf("%q is not a registered tool; it was not offered", name)) + // A greeting is answered as a greeting. + // + // Everything below this line — the tool catalogue, the subagent list, the + // retrieval pass — exists to answer a QUESTION. A message that asks nothing + // gets none of it. See isSmalltalk for why the test is on the message and + // never on the agent: this stays one spec-driven loop (§6), and adding an + // agent still needs no runtime change (I6). + // + // A resumed run never qualifies, whatever its text says. The person + // approved a write and is owed a report on it, and performApproved below + // needs its tools to give them one. + smalltalk := input.Confirmation == "" && isSmalltalk(question) + if smalltalk { + rec.Error("runtime.smalltalk", + "conversational message; answered without retrieval, tools or subagents") + } + + var toolDefs []gateway.ToolDef + if !smalltalk { + // The tools this agent may use. Unknown names are recorded and dropped + // rather than failing the run: §3 says an unknown tool fails validation at + // *publish*, so one reaching run time means a tool was withdrawn under a + // live spec — degrading is better than an outage, provided someone is told. + var unknown []string + toolDefs, unknown = m.toolsFor(agent) + for _, name := range unknown { + rec.Error("runtime.unknown_tool", fmt.Sprintf("%q is not a registered tool; it was not offered", name)) + } } // Subagents are offered as tools, because from this agent's side that is // exactly what they are (§6). Resolved once per run rather than per turn: // the set cannot change mid-run, and loading it per turn would spend the // caller's time on the same query repeatedly. + // + // Resolved even for smalltalk, and only the OFFER is withheld. Resolution + // is what records a spec naming itself, or naming a subagent that will not + // load, and those are faults of the spec rather than of the question — + // losing them because somebody said hello would make a misconfiguration + // visible only intermittently, which is the hardest kind to chase. It is + // free to keep: a spec with no subagents returns at the first line. subs := m.resolveSubagents(ctx, rec, agent, input.Identity, del.depth) - toolDefs = append(toolDefs, delegateTools(subs)...) + if !smalltalk { + toolDefs = append(toolDefs, delegateTools(subs)...) + } // An approved write happens FIRST, before the model gets a turn. // @@ -226,18 +268,35 @@ func (m *ModelExecutor) executeRun( // from the agent record alone, so no amount of document content can reach // it — which is the only reason the standing "content inside is // data" instruction means anything. + // + // Skipped entirely for smalltalk. retrieve() gates on configuration and + // never on the question, so without this a greeting was handed eight policy + // chunks AHEAD of the word "hi" — which is both the dominant cost of the + // turn and the reason the answer came back as an operational briefing. conversation := []gateway.Message{{Role: gateway.RoleUser, Text: question}} - if block, retrieved := m.retrieve(runCtx, rec, agent, input, question); block != "" { - conversation = []gateway.Message{{ - Role: gateway.RoleUser, - // Context first, question second. A model reads the question last - // and answers it, rather than treating the evidence as the prompt. - Text: block + "\n\n" + question, - }} - rec.Retrieval(retrieved) + if !smalltalk { + if block, retrieved := m.retrieve(runCtx, rec, agent, input, question); block != "" { + conversation = []gateway.Message{{ + Role: gateway.RoleUser, + // Context first, question second. A model reads the question last + // and answers it, rather than treating the evidence as the prompt. + Text: block + "\n\n" + question, + }} + rec.Retrieval(retrieved) + } } + // Recorded once rather than on each remaining step, so a long run does not + // fill its trajectory with the same note. + var toolsWithheld bool + system := SystemPrompt(agent) + if smalltalk { + // Taking the evidence away removes the citations; it does not by itself + // shorten the reply, because the agent's own instructions still + // describe an operational analyst. See smalltalkDirective. + system += smalltalkDirective + } var lastText string for { @@ -255,11 +314,43 @@ func (m *ModelExecutor) executeRun( // halves go through StreamComplete, so the loop has one call site and // no branch on transport — a run behaves identically whether its text // arrived in one piece or a hundred. + // The catalogue has to be resent on every call — the wire protocol has + // no way to refer back to one already sent — but a catalogue the model + // is no longer ALLOWED to use is pure waste. Once the tool-call budget + // is spent, every definition describes a call that would be refused, + // and the step it is being sent on is the synthesis turn that just + // needs to write the answer up. + // + // Measured on this deployment's control-center agent: seven tools, + // ~1.2k tokens, resent on the final call of every tool-using run. + stepTools := toolDefs + if len(stepTools) > 0 && budget.Snapshot().ToolCallsLeft <= 0 { + stepTools = nil + if !toolsWithheld { + toolsWithheld = true + rec.Error("runtime.tools_withheld", + "tool-call budget spent; catalogue not resent on the remaining steps") + } + } + + // Smalltalk is capped far below the tier's ceiling. The directive in + // the system prompt is what actually shortens the reply; this only + // bounds the bill for a model that ignores it. + maxOut := budget.Limits().MaxOutputTokens + if smalltalk && (maxOut <= 0 || maxOut > smalltalkMaxOutputTokens) { + maxOut = smalltalkMaxOutputTokens + } + resp, err := gateway.StreamComplete(runCtx, m.gw, gateway.Request{ Tier: tier, System: system, Messages: conversation, - Tools: toolDefs, + Tools: stepTools, + // Per STEP, and taken from the budget rather than from the tier — + // so a delegated run inherits the parent's ceiling along with the + // parent's budget instead of reading its own tier and quietly + // buying a longer answer than the parent was allowed. + MaxOutputTokens: maxOut, }, input.OnDelta) // Charged whatever happened. A refused or failed call was still billed, @@ -578,9 +669,13 @@ func terminationFor(err error) Termination { // finish closes the trajectory, persists it, and builds the caller's result. // -// Persistence uses the *caller's* context, not the run's: the run context is +// Persistence deliberately outlives the run's context: that context is // cancelled at the deadline, and a run that ended by running out of time is -// exactly the one whose record is most worth keeping. +// exactly the one whose record is most worth keeping. It does NOT outlive the +// caller's patience — trajectoryPersistTimeout bounds the whole write, because +// a run that answered inside its deadline and then sat in the sink for a minute +// is, to the person waiting, a slow run. I3 bounds the run; this bounds its +// tail. func (m *ModelExecutor) finish( ctx context.Context, rec *Recorder, @@ -594,7 +689,16 @@ func (m *ModelExecutor) finish( if cause != nil { var gwErr *gateway.Error if errors.As(cause, &gwErr) { - rec.Error(gwErr.Code, gwErr.Message) + // The status rides in the message because `entries` has no column + // for it and E5 forbids applying a migration from here. It matters: + // `gateway.upstream` alone cannot tell a provider shedding load + // (5xx, clears by itself) from an endpoint rejecting the request + // (4xx, needs an administrator), and those are opposite actions. + msg := gwErr.Message + if gwErr.Status > 0 { + msg = fmt.Sprintf("http %d: %s", gwErr.Status, msg) + } + rec.Error(gwErr.Code, msg) } else { rec.Error("runtime.failed", cause.Error()) } @@ -602,11 +706,19 @@ func (m *ModelExecutor) finish( rec.Budget(budget.Snapshot()) traj := rec.Finish(term) + // WithoutCancel so a deadline-terminated run still records itself; the + // timeout so it cannot record itself forever. One budget covers the parent + // and every child, since writing the tree is one logical act and a + // per-trajectory timeout would multiply by the number of subagents. + persistCtx, cancelPersist := context.WithTimeout( + context.WithoutCancel(ctx), trajectoryPersistTimeout) + defer cancelPersist() + // A sink that fails must not fail the run — the answer was already // produced. It is recorded in the trajectory we could not save, which is // the best available place for it. var unsaved []string - if err := m.sink.Save(ctx, traj); err != nil { + if err := m.sink.Save(persistCtx, traj); err != nil { rec.Error("runtime.trajectory_unsaved", err.Error()) unsaved = append(unsaved, traj.RunID+": "+err.Error()) } @@ -616,7 +728,7 @@ func (m *ModelExecutor) finish( // own descendants already ordered behind it, so one pass here writes a // whole tree parent-first. for _, child := range rec.Children() { - if err := m.sink.Save(ctx, child); err != nil { + if err := m.sink.Save(persistCtx, child); err != nil { rec.Error("runtime.subrun_unsaved", fmt.Sprintf("%s: %s", child.RunID, err.Error())) unsaved = append(unsaved, child.RunID+": "+err.Error()) diff --git a/go-api/internal/runtime/output_cap_test.go b/go-api/internal/runtime/output_cap_test.go new file mode 100644 index 0000000..ca3a3fd --- /dev/null +++ b/go-api/internal/runtime/output_cap_test.go @@ -0,0 +1,139 @@ +package runtime + +import ( + "context" + "testing" + "time" +) + +// These use a real question, never a greeting: a conversational message takes +// the smalltalk path and is capped well below its tier (see smalltalk.go), so +// "hi" here would assert the greeting cap while appearing to assert the tier's. +// +// A per-call cap is the difference between "the run may spend 120k tokens" and +// "any one answer may be 16k tokens long, eight times over". These assert the +// cap actually reaches the gateway, because it is inert until it does. + +func TestOutputCapReachesGatewayPerTier(t *testing.T) { + for _, tc := range []struct { + reasoning string + want int64 + }{ + {"fast", 1_500}, + {"balanced", 3_000}, + {"deep", 4_000}, + {"nonsense-tier", 3_000}, // normalised to balanced, still capped + } { + t.Run(tc.reasoning, func(t *testing.T) { + gw := &fakeGateway{text: "done"} + exec := NewModelExecutor(gw, &MemorySink{}, nil) + + agent := testAgent() + agent.Reasoning = tc.reasoning + + if _, err := exec.ExecuteAgent(context.Background(), agent, testInput("what happened today?")); err != nil { + t.Fatal(err) + } + if got := gw.lastReq.MaxOutputTokens; got != tc.want { + t.Errorf("MaxOutputTokens = %d, want %d", got, tc.want) + } + }) + } +} + +// The cap comes from the budget, not from the agent's own tier. That is what +// makes a subagent inherit the parent's ceiling along with the parent's budget +// instead of reading its own tier and buying a longer answer. +func TestOutputCapComesFromTheBudgetNotTheTier(t *testing.T) { + gw := &fakeGateway{text: "done"} + exec := NewModelExecutor(gw, &MemorySink{}, nil) + + agent := testAgent() + agent.Reasoning = "deep" // would be 4_000 if the tier decided + + limits := LimitsForTier("deep") + limits.MaxOutputTokens = 777 + + if _, err := exec.executeWithLimits(context.Background(), agent, testInput("what happened today?"), limits); err != nil { + t.Fatal(err) + } + if got := gw.lastReq.MaxOutputTokens; got != 777 { + t.Errorf("MaxOutputTokens = %d, want 777 from the supplied limits", got) + } +} + +// blockingSink is a database that has stopped answering. It returns only when +// its context ends, and records what ended it. +type blockingSink struct { + ctxErr error + elapsed time.Duration +} + +func (b *blockingSink) Save(ctx context.Context, _ *Trajectory) error { + start := time.Now() + <-ctx.Done() + b.elapsed = time.Since(start) + b.ctxErr = ctx.Err() + return ctx.Err() +} + +// A run that answered must not then wait on the sink indefinitely: to the person +// watching, a slow write is a slow run. I3 bounds the run; this bounds its tail. +func TestTrajectoryPersistenceIsBounded(t *testing.T) { + restore := trajectoryPersistTimeout + trajectoryPersistTimeout = 50 * time.Millisecond + t.Cleanup(func() { trajectoryPersistTimeout = restore }) + + sink := &blockingSink{} + exec := NewModelExecutor(&fakeGateway{text: "answered"}, sink, nil) + + start := time.Now() + res, err := exec.ExecuteAgent(context.Background(), testAgent(), testInput("hi")) + total := time.Since(start) + + // The answer survives the sink failing — §6 keeps the two separate. + if err != nil { + t.Fatalf("a failed save must not fail the run: %v", err) + } + if res.Output != "answered" { + t.Errorf("Output = %q, want the model's answer", res.Output) + } + if len(res.Unsaved) != 1 { + t.Errorf("Unsaved = %v, want the one trajectory that could not be written", res.Unsaved) + } + + if sink.ctxErr != context.DeadlineExceeded { + t.Errorf("sink ctx ended with %v, want DeadlineExceeded — the write was not bounded", sink.ctxErr) + } + // Generous slack: the assertion is "bounded", not "fast". + if total > 2*time.Second { + t.Errorf("run took %v with a hung sink, want the persist timeout to cut it", total) + } +} + +// The bound must not become a cancellation: a run terminated by its own deadline +// is the one whose record matters most, so the write starts from a live context +// even when the caller's is already dead. +func TestTrajectoryPersistenceOutlivesACancelledCaller(t *testing.T) { + sink := &MemorySink{} + exec := NewModelExecutor(&fakeGateway{text: "answered", delay: 50 * time.Millisecond}, sink, nil) + + ctx, cancel := context.WithCancel(context.Background()) + // Cancelled while the model call is in flight: the run ends unhappily and + // must still be recorded. + go func() { + time.Sleep(10 * time.Millisecond) + cancel() + }() + + res, _ := exec.ExecuteAgent(ctx, testAgent(), testInput("hi")) + if res == nil { + t.Fatal("no result") + } + if len(res.Unsaved) != 0 { + t.Errorf("Unsaved = %v, want the trajectory written despite the cancelled caller", res.Unsaved) + } + if got := len(sink.Runs); got != 1 { + t.Errorf("sink holds %d trajectories, want 1", got) + } +} diff --git a/go-api/internal/runtime/smalltalk.go b/go-api/internal/runtime/smalltalk.go new file mode 100644 index 0000000..1052526 --- /dev/null +++ b/go-api/internal/runtime/smalltalk.go @@ -0,0 +1,126 @@ +package runtime + +import "strings" + +// isSmalltalk reports a message that cannot be answered any better by looking +// something up. +// +// "Hi" used to cost a full operational turn. Retrieval is gated on +// configuration and never on the question (see retrieve), so a greeting arrived +// at the model wrapped in eight policy chunks, with the whole tool catalogue +// attached and the evidence placed BEFORE the question. A model handed that +// reasonably concludes it was asked for an operational briefing, and answers +// with one — screening backlog, uncovered shifts, citations and all. Measured +// against production: 6,174 tokens over two model calls, for the word "hi". +// +// The cost is the smaller half. The real damage is that the product appears not +// to understand being greeted, which is the first thing anybody tries. +// +// This is deliberately NOT a per-agent rule, and not an `if agent_key == ...` +// — §13 lists that as the anti-pattern it is. It is a property of the MESSAGE, +// applied identically to every spec, so adding an agent still requires no +// runtime change (I6). +// +// Conservative by construction: the normalised message must match a phrase in +// the set EXACTLY. Nothing substring-matches, so "hi, which shifts are +// uncovered?" is an operational question and keeps its tools and its evidence. +// A false negative costs a few thousand tokens; a false positive answers a real +// question with a greeting, so the set only holds phrases that carry no request +// at all. +func isSmalltalk(q string) bool { + n := normaliseSmalltalk(q) + if n == "" { + return false + } + _, ok := smalltalkPhrases[n] + return ok +} + +// normaliseSmalltalk reduces a message to lowercase letters and single spaces. +// +// Punctuation and emoji are dropped rather than enumerated, so "Hi!", "hi :)" +// and "HI 👋" all arrive as "hi" without the set needing a row for each. Digits +// are NOT letters and so are dropped too, which is harmless here: no phrase in +// the set contains one, and a message that does — "shift 12?" — fails the exact +// match either way. +func normaliseSmalltalk(q string) string { + var b strings.Builder + b.Grow(len(q)) + space := false + for _, r := range strings.ToLower(strings.TrimSpace(q)) { + switch { + case r >= 'a' && r <= 'z': + if space && b.Len() > 0 { + b.WriteByte(' ') + } + space = false + b.WriteRune(r) + case r == '\'' || r == '’': + // Dropped outright rather than treated as a separator, so "how's" + // stays one word. Both the ASCII quote and the curly one a phone + // keyboard substitutes — the same character to whoever typed it, + // and not to the first version of this function. + default: + // Any other run of non-letters is one separator, so "thank-you" + // and "thank you" normalise alike. + space = true + } + } + return b.String() +} + +// smalltalkPhrases is the whole rule, as data. +// +// Greetings, thanks and farewells only. Each is a complete message that asks +// for nothing, which is what makes skipping retrieval and tools safe rather +// than merely cheap. Acknowledgements like "ok" and "cool" are deliberately +// absent: they are plausible smalltalk but also plausible answers to a +// question the agent just asked, and the cost of being wrong is higher than +// the tokens being saved. +var smalltalkPhrases = map[string]struct{}{ + "hi": {}, "hii": {}, "hiya": {}, "hello": {}, "helo": {}, "hey": {}, + "yo": {}, "howdy": {}, "greetings": {}, "hi there": {}, + "hello there": {}, "hey there": {}, "hi owliver": {}, + "hello owliver": {}, "hey owliver": {}, + + "good morning": {}, "good afternoon": {}, "good evening": {}, + "good day": {}, "morning": {}, "afternoon": {}, "evening": {}, + "gm": {}, "ge": {}, + + "how are you": {}, "how are you doing": {}, "hows it going": {}, + "how is it going": {}, "you there": {}, "are you there": {}, + + "thanks": {}, "thank you": {}, "thanks a lot": {}, + "thank you very much": {}, "thanks very much": {}, "many thanks": {}, + "ty": {}, "cheers": {}, "thank u": {}, + + "bye": {}, "goodbye": {}, "good bye": {}, "see you": {}, + "see ya": {}, "good night": {}, "goodnight": {}, "later": {}, +} + +// smalltalkDirective is appended to the system prompt for a smalltalk turn. +// +// Needed because the agent's own instructions describe an operational analyst, +// and an operational analyst greeted with "hi" and given no tools will still +// reach for the longest answer it can justify. Removing the evidence removes +// the citations; it does not by itself shorten the reply. +// +// Appended to the SYSTEM prompt rather than wrapped around the user's message: +// it is a standing instruction from the platform, not something the person +// said, and putting words in their mouth is how a transcript stops matching +// what was typed. I7 is untouched — this is the runtime's own text, not +// retrieved content, and nothing retrieved can reach here because retrieval did +// not run. +const smalltalkDirective = "\n\nThe person has greeted you or said something " + + "conversational. Reply in one or two short sentences: greet them back and " + + "offer to help. Do not summarise data, do not list findings or next steps, " + + "and do not cite sources — you have not looked anything up." + +// smalltalkMaxOutputTokens caps a greeting's reply. +// +// A ceiling the model is not told about truncates mid-sentence rather than +// winding down, so this sits well above any sane greeting (a sentence or two is +// well under 100 tokens) and acts only as a backstop for a model that ignores +// the directive above. The directive does the shortening; this bounds the bill +// when it does not. +const smalltalkMaxOutputTokens = 256 diff --git a/go-api/internal/runtime/smalltalk_test.go b/go-api/internal/runtime/smalltalk_test.go new file mode 100644 index 0000000..73f7e70 --- /dev/null +++ b/go-api/internal/runtime/smalltalk_test.go @@ -0,0 +1,178 @@ +package runtime + +import ( + "context" + "encoding/json" + "strings" + "testing" + "time" + + "github.com/krow/krow-backend/go-api/internal/gateway" + "github.com/krow/krow-backend/go-api/internal/tools" +) + +// TestIsSmalltalkGreetings covers what the fix is for: the messages that were +// costing a full operational turn. +func TestIsSmalltalkGreetings(t *testing.T) { + for _, q := range []string{ + "hi", "Hi", "HI", "hi!", "hi.", " hi ", "hi :)", "hi 👋", + "hello", "Hello!", "hey", "Hey there", "hi there", + "good morning", "Good Morning!", "good evening", + "thanks", "Thank you", "thank-you", "thank you", + "bye", "Goodbye", "good night", + "how are you", "How's it going?", + } { + if !isSmalltalk(q) { + t.Errorf("isSmalltalk(%q) = false, want true", q) + } + } +} + +// TestIsSmalltalkRealQuestions is the half that matters for correctness. A +// false positive answers a real operational question with a greeting, so +// anything carrying a request must fall through — including the ones that +// merely START with a greeting. +func TestIsSmalltalkRealQuestions(t *testing.T) { + for _, q := range []string{ + "", " ", + "how many open positions?", + "what happened today", + "hi, how many open positions?", + "hello there, which shifts are uncovered?", + "hey can you check the screening backlog", + "thanks — now show me the overtime report", + "good morning, what needs attention right now?", + "say hi to the new starters", + "how are you handling the uncovered shifts", + "bye week coverage", + } { + if isSmalltalk(q) { + t.Errorf("isSmalltalk(%q) = true, want false", q) + } + } +} + +func TestNormaliseSmalltalk(t *testing.T) { + cases := map[string]string{ + "Hi!": "hi", + " HELLO ": "hello", + "thank-you": "thank you", + "thank you": "thank you", + "How's it go?": "hows it go", + "👋": "", + "shift 12": "shift", + } + for in, want := range cases { + if got := normaliseSmalltalk(in); got != want { + t.Errorf("normaliseSmalltalk(%q) = %q, want %q", in, got, want) + } + } +} + +// TestSmalltalkSendsNoToolsAndSkipsRetrieval is the fix as the user meets it: +// "hi" reaches the model as "hi", with nothing attached. +func TestSmalltalkSendsNoToolsAndSkipsRetrieval(t *testing.T) { + var toolCalls int + reg := tools.NewRegistry() + reg.MustRegister(countingTool("activity_breakdown", &toolCalls)) + + ret := &scriptedRetriever{results: onePassage("Staff must arrive fifteen minutes early.")} + gw := &scriptedGateway{} + + agent := knowledgeAgent() + agent.Tools = []string{"activity_breakdown"} + + exec := NewModelExecutor(gw, &MemorySink{}, reg).WithRetriever(ret) + res, err := exec.ExecuteAgent(context.Background(), agent, testInput("Hi")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Termination != TerminationCompleted { + t.Fatalf("Termination = %q, want Completed", res.Termination) + } + + if ret.calls != 0 { + t.Errorf("retriever was called %d times for a greeting; want 0", ret.calls) + } + if len(gw.seen) != 1 { + t.Fatalf("a greeting took %d model calls, want 1", len(gw.seen)) + } + req := gw.seen[0] + if len(req.Tools) != 0 { + t.Errorf("greeting carried %d tool definitions, want 0", len(req.Tools)) + } + if req.Messages[0].Text != "Hi" { + t.Errorf("model saw %q, want the bare greeting", req.Messages[0].Text) + } + if !strings.Contains(req.System, "greeted you") { + t.Error("the smalltalk directive did not reach the system prompt") + } + if req.MaxOutputTokens != smalltalkMaxOutputTokens { + t.Errorf("MaxOutputTokens = %d, want the smalltalk cap %d", + req.MaxOutputTokens, smalltalkMaxOutputTokens) + } +} + +// TestOperationalQuestionKeepsToolsAndRetrieval is the guard on the fix above. +// The cheap path must not swallow a question that needs evidence. +func TestOperationalQuestionKeepsToolsAndRetrieval(t *testing.T) { + var toolCalls int + reg := tools.NewRegistry() + reg.MustRegister(countingTool("activity_breakdown", &toolCalls)) + + ret := &scriptedRetriever{results: onePassage("Staff must arrive fifteen minutes early.")} + gw := &scriptedGateway{} + + agent := knowledgeAgent() + agent.Tools = []string{"activity_breakdown"} + + exec := NewModelExecutor(gw, &MemorySink{}, reg).WithRetriever(ret) + if _, err := exec.ExecuteAgent( + context.Background(), agent, testInput("hi, how many open positions?"), + ); err != nil { + t.Fatalf("unexpected error: %v", err) + } + + if ret.calls != 1 { + t.Errorf("retriever called %d times for a real question, want 1", ret.calls) + } + if len(gw.seen[0].Tools) == 0 { + t.Error("a real question was sent with no tools") + } +} + +// TestCatalogueWithheldOnceToolBudgetIsSpent covers the other half of the cost +// work: the synthesis turn at the end of a tool-using run is sent without a +// catalogue the model is no longer permitted to use. +func TestCatalogueWithheldOnceToolBudgetIsSpent(t *testing.T) { + var toolCalls int + reg := tools.NewRegistry() + reg.MustRegister(countingTool("activity_breakdown", &toolCalls)) + + gw := &scriptedGateway{steps: []*gateway.Response{{ + ToolCalls: []gateway.ToolCall{{ID: "call_1", Name: "activity_breakdown", Input: json.RawMessage(`{}`)}}, + StopReason: "tool_use", Model: "fake-model", + }}} + + agent := testAgent() + agent.Tools = []string{"activity_breakdown"} + + exec := NewModelExecutor(gw, &MemorySink{}, reg) + res, err := exec.executeWithLimits(context.Background(), agent, testInput("what happened today"), + Limits{MaxSteps: 4, MaxToolCalls: 1, MaxTokens: 100_000, Deadline: 30 * time.Second}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Termination != TerminationCompleted { + t.Fatalf("Termination = %q, want Completed", res.Termination) + } + if len(gw.seen) < 2 { + t.Fatalf("expected at least 2 model calls, got %d", len(gw.seen)) + } + if len(gw.seen[0].Tools) == 0 { + t.Error("the first call must offer the catalogue") + } + if n := len(gw.seen[len(gw.seen)-1].Tools); n != 0 { + t.Errorf("the final call carried %d tool definitions; the tool budget was spent", n) + } +}