4 Commits

Author SHA1 Message Date
a372340281 add the greeting msg
Some checks failed
CI / test (push) Failing after 4m41s
CI / fixture (push) Failing after 9s
2026-10-05 16:20:17 +05:30
939598a187 Add the deploy runbook for db4803c and the Gemini switch
Some checks failed
CI / test (push) Failing after 4m38s
CI / fixture (push) Failing after 10s
Image first, config second -- the old binary on Gemini config fails
every tool-using run on its second model call, and the runbook says
why, what was verified, and how to roll back either half.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
2026-09-22 13:16:46 +05:30
db4803c557 Report an unsaved trajectory somewhere that survives
A failed save must not fail the run -- the answer already exists --
but until now the only record of the loss was an error entry
appended to the trajectory that had just failed to save. §6 says
the trajectory is not optional telemetry; losing one silently is
the worst version of losing one.

The runtime has no logger by design, so ExecutionResult gains an
Unsaved list the surface reads and turns into a §10 log line with
run_id, tenant_id, agent_key and agent_version. Never serialised
to the client. Covered on all three run paths, streaming included.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
2026-09-22 13:15:39 +05:30
797ee5f2d2 Round-trip provider metadata on tool calls; Gemini requires it
Gemini 3 models attach a thought signature to every function call
and reject the follow-up -- 400, "Function call is missing a
thought_signature in functionCall parts" -- when the assistant
message echoing that call does not carry it back. The gateway
rebuilt the assistant turn from id, name and arguments alone, so
every tool-using run on Gemini died on its second model call, after
a first call that looked perfectly healthy. Found when production
was pointed at Gemini on 2026-09-22; rolled back to Groq within
minutes.

ToolCall gains an opaque Extra field: the raw JSON of the wire's
extra_content, captured on both the streaming and non-streaming
paths and emitted verbatim on the next request. The gateway does
not read it and must not -- the point of one wire shape is that a
vendor's private fields pass through untouched. Absent stays
absent; no provider receives a null it never sent.

Also: Gemini wraps its error body in a one-element array, which
the message parser read as "no detail". The trajectory therefore
said only "the model rejected the request" where the body named
the missing signature outright. Unwrapped now, so the next
provider quirk is legible in the trajectory instead of costing a
day of proxy captures.

Verified end to end with the real gateway against real Gemini: a
three-turn tool-calling run completed and the proxy confirmed the
signature on every echoed call.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
2026-09-22 13:15:39 +05:30
15 changed files with 1439 additions and 35 deletions

176
docs/deploy-db4803c.md Normal file
View File

@@ -0,0 +1,176 @@
# Deploying `db4803c` and switching the model vendor to Gemini
Two things land together, and the order matters: the image **must** be
running before the configuration switches vendor. Old binary on Gemini
config = every tool-using run dies on its second model call (§2). New binary
on Groq config = works exactly as today. So: image first, config second.
Live today: Groq free tier, `8000 TPM`, **35% of runs since Sep 9 end in
`gateway.rate_limited`**. That number is why this deploy exists.
---
## 1. What is being deployed
`db4803c`, on `origin/main`. Five commits since `8e36faf`:
| Commit | What |
| --- | --- |
| `822b3b1` | `.env` untracked again; ignore rules `3455ad0` deleted are back |
| `b765495` | Archiving an agent a published agent delegates to → 409 |
| `5166fde` | `GatewayFailure` termination; **migration 000016** |
| `797ee5f` | Provider metadata round-tripped on tool calls — **Gemini needs this** |
| `db4803c` | An unsaved trajectory is logged, not just noted in itself |
**One migration.** `000016` widens `agent_runs.termination_check` to admit
`GatewayFailure`. The `migrate` init container applies it before the API
starts. Its down migration is verified (folds rows to `ToolFailure` before
narrowing the CHECK), so a rollback of the image is safe.
## 2. Why the image must go first
Gemini 3 models attach a `thought_signature` to every function call and
reject the follow-up request without it. `797ee5f` teaches the gateway to
carry it back. The image on the cluster today does not, and it was proved on
2026-09-22: pointed at Gemini, the first model call succeeded, the tool ran,
the second call answered `400 Function call is missing a thought_signature`.
Rolled back to Groq within minutes.
## 3. What was verified before writing this
With the `797ee5f` binary, locally, against the production database and the
real Gemini API through a request-logging proxy:
| Check | Result |
| --- | --- |
| `positions-agent`, 2 model calls | `Completed`, signature present on the echoed call |
| `krow-workforce-agent`, 3 tool-calling turns, **12,123 tokens** | `Completed`, correct answer. This exceeds Groq's entire per-minute ceiling |
| Streaming path carries the signature | unit test + live run |
| `GatewayFailure` reaches the surface with its own wording | seen live before rollback |
| Full Go suite against a real Postgres (the DB tests skip without one) | 18/18 packages |
| Migration 000016 up → down → up on a scratch database | clean |
Model reliability, six bare calls each, 2026-09-22 ~13:00 IST:
| Model | HTTP codes |
| --- | --- |
| `gemini-3.8-flash` | 503 503 200 200 503 503 |
| `gemini-3.5-flash` | 503 200 503 200 200 200 |
| `gemini-3.5-flash-lite` | 200 200 200 200 200 200 |
The gateway retries a 503 three times with backoff; at the rates above a
three-call run on either larger model still fails often. **All three tiers
run `gemini-3.5-flash-lite`** until the larger models stop shedding load or
the key is on a paid tier. `gemini-3.1-pro-preview` answers 429 (pro is not
on the free tier); the 2.5 family is listed but blocked for new keys.
## 4. Build and push the image (the other machine)
The Dockerfile cross-compiles, so any host with `buildx` and a Docker Hub
login works:
```bash
git checkout db4803c
docker buildx build --platform linux/amd64 \
-f infrastructure/Dockerfile.api \
-t doormile/krowbackend:db4803c -t doormile/krowbackend:latest \
--push .
```
Two tags on purpose: `:latest` is what the StatefulSet pulls; `:db4803c` is
what you roll back **to** if you need to (§7). Confirm before touching the
cluster:
```bash
docker buildx imagetools inspect doormile/krowbackend:latest | grep -E 'Platform|Digest' | head -3
```
## 5. Switch the cluster (on the server, as root)
The Gemini key is already staged in the Secret as `MODEL_API_KEY_GEMINI`;
the Groq key stays as `MODEL_API_KEY_GROQ`. The manifests in
`/opt/kubernetes/manifests/krow` already describe the Gemini configuration
(committed `pending`), so the config half is an `apply`.
```bash
# 1. the credential the API reads becomes the Gemini one
kubectl -n krow patch secret krow-model --type=json \
-p '[{"op":"copy","from":"/data/MODEL_API_KEY_GEMINI","path":"/data/MODEL_API_KEY"}]'
# 2. configmap → Gemini base URL and model ids
kubectl apply -k /opt/kubernetes/manifests/krow/
# 3. new pods: pull :latest, run migration 16, boot on the new config
kubectl -n krow rollout restart statefulset/krow
kubectl -n krow rollout status statefulset/krow --timeout=5m
```
`rollout status` waits for krow-2, then krow-1, each gated on readiness. If
krow-2 does not come up, krow-1 is still serving on the old image and Groq.
## 6. Smoke test
```bash
# migration 16 applied?
kubectl -n krow logs krow-2 -c migrate | tail -2 # want: 16/u gateway_failure_termination
# boot on the right vendor?
kubectl -n krow logs krow-2 -c api | grep -m1 '"listening"' | grep -o '"endpoints":[0-9]*' # want 71
# a real run, through the public URL
T=$(curl -s -D - -o /dev/null -X POST https://mcp.krowforce.com/api/v1/auth/login \
-H 'Content-Type: application/json' \
-d '{"email":"demo@krow.app","password":"<demo password>"}' \
| sed -n 's/^[Ss]et-[Cc]ookie: krow_session=\([^;]*\).*/\1/p')
curl -s -X POST https://mcp.krowforce.com/api/v1/agents/65bfd77d-2f74-4548-ab52-4e720e153397/runs \
-H "Cookie: krow_session=$T" -H 'Content-Type: application/json' \
-d '{"input":"How many open positions are there?"}' | grep -E '"(termination|output)"'
```
Want `"termination": "Completed"` and a count. `GatewayFailure` with
"usually it is busy" is Gemini shedding load — retry once. `ToolFailure` on
the **second** model call means the old image is still running (§2).
Then check the trajectory landed with the right model:
```bash
# from a machine with psql / the postgres image; DATABASE_URL from secret/krow-db
psql "$DATABASE_URL" -Atc "SELECT model, termination FROM agent_runs ORDER BY started_at DESC LIMIT 1"
```
Want `gemini-3.5-flash-lite | Completed`.
## 7. Rollback
Config only (image stays — it works on Groq too):
```bash
kubectl -n krow patch secret krow-model --type=json \
-p '[{"op":"copy","from":"/data/MODEL_API_KEY_GROQ","path":"/data/MODEL_API_KEY"}]'
kubectl -n krow patch cm krow-config --type merge -p '{"data":{
"MODEL_BASE_URL":"https://api.groq.com/openai/v1",
"MODEL_FAST":"openai/gpt-oss-20b","MODEL_BALANCED":"openai/gpt-oss-120b","MODEL_DEEP":"openai/gpt-oss-120b"}}'
kubectl -n krow rollout restart statefulset/krow
```
Image too (only if `db4803c` itself misbehaves):
```bash
kubectl -n krow set image statefulset/krow api=doormile/krowbackend:<previous tag>
```
Migration 16 stays applied; the old binary never writes `GatewayFailure`,
so the wider CHECK is harmless to it. Reverse it only if you must:
`migrate ... down 1` — it folds existing `GatewayFailure` rows to `ToolFailure`.
## 8. Not covered here
- **`activity-agent` is archived while `krow-workforce-agent v2` delegates to
it.** `b765495` prevents this happening again; it does not repair the
existing case. Either unarchive `activity-agent` or publish workforce v3
without it — a product decision.
- **`OAUTH_LOGIN_PATH` (`/login`) 404s on `mcp.krowforce.com`.** A signed-out
MCP consent redirect goes nowhere. Signed-in users are unaffected.
- **Rotation.** The Anthropic key in `3455ad0`'s history, the Groq key, the
Gemini key (pasted in a chat), the DB admin password (8 chars, public IP,
no TLS), the root SSH password, the demo login.

View File

@@ -85,6 +85,25 @@ type ToolCall struct {
ID string
Name string
Input json.RawMessage
// Extra is provider metadata attached to the call, carried back to the
// provider verbatim on the next turn and never read here.
//
// It exists because at least one provider requires it. Gemini 3 models
// attach a "thought signature" to every function call and REJECT the
// follow-up request — 400, "Function call is missing a thought_signature"
// — if the assistant message that echoes the call does not carry it back.
// A gateway that rebuilds the assistant turn from ID, Name and Input alone
// drops it, and every tool-using run dies on its second model call while
// the first one looked perfectly healthy. That is exactly what happened
// on 2026-09-22 when production was pointed at Gemini.
//
// The gateway does not know what is in it and must not: the whole point
// of speaking one wire shape is that a vendor's private fields pass
// through untouched. It is the raw JSON of the call's extra_content
// object, or nil when the provider sent none, in which case it is omitted
// from the request again.
Extra json.RawMessage
}
// ToolResult is what came back, on its way to the model.

View File

@@ -247,6 +247,7 @@ func (g *OpenAIGateway) decode(
// The raw JSON, not a parsed value — handed to the handler's own
// decoder rather than matched on as a string here.
Input: json.RawMessage(args),
Extra: c.ExtraContent,
})
}
@@ -301,6 +302,10 @@ type oaiToolCall struct {
ID string `json:"id,omitempty"`
Type string `json:"type,omitempty"`
Function oaiFunctionRef `json:"function"`
// ExtraContent is the provider's own metadata on the call, round-tripped
// as raw JSON. See ToolCall.Extra for why it is not optional.
ExtraContent json.RawMessage `json:"extra_content,omitempty"`
}
type oaiFunctionRef struct {
@@ -493,9 +498,10 @@ func encodeOpenAIMessages(system string, msgs []Message) []oaiMessage {
args = "{}"
}
msg.ToolCalls = append(msg.ToolCalls, oaiToolCall{
ID: c.ID,
Type: "function",
Function: oaiFunctionRef{Name: c.Name, Arguments: args},
ID: c.ID,
Type: "function",
Function: oaiFunctionRef{Name: c.Name, Arguments: args},
ExtraContent: c.Extra,
})
}
out = append(out, msg)
@@ -549,6 +555,20 @@ func translateOpenAI(status int, body []byte) error {
// worth reading and not often enough to depend on, so an unparseable body
// yields nothing rather than failing a failure.
func openAIErrorMessage(body []byte) string {
// Gemini wraps its error in a one-element ARRAY — `[{"error":{...}}]` —
// where OpenAI, Groq and the rest send the object bare. Unwrapped here
// rather than tolerated as "no detail", because the detail is the whole
// value of the field: for two weeks the trajectory said only "the model
// rejected the request" when the body said "Function call is missing a
// thought_signature", and the difference was a day of diagnosis.
body = bytes.TrimSpace(body)
if bytes.HasPrefix(body, []byte("[")) {
var many []json.RawMessage
if err := json.Unmarshal(body, &many); err != nil || len(many) == 0 {
return ""
}
body = many[0]
}
var envelope struct {
Error struct {
Message string `json:"message"`
@@ -744,6 +764,11 @@ func (a *streamAccumulator) addToolCallDeltas(deltas []oaiToolCall) {
if d.Function.Name != "" {
call.Function.Name = d.Function.Name
}
// Provider metadata arrives whole on one fragment, like the id. Kept
// when non-empty so a later empty fragment does not erase it.
if len(d.ExtraContent) > 0 {
call.ExtraContent = d.ExtraContent
}
// Arguments are the fragmented field: concatenated, never replaced.
call.Function.Arguments += d.Function.Arguments
}

View File

@@ -209,3 +209,29 @@ func TestStreamCompleteUsesTheStreamingPath(t *testing.T) {
t.Errorf("Text = %q", resp.Text)
}
}
// The streamed shape of TestToolCallProviderMetadataIsRoundTripped: the
// metadata arrives on one fragment, and later fragments that carry only
// argument text must not erase it.
func TestStreamKeepsToolCallProviderMetadata(t *testing.T) {
const sig = `{"google":{"thought_signature":"El4KXAFpFH0T4CM3"}}`
acc, err := accumulateSSE(strings.NewReader(strings.Join([]string{
`data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_a","function":{"name":"open_positions","arguments":""},"extra_content":` + sig + `}]}}]}`,
`data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{}"}}]}}]}`,
`data: {"choices":[{"delta":{},"finish_reason":"tool_calls"}]}`,
`data: [DONE]`,
}, "\n\n")), func(string) {})
if err != nil {
t.Fatalf("accumulateSSE: %v", err)
}
msg := acc.message()
if len(msg.ToolCalls) != 1 {
t.Fatalf("got %d tool calls, want 1", len(msg.ToolCalls))
}
if string(msg.ToolCalls[0].ExtraContent) != sig {
t.Errorf("extra_content after streaming = %s, want %s", msg.ToolCalls[0].ExtraContent, sig)
}
if msg.ToolCalls[0].Function.Arguments != "{}" {
t.Errorf("arguments = %q, want {}", msg.ToolCalls[0].Function.Arguments)
}
}

View File

@@ -339,3 +339,96 @@ func TestBaseURLDefaultsAndTrimsSlash(t *testing.T) {
t.Errorf("endpoint = %q, want the trailing slash collapsed", got)
}
}
// THE FAILURE THIS EXISTS FOR: a provider that attaches private metadata to a
// tool call and refuses the follow-up without it. Gemini 3 does exactly this
// ("Function call is missing a thought_signature"), and a gateway that rebuilt
// the assistant turn from id, name and arguments alone killed every tool-using
// run on its second model call — after a first call that looked healthy.
//
// The round trip is tested end to end: the provider's extra_content on the
// response must reappear, byte for byte, on the next request's echo of that
// call. The gateway must not care what is inside it.
func TestToolCallProviderMetadataIsRoundTripped(t *testing.T) {
const sig = `{"google":{"thought_signature":"El4KXAFpFH0T4CM3"}}`
gw, captured := serve(t, func(w http.ResponseWriter, _ *oaiRequest) {
_, _ = io.WriteString(w, `{
"choices":[{"message":{"role":"assistant","tool_calls":[
{"id":"call_x","type":"function",
"function":{"name":"open_positions","arguments":"{}"},
"extra_content":`+sig+`}]},
"finish_reason":"tool_calls"}],
"usage":{"prompt_tokens":10,"completion_tokens":5}
}`)
})
resp, err := gw.Complete(context.Background(), ask("how many open positions?"))
if err != nil {
t.Fatalf("Complete: %v", err)
}
if len(resp.ToolCalls) != 1 {
t.Fatalf("got %d tool calls, want 1", len(resp.ToolCalls))
}
if string(resp.ToolCalls[0].Extra) != sig {
t.Fatalf("Extra = %s, want the provider's extra_content verbatim", resp.ToolCalls[0].Extra)
}
// Second turn: the loop echoes the assistant's call and adds the result.
// This is the request Gemini rejects when the signature is missing.
_, err = gw.Complete(context.Background(), Request{Tier: TierBalanced, Messages: []Message{
{Role: RoleUser, Text: "how many open positions?"},
{Role: RoleAssistant, ToolCalls: resp.ToolCalls},
{Role: RoleUser, ToolResults: []ToolResult{{CallID: "call_x", Content: `{"count":14}`}}},
}})
if err != nil {
t.Fatalf("second Complete: %v", err)
}
var echoed *oaiToolCall
for i := range captured.Messages {
if len(captured.Messages[i].ToolCalls) > 0 {
echoed = &captured.Messages[i].ToolCalls[0]
}
}
if echoed == nil {
t.Fatalf("the second request did not echo the assistant's tool call: %+v", captured.Messages)
}
if string(echoed.ExtraContent) != sig {
t.Errorf("echoed extra_content = %s, want %s", echoed.ExtraContent, sig)
}
}
// A provider that sends no metadata must not receive an "extra_content": null
// it never asked for. Absent stays absent.
func TestToolCallWithoutProviderMetadataOmitsTheField(t *testing.T) {
msgs := []Message{
{Role: RoleAssistant, ToolCalls: []ToolCall{{ID: "call_1", Name: "open_positions", Input: json.RawMessage(`{}`)}}},
}
raw, err := json.Marshal(encodeOpenAIMessages("", msgs))
if err != nil {
t.Fatal(err)
}
if strings.Contains(string(raw), "extra_content") {
t.Errorf("extra_content was emitted for a call that had none: %s", raw)
}
}
// Gemini wraps its error in a one-element array. The detail must survive,
// because a bare "the model rejected the request" is the difference between a
// one-line diagnosis and a day of one.
func TestProviderErrorDetailSurvivesArrayEnvelope(t *testing.T) {
cases := map[string]string{
`{"error":{"message":"bare object"}}`: "bare object",
`[{"error":{"message":"array wrapped"}}]`: "array wrapped",
` [ {"error":{"message":"padded"}} ] `: "padded",
`{"message":"top level"}`: "top level",
`[]`: "",
`not json`: "",
}
for body, want := range cases {
if got := openAIErrorMessage([]byte(body)); got != want {
t.Errorf("openAIErrorMessage(%s) = %q, want %q", body, got, want)
}
}
}

View File

@@ -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)
}
})
}
}

View File

@@ -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")
}
}

View File

@@ -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"
)
@@ -162,10 +163,62 @@ func (s *Server) handleAgentRun(w http.ResponseWriter, r *http.Request) {
writeError(w, s.log, runLoadError(runErr))
return
}
s.logUnsaved(ident, res)
s.logGatewayFailure(ident, res)
writeJSON(w, http.StatusOK, buildRunResponse(res))
}
// logUnsaved is the operator's record of a run whose trajectory did not
// persist. The run itself already answered; §6 says the trajectory is not
// optional telemetry, so losing one is an error even when nothing else went
// wrong, and it carries every field §10 asks a log line to carry.
func (s *Server) logUnsaved(ident authctx.Identity, res *runtime.ExecutionResult) {
for _, detail := range res.Unsaved {
s.log.Error("trajectory unsaved",
"run_id", res.RunID, "tenant_id", ident.OrgID,
"agent_key", res.AgentID, "agent_version", res.AgentVersion,
"detail", detail)
}
}
// 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
@@ -192,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
}
@@ -208,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 ""
@@ -225,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 —
@@ -336,6 +445,8 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au
writeError(w, s.log, runLoadError(runErr))
return
}
s.logUnsaved(ident, res)
s.logGatewayFailure(ident, res)
writeJSON(w, http.StatusOK, buildRunResponse(res))
return
}
@@ -383,6 +494,8 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au
return
}
s.logUnsaved(ident, res)
s.logGatewayFailure(ident, res)
send(map[string]any{"run": buildRunResponse(res)})
fmt.Fprint(w, "data: [DONE]\n\n")
flusher.Flush()

View File

@@ -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}
}
}

View File

@@ -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 <context> 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,21 @@ 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.
if err := m.sink.Save(ctx, traj); err != nil {
var unsaved []string
if err := m.sink.Save(persistCtx, traj); err != nil {
rec.Error("runtime.trajectory_unsaved", err.Error())
unsaved = append(unsaved, traj.RunID+": "+err.Error())
}
// Delegated runs are written AFTER this one, because parent_run_id is a
@@ -614,13 +728,15 @@ 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())
}
}
res := &ExecutionResult{
Unsaved: unsaved,
Success: term == TerminationCompleted,
Output: output,
AgentID: agent.ID,

View File

@@ -265,6 +265,28 @@ func TestSinkFailureDoesNotFailTheRun(t *testing.T) {
if res.Output != "the answer" {
t.Errorf("Output = %q, want the answer through", res.Output)
}
// ...but the loss must be reported somewhere that survives. The runtime
// has no logger; the surface reads this and writes the operator's line.
// Before this field existed the only record was an entry in the very
// trajectory that had failed to save.
if len(res.Unsaved) != 1 || !strings.Contains(res.Unsaved[0], res.RunID) ||
!strings.Contains(res.Unsaved[0], context.DeadlineExceeded.Error()) {
t.Errorf("Unsaved = %v, want one entry naming run %s and the error", res.Unsaved, res.RunID)
}
}
// A healthy run reports nothing unsaved. Guarded so the surface never logs a
// phantom loss.
func TestHealthySaveReportsNothingUnsaved(t *testing.T) {
gw := &fakeGateway{text: "the answer"}
exec := NewModelExecutor(gw, &MemorySink{}, nil)
res, err := exec.ExecuteAgent(context.Background(), testAgent(), testInput("q"))
if err != nil {
t.Fatal(err)
}
if len(res.Unsaved) != 0 {
t.Errorf("Unsaved = %v on a healthy run, want none", res.Unsaved)
}
}
type failingSink struct{}

View File

@@ -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)
}
}

View File

@@ -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

View File

@@ -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)
}
}

View File

@@ -167,6 +167,18 @@ type ExecutionResult struct {
// the token of whichever the person approves as ExecutionInput.Confirmation
// on the next call.
Confirmations []*tools.Confirmation `json:"confirmations,omitempty"`
// Unsaved names the trajectories this run produced that could not be
// persisted, as "<run id>: <error>". Empty on every healthy run.
//
// A failed save must not fail the run — the answer already exists — but
// it must not vanish either. Until 2026-09-22 the only record of it was
// an entry appended to the trajectory that had just failed to save, which
// is a note left in a bottle that sank. The runtime has no logger by
// design; the surface does, and reads this to write the §10 line an
// operator can grep for. Never serialised to the client: it is an
// operator's concern, not the caller's.
Unsaved []string `json:"-"`
}
// RuntimeError is a structured error containing context for execution failures.