Compare commits
6 Commits
feat/expor
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 939598a187 | |||
| db4803c557 | |||
| 797ee5f2d2 | |||
| 5166fde764 | |||
| b765495eb7 | |||
| 822b3b1edb |
67
.env
67
.env
@@ -1,67 +0,0 @@
|
|||||||
# ============================================================================
|
|
||||||
# Krow backend — example environment
|
|
||||||
#
|
|
||||||
# Copy to .env and fill in. .env is gitignored and must never be committed.
|
|
||||||
# Every value below is a placeholder or a safe local default: no real password,
|
|
||||||
# API key or token belongs in this file.
|
|
||||||
#
|
|
||||||
# cp .env.example .env
|
|
||||||
# ============================================================================
|
|
||||||
|
|
||||||
# ── Application ─────────────────────────────────────────────────────────────
|
|
||||||
APP_ENV=development
|
|
||||||
LOG_LEVEL=info # debug | info | warn | error
|
|
||||||
|
|
||||||
# ── HTTP server ─────────────────────────────────────────────────────────────
|
|
||||||
HTTP_HOST=127.0.0.1
|
|
||||||
HTTP_PORT=8080
|
|
||||||
HTTP_READ_TIMEOUT=15s
|
|
||||||
# An agent run on the deep tier may legally take 2m0s. A 30s write timeout
|
|
||||||
# aborted the response mid-run and the proxy reported it as 502.
|
|
||||||
HTTP_WRITE_TIMEOUT=2m30s
|
|
||||||
HTTP_IDLE_TIMEOUT=60s
|
|
||||||
HTTP_SHUTDOWN_TIMEOUT=10s
|
|
||||||
|
|
||||||
# ── PostgreSQL ──────────────────────────────────────────────────────────────
|
|
||||||
# The local development database. DATABASE_NAME is mixed-case and hyphenated,
|
|
||||||
# so anything that interpolates it into SQL must quote it: "Krow-force".
|
|
||||||
DATABASE_HOST=127.0.0.1
|
|
||||||
DATABASE_PORT=5432
|
|
||||||
DATABASE_NAME=Krow-force
|
|
||||||
DATABASE_USER=postgres
|
|
||||||
DATABASE_PASSWORD=
|
|
||||||
DATABASE_SCHEMA=public
|
|
||||||
|
|
||||||
# sslmode: disable is fine for a loopback dev database. APP_ENV=production
|
|
||||||
# rejects `disable` at startup — use require or verify-full there.
|
|
||||||
DATABASE_SSLMODE=disable
|
|
||||||
|
|
||||||
# Pool and timeout tuning.
|
|
||||||
DATABASE_MAX_OPEN_CONNS=25
|
|
||||||
DATABASE_MIN_IDLE_CONNS=2
|
|
||||||
DATABASE_CONN_MAX_LIFETIME=30m
|
|
||||||
DATABASE_CONNECT_TIMEOUT=5s
|
|
||||||
DATABASE_STATEMENT_TIMEOUT=10s
|
|
||||||
|
|
||||||
# ── Migrations ──────────────────────────────────────────────────────────────
|
|
||||||
# Consumed by the Makefile, which builds the golang-migrate URL from the
|
|
||||||
# DATABASE_* values above. Keep it pointed at the repository's migrations/.
|
|
||||||
MIGRATIONS_DIR=./migrations
|
|
||||||
|
|
||||||
# ── Seed ────────────────────────────────────────────────────────────────────
|
|
||||||
# The demo fixture, generated from the frontend repository's src/api/seed.js.
|
|
||||||
# The Makefile passes an absolute path; this default suits running from the
|
|
||||||
# repository root.
|
|
||||||
SEED_FIXTURE_PATH=./seed/fixtures/seed.json
|
|
||||||
|
|
||||||
# ── Model gateway ───────────────────────────────────────────────────────────
|
|
||||||
ANTHROPIC_API_KEY=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
|
|
||||||
|
|
||||||
# ── Knowledge layer (retrieval) ─────────────────────────────────────────────
|
|
||||||
# A real semantic embedding model, running locally. No credential, no per-token
|
|
||||||
# cost, and no tenant text leaving this machine. Change either of the first two
|
|
||||||
# and the stored vectors stop being searched — run `make reembed ORG=<slug>`.
|
|
||||||
EMBED_PROVIDER=ollama
|
|
||||||
EMBED_MODEL=nomic-embed-text
|
|
||||||
EMBED_DIMENSIONS=768
|
|
||||||
EMBED_BASE_URL=http://localhost:11434
|
|
||||||
3
.gitignore
vendored
3
.gitignore
vendored
@@ -1,5 +1,8 @@
|
|||||||
# Secrets and local configuration
|
# Secrets and local configuration
|
||||||
# A bare pattern matches at any depth, so this covers infrastructure/.env too.
|
# A bare pattern matches at any depth, so this covers infrastructure/.env too.
|
||||||
|
.env
|
||||||
|
.env.local
|
||||||
|
.env.*.local
|
||||||
|
|
||||||
# TLS material. The local-db overlay generates a self-signed pair inside the
|
# TLS material. The local-db overlay generates a self-signed pair inside the
|
||||||
# postgres volume, but nothing stops someone dropping certs here by hand.
|
# postgres volume, but nothing stops someone dropping certs here by hand.
|
||||||
|
|||||||
31
CLAUDE.md
31
CLAUDE.md
@@ -136,7 +136,7 @@ The agent loop lives in `src/runtime/loop.py`. It is the highest-risk file in th
|
|||||||
- Single loop, spec-driven. No per-agent branching.
|
- Single loop, spec-driven. No per-agent branching.
|
||||||
- Decrement budgets **before** dispatch, not after, so a hung tool cannot overrun.
|
- Decrement budgets **before** dispatch, not after, so a hung tool cannot overrun.
|
||||||
- Stream partial assistant text as it arrives; buffer tool calls until complete.
|
- Stream partial assistant text as it arrives; buffer tool calls until complete.
|
||||||
- Termination reasons are an enum: `Completed | BudgetExceeded | Deadline | ConfirmationPending | ToolFailure | Refused`. Every run ends with exactly one.
|
- Termination reasons are an enum: `Completed | BudgetExceeded | Deadline | ConfirmationPending | ToolFailure | GatewayFailure | Refused`. Every run ends with exactly one. `GatewayFailure` is the model provider not answering (rate limited, request rejected, credential refused, unreachable) and is deliberately not `ToolFailure`: the two are different operational questions, and until 2026-09-22 the enum could not tell them apart.
|
||||||
- Persist a full trajectory per run: every message, tool call, tool result, and budget snapshot. This is what makes debugging and evals possible — it is not optional telemetry.
|
- Persist a full trajectory per run: every message, tool call, tool result, and budget snapshot. This is what makes debugging and evals possible — it is not optional telemetry.
|
||||||
- Delegation is a tool call from the parent's perspective. Subagent runs get their own trajectory, linked by `parent_run_id`.
|
- Delegation is a tool call from the parent's perspective. Subagent runs get their own trajectory, linked by `parent_run_id`.
|
||||||
|
|
||||||
@@ -220,21 +220,11 @@ depends on the curated-versus-self-serve decision and is not settled.
|
|||||||
| Layer | State |
|
| Layer | State |
|
||||||
|---|---|
|
|---|---|
|
||||||
| Surfaces | `POST /api/v1/agents/{id}/runs` (streams over SSE on `Accept: text/event-stream`), `GET /api/v1/runs/{id}`; the chat panel is the only answering path — the browser simulator is deleted |
|
| Surfaces | `POST /api/v1/agents/{id}/runs` (streams over SSE on `Accept: text/event-stream`), `GET /api/v1/runs/{id}`; the chat panel is the only answering path — the browser simulator is deleted |
|
||||||
| Orchestration | spec-driven loop, four bounds claimed before dispatch, six terminations, trajectories in `agent_runs`; delegation per §6 — a subagent is a tool call, runs as the caller, shares the parent budget, capped at depth 2, and writes its own trajectory linked by `parent_run_id` |
|
| Orchestration | spec-driven loop, four bounds claimed before dispatch, seven terminations, trajectories in `agent_runs`; delegation per §6 — a subagent is a tool call, runs as the caller, shares the parent budget, capped at depth 2, and writes its own trajectory linked by `parent_run_id` |
|
||||||
| Registry | 9 agents + 24 skills as rows; published versions immutable (append-only, trigger-enforced); runs pin the version they started with |
|
| Registry | 9 agents + 23 skills as rows; published versions immutable (append-only, trigger-enforced); runs pin the version they started with |
|
||||||
| Tools | 19, two of which write (`move_application`, `assign_worker`), behind a bound single-use confirmation |
|
| Tools | 19, two of which write (`move_application`, `assign_worker`), behind a bound single-use confirmation |
|
||||||
| Knowledge | ACL-tagged ingest, hybrid dense + BM25 fused with RRF, pre-filtered |
|
| Knowledge | ACL-tagged ingest, hybrid dense + BM25 fused with RRF, pre-filtered |
|
||||||
| Gateway | tier → model + effort, token accounting, refusal as an outcome; one wire protocol — `openai`, the chat-completions shape that Groq (the default), Gemini, OpenRouter, Together, vLLM and a local Ollama all serve. The Anthropic path was removed; `MODEL_PROVIDER=anthropic`, a stale `ANTHROPIC_API_KEY` and a leftover `claude-*` id are each refused at startup rather than ignored |
|
| Gateway | tier → model + effort, token accounting, refusal as an outcome |
|
||||||
|
|
||||||
**Conversational writes are not agent tool calls.** Two skills — `create-position`
|
|
||||||
and `create-employee-role` — collect a record through the chat panel and then
|
|
||||||
write it with the same REST call the manual form uses, as the signed-in user.
|
|
||||||
They are therefore outside I4's confirmation-token mechanism, which governs
|
|
||||||
tools an AGENT invokes on a caller's behalf. The person is making the request
|
|
||||||
themselves, and the flow's review step ("Ready to create this position?") is
|
|
||||||
where they agree to it. Worth knowing rather than worth fixing: if a write is
|
|
||||||
ever moved from the panel into an agent tool, it acquires I4's bound single-use
|
|
||||||
confirmation at that point and not before.
|
|
||||||
|
|
||||||
**Deviations from this document, all deliberate and all flagged in code:**
|
**Deviations from this document, all deliberate and all flagged in code:**
|
||||||
|
|
||||||
@@ -260,18 +250,7 @@ confirmation at that point and not before.
|
|||||||
Do not resolve these unilaterally. Flag them and ask.
|
Do not resolve these unilaterally. Flag them and ask.
|
||||||
|
|
||||||
- **Who authors agents?** Curated (the team ships specs) vs. self-serve (tenants author their own). Self-serve requires prompt-injection hardening at the authoring boundary, per-tenant cost caps, an approval workflow, and a sandbox — roughly 3× the platform. Current assumption: **curated**, with the registry designed so self-serve is additive later.
|
- **Who authors agents?** Curated (the team ships specs) vs. self-serve (tenants author their own). Self-serve requires prompt-injection hardening at the authoring boundary, per-tenant cost caps, an approval workflow, and a sandbox — roughly 3× the platform. Current assumption: **curated**, with the registry designed so self-serve is additive later.
|
||||||
- **Model hosting.** Self-hosted vs. API vs. mixed by tier. **Still open** —
|
- **Model hosting.** Self-hosted vs. API vs. mixed by tier.
|
||||||
but no longer expensive to change: `MODEL_BASE_URL` + the three `MODEL_*` ids
|
|
||||||
move the whole platform between Groq (the default), Gemini, OpenRouter,
|
|
||||||
Together, vLLM and a local Ollama without a code change, and `make eval-live`
|
|
||||||
runs the suite against whichever is configured. Decide it on the eval
|
|
||||||
evidence, and weigh the I7 case heaviest: a cheaper model that follows the
|
|
||||||
planted injection is a security regression, not a saving.
|
|
||||||
|
|
||||||
The gap that the Anthropic removal opened here is closed: `openai/gpt-oss-120b`
|
|
||||||
on Groq has been through `make eval-live` and passes all three cases including
|
|
||||||
I7 (2026-09-07). The decision itself — self-hosted vs. API vs. mixed by tier —
|
|
||||||
is still open and still not mine to settle.
|
|
||||||
- **Confirmation UX.** Inline in-chat vs. an approval queue.
|
- **Confirmation UX.** Inline in-chat vs. an approval queue.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|||||||
176
docs/deploy-db4803c.md
Normal file
176
docs/deploy-db4803c.md
Normal 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.
|
||||||
@@ -733,6 +733,9 @@ func TestMigrationPairsAreComplete(t *testing.T) {
|
|||||||
// Phase 5: shared rate limit counters, so a limit means the same thing
|
// Phase 5: shared rate limit counters, so a limit means the same thing
|
||||||
// behind one instance and behind ten.
|
// behind one instance and behind ten.
|
||||||
"000015_rate_limits.up.sql",
|
"000015_rate_limits.up.sql",
|
||||||
|
// A seventh termination reason. The CHECK in 000006 was chosen so
|
||||||
|
// this would be a migration rather than an ALTER TYPE; this is it.
|
||||||
|
"000016_gateway_failure_termination.up.sql",
|
||||||
}
|
}
|
||||||
if len(ups) != len(want) {
|
if len(ups) != len(want) {
|
||||||
t.Fatalf("%d migrations, want %d — update this list deliberately", len(ups), len(want))
|
t.Fatalf("%d migrations, want %d — update this list deliberately", len(ups), len(want))
|
||||||
|
|||||||
@@ -85,6 +85,25 @@ type ToolCall struct {
|
|||||||
ID string
|
ID string
|
||||||
Name string
|
Name string
|
||||||
Input json.RawMessage
|
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.
|
// ToolResult is what came back, on its way to the model.
|
||||||
|
|||||||
@@ -247,6 +247,7 @@ func (g *OpenAIGateway) decode(
|
|||||||
// The raw JSON, not a parsed value — handed to the handler's own
|
// The raw JSON, not a parsed value — handed to the handler's own
|
||||||
// decoder rather than matched on as a string here.
|
// decoder rather than matched on as a string here.
|
||||||
Input: json.RawMessage(args),
|
Input: json.RawMessage(args),
|
||||||
|
Extra: c.ExtraContent,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -301,6 +302,10 @@ type oaiToolCall struct {
|
|||||||
ID string `json:"id,omitempty"`
|
ID string `json:"id,omitempty"`
|
||||||
Type string `json:"type,omitempty"`
|
Type string `json:"type,omitempty"`
|
||||||
Function oaiFunctionRef `json:"function"`
|
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 {
|
type oaiFunctionRef struct {
|
||||||
@@ -493,9 +498,10 @@ func encodeOpenAIMessages(system string, msgs []Message) []oaiMessage {
|
|||||||
args = "{}"
|
args = "{}"
|
||||||
}
|
}
|
||||||
msg.ToolCalls = append(msg.ToolCalls, oaiToolCall{
|
msg.ToolCalls = append(msg.ToolCalls, oaiToolCall{
|
||||||
ID: c.ID,
|
ID: c.ID,
|
||||||
Type: "function",
|
Type: "function",
|
||||||
Function: oaiFunctionRef{Name: c.Name, Arguments: args},
|
Function: oaiFunctionRef{Name: c.Name, Arguments: args},
|
||||||
|
ExtraContent: c.Extra,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
out = append(out, msg)
|
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
|
// worth reading and not often enough to depend on, so an unparseable body
|
||||||
// yields nothing rather than failing a failure.
|
// yields nothing rather than failing a failure.
|
||||||
func openAIErrorMessage(body []byte) string {
|
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 {
|
var envelope struct {
|
||||||
Error struct {
|
Error struct {
|
||||||
Message string `json:"message"`
|
Message string `json:"message"`
|
||||||
@@ -744,6 +764,11 @@ func (a *streamAccumulator) addToolCallDeltas(deltas []oaiToolCall) {
|
|||||||
if d.Function.Name != "" {
|
if d.Function.Name != "" {
|
||||||
call.Function.Name = 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.
|
// Arguments are the fragmented field: concatenated, never replaced.
|
||||||
call.Function.Arguments += d.Function.Arguments
|
call.Function.Arguments += d.Function.Arguments
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -209,3 +209,29 @@ func TestStreamCompleteUsesTheStreamingPath(t *testing.T) {
|
|||||||
t.Errorf("Text = %q", resp.Text)
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -339,3 +339,96 @@ func TestBaseURLDefaultsAndTrimsSlash(t *testing.T) {
|
|||||||
t.Errorf("endpoint = %q, want the trailing slash collapsed", got)
|
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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -1520,3 +1520,84 @@ Find the shifts nobody has taken.
|
|||||||
t.Errorf("version moved on a restore: %v -> %v", arc["version"], got)
|
t.Errorf("version moved on a restore: %v -> %v", arc["version"], got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The reverse of the cycle and unknown-key checks. Those prove an edge is
|
||||||
|
// valid when the PARENT is written; this proves the edge stays valid when the
|
||||||
|
// CHILD is archived. Without it a published parent keeps delegating into
|
||||||
|
// nothing — exactly what krow-workforce-agent did after activity-agent was
|
||||||
|
// archived under it on 2026-09-15.
|
||||||
|
func TestArchivingADelegatedSubagentIsRefused(t *testing.T) {
|
||||||
|
r := newRBAC(t)
|
||||||
|
|
||||||
|
agent := func(id, name, status string, version int, subagents ...string) string {
|
||||||
|
var sub string
|
||||||
|
if len(subagents) > 0 {
|
||||||
|
sub = "subagents:\n"
|
||||||
|
for _, s := range subagents {
|
||||||
|
sub += " - " + s + "\n"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return fmt.Sprintf(`---
|
||||||
|
id: %s
|
||||||
|
name: %s
|
||||||
|
description: part of a delegation graph
|
||||||
|
status: %s
|
||||||
|
version: %d
|
||||||
|
pages:
|
||||||
|
- candidates
|
||||||
|
%s---
|
||||||
|
|
||||||
|
## Instructions
|
||||||
|
Delegate.
|
||||||
|
`, id, name, status, version, sub)
|
||||||
|
}
|
||||||
|
|
||||||
|
res := r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||||
|
"markdown": agent("dep-child", "Child", "published", 1), "visibility": "organization",
|
||||||
|
})
|
||||||
|
if res.code != http.StatusCreated {
|
||||||
|
t.Fatalf("create child: status %d (%v)", res.code, res.body)
|
||||||
|
}
|
||||||
|
childID, _ := res.record(t)["id"].(string)
|
||||||
|
|
||||||
|
res = r.as(r.admin, "POST", "/api/v1/agent-definitions", map[string]any{
|
||||||
|
"markdown": agent("dep-parent", "Parent", "published", 1, "dep-child"), "visibility": "organization",
|
||||||
|
})
|
||||||
|
if res.code != http.StatusCreated {
|
||||||
|
t.Fatalf("create parent: status %d (%v)", res.code, res.body)
|
||||||
|
}
|
||||||
|
parentID, _ := res.record(t)["id"].(string)
|
||||||
|
|
||||||
|
// Both ways of archiving must be refused: the status-only patch the UI
|
||||||
|
// sends, and a markdown save whose frontmatter says archived.
|
||||||
|
for name, patch := range map[string]map[string]any{
|
||||||
|
"status patch": {"status": "archived"},
|
||||||
|
"markdown save": {"markdown": agent("dep-child", "Child", "archived", 2)},
|
||||||
|
} {
|
||||||
|
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+childID, patch)
|
||||||
|
if res.code != http.StatusConflict {
|
||||||
|
t.Fatalf("%s: archiving a delegated-to agent: status %d, want 409 (%v)", name, res.code, res.body)
|
||||||
|
}
|
||||||
|
if body := fmt.Sprint(res.body); !strings.Contains(body, "dep-parent") {
|
||||||
|
t.Errorf("%s: the refusal did not name the dependent: %v", name, res.body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Refused, not half-applied.
|
||||||
|
res = r.as(r.admin, "GET", "/api/v1/agent-definitions/"+childID, nil)
|
||||||
|
if st, _ := res.record(t)["status"].(string); st != "published" {
|
||||||
|
t.Fatalf("child status after refused archives = %q, want published", st)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A DRAFT parent does not pin the child. Move the parent to draft and the
|
||||||
|
// archive goes through: an abandoned experiment must not hold a
|
||||||
|
// production agent in place.
|
||||||
|
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+parentID, map[string]any{"status": "draft"})
|
||||||
|
if res.code != http.StatusOK {
|
||||||
|
t.Fatalf("draft the parent: status %d (%v)", res.code, res.body)
|
||||||
|
}
|
||||||
|
res = r.as(r.admin, "PATCH", "/api/v1/agent-definitions/"+childID, map[string]any{"status": "archived"})
|
||||||
|
if res.code != http.StatusOK {
|
||||||
|
t.Fatalf("archive with only a draft dependent: status %d, want 200 (%v)", res.code, res.body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -162,10 +162,24 @@ func (s *Server) handleAgentRun(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, s.log, runLoadError(runErr))
|
writeError(w, s.log, runLoadError(runErr))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
s.logUnsaved(ident, res)
|
||||||
|
|
||||||
writeJSON(w, http.StatusOK, buildRunResponse(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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// buildRunResponse turns a runtime result into the client's shape.
|
// buildRunResponse turns a runtime result into the client's shape.
|
||||||
//
|
//
|
||||||
// Every termination answers 200. That looks wrong at first and is not: the
|
// Every termination answers 200. That looks wrong at first and is not: the
|
||||||
@@ -205,7 +219,7 @@ func buildRunResponse(res *runtime.ExecutionResult) runResponse {
|
|||||||
// reached its budget before finishing" is accurate and means nothing to
|
// reached its budget before finishing" is accurate and means nothing to
|
||||||
// somebody who has never heard of a token budget.
|
// somebody who has never heard of a token budget.
|
||||||
//
|
//
|
||||||
// Every one of the six is spelled out. A default that said "something went
|
// 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
|
// 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.
|
// indistinguishable to the person best placed to tell us which it was.
|
||||||
func terminationMessage(t runtime.Termination) string {
|
func terminationMessage(t runtime.Termination) string {
|
||||||
@@ -224,6 +238,12 @@ func terminationMessage(t runtime.Termination) string {
|
|||||||
"Nothing was changed."
|
"Nothing was changed."
|
||||||
case runtime.TerminationRefused:
|
case runtime.TerminationRefused:
|
||||||
return "The agent declined to answer this one."
|
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."
|
||||||
default:
|
default:
|
||||||
return "The agent did not finish."
|
return "The agent did not finish."
|
||||||
}
|
}
|
||||||
@@ -330,6 +350,7 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au
|
|||||||
writeError(w, s.log, runLoadError(runErr))
|
writeError(w, s.log, runLoadError(runErr))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
s.logUnsaved(ident, res)
|
||||||
writeJSON(w, http.StatusOK, buildRunResponse(res))
|
writeJSON(w, http.StatusOK, buildRunResponse(res))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -377,6 +398,7 @@ func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident au
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
s.logUnsaved(ident, res)
|
||||||
send(map[string]any{"run": buildRunResponse(res)})
|
send(map[string]any{"run": buildRunResponse(res)})
|
||||||
fmt.Fprint(w, "data: [DONE]\n\n")
|
fmt.Fprint(w, "data: [DONE]\n\n")
|
||||||
flusher.Flush()
|
flusher.Flush()
|
||||||
|
|||||||
@@ -22,13 +22,26 @@ const (
|
|||||||
TerminationConfirmationPending Termination = "ConfirmationPending"
|
TerminationConfirmationPending Termination = "ConfirmationPending"
|
||||||
TerminationToolFailure Termination = "ToolFailure"
|
TerminationToolFailure Termination = "ToolFailure"
|
||||||
TerminationRefused Termination = "Refused"
|
TerminationRefused Termination = "Refused"
|
||||||
|
|
||||||
|
// GatewayFailure is the model provider failing to answer at all: rate
|
||||||
|
// limited, rejected the request, refused the credential, or unreachable.
|
||||||
|
// Added 2026-09-22 because until then every one of those was recorded as
|
||||||
|
// ToolFailure, and 131 of 318 production runs read as "a tool is broken"
|
||||||
|
// when no tool had failed — 45 of them were Groq's free-tier rate limit,
|
||||||
|
// which is a capacity decision, not a bug. The two are different questions
|
||||||
|
// to an operator ("what did we break" versus "what are we not paying
|
||||||
|
// for"), and an enum that could not tell them apart hid the answer for
|
||||||
|
// two weeks. Refused and Deadline keep their own reasons; this is the
|
||||||
|
// rest of the gateway's vocabulary.
|
||||||
|
TerminationGatewayFailure Termination = "GatewayFailure"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Valid reports whether t is one of the six.
|
// Valid reports whether t is one of the seven.
|
||||||
func (t Termination) Valid() bool {
|
func (t Termination) Valid() bool {
|
||||||
switch t {
|
switch t {
|
||||||
case TerminationCompleted, TerminationBudgetExceeded, TerminationDeadline,
|
case TerminationCompleted, TerminationBudgetExceeded, TerminationDeadline,
|
||||||
TerminationConfirmationPending, TerminationToolFailure, TerminationRefused:
|
TerminationConfirmationPending, TerminationToolFailure, TerminationRefused,
|
||||||
|
TerminationGatewayFailure:
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
|
|||||||
@@ -241,12 +241,16 @@ func (m *ModelExecutor) delegate(
|
|||||||
rec.AddChildren(collected.Runs)
|
rec.AddChildren(collected.Runs)
|
||||||
|
|
||||||
if res == nil {
|
if res == nil {
|
||||||
msg := "the subagent returned nothing"
|
// The parent sees a delegation as a tool call, but the REASON it
|
||||||
|
// failed is still worth carrying: a subagent the provider rate limited
|
||||||
|
// should read as GatewayFailure in the parent's trajectory too, or the
|
||||||
|
// parent's operator goes looking for a tool that never broke.
|
||||||
|
msg, term := "the subagent returned nothing", TerminationToolFailure
|
||||||
if err != nil {
|
if err != nil {
|
||||||
msg = err.Error()
|
msg, term = err.Error(), terminationFor(err)
|
||||||
}
|
}
|
||||||
return delegationAnswer{
|
return delegationAnswer{
|
||||||
Agent: sub.ID, Termination: string(TerminationToolFailure), Error: msg,
|
Agent: sub.ID, Termination: string(term), Error: msg,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -555,6 +555,12 @@ func (m *ModelExecutor) runTools(
|
|||||||
// a Refused run is one that must not be retried. Flattening them into a single
|
// a Refused run is one that must not be retried. Flattening them into a single
|
||||||
// failure reason would make every one of those distinctions unanswerable from
|
// failure reason would make every one of those distinctions unanswerable from
|
||||||
// the trajectory.
|
// the trajectory.
|
||||||
|
//
|
||||||
|
// Any other gateway error is GatewayFailure, not ToolFailure. The provider
|
||||||
|
// being rate limited, rejecting the request or refusing the key is not a tool
|
||||||
|
// failing, and calling it one sent two weeks of operators looking for a broken
|
||||||
|
// tool that did not exist. A non-gateway error — a tool the model invented, a
|
||||||
|
// result that could not be encoded — is still the tool layer's.
|
||||||
func terminationFor(err error) Termination {
|
func terminationFor(err error) Termination {
|
||||||
var gwErr *gateway.Error
|
var gwErr *gateway.Error
|
||||||
if !errors.As(err, &gwErr) {
|
if !errors.As(err, &gwErr) {
|
||||||
@@ -566,7 +572,7 @@ func terminationFor(err error) Termination {
|
|||||||
case gateway.CodeTimeout:
|
case gateway.CodeTimeout:
|
||||||
return TerminationDeadline
|
return TerminationDeadline
|
||||||
default:
|
default:
|
||||||
return TerminationToolFailure
|
return TerminationGatewayFailure
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -599,8 +605,10 @@ func (m *ModelExecutor) finish(
|
|||||||
// A sink that fails must not fail the run — the answer was already
|
// 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
|
// produced. It is recorded in the trajectory we could not save, which is
|
||||||
// the best available place for it.
|
// the best available place for it.
|
||||||
|
var unsaved []string
|
||||||
if err := m.sink.Save(ctx, traj); err != nil {
|
if err := m.sink.Save(ctx, traj); err != nil {
|
||||||
rec.Error("runtime.trajectory_unsaved", err.Error())
|
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
|
// Delegated runs are written AFTER this one, because parent_run_id is a
|
||||||
@@ -611,10 +619,12 @@ func (m *ModelExecutor) finish(
|
|||||||
if err := m.sink.Save(ctx, child); err != nil {
|
if err := m.sink.Save(ctx, child); err != nil {
|
||||||
rec.Error("runtime.subrun_unsaved",
|
rec.Error("runtime.subrun_unsaved",
|
||||||
fmt.Sprintf("%s: %s", child.RunID, err.Error()))
|
fmt.Sprintf("%s: %s", child.RunID, err.Error()))
|
||||||
|
unsaved = append(unsaved, child.RunID+": "+err.Error())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
res := &ExecutionResult{
|
res := &ExecutionResult{
|
||||||
|
Unsaved: unsaved,
|
||||||
Success: term == TerminationCompleted,
|
Success: term == TerminationCompleted,
|
||||||
Output: output,
|
Output: output,
|
||||||
AgentID: agent.ID,
|
AgentID: agent.ID,
|
||||||
@@ -657,6 +667,8 @@ func terminationMessage(t Termination) string {
|
|||||||
return "the run is waiting on a confirmation"
|
return "the run is waiting on a confirmation"
|
||||||
case TerminationToolFailure:
|
case TerminationToolFailure:
|
||||||
return "the run failed"
|
return "the run failed"
|
||||||
|
case TerminationGatewayFailure:
|
||||||
|
return "the model provider did not answer"
|
||||||
default:
|
default:
|
||||||
return string(t)
|
return string(t)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -265,6 +265,28 @@ func TestSinkFailureDoesNotFailTheRun(t *testing.T) {
|
|||||||
if res.Output != "the answer" {
|
if res.Output != "the answer" {
|
||||||
t.Errorf("Output = %q, want the answer through", res.Output)
|
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{}
|
type failingSink struct{}
|
||||||
@@ -277,6 +299,7 @@ func TestTerminationValidRejectsInvented(t *testing.T) {
|
|||||||
for _, ok := range []Termination{
|
for _, ok := range []Termination{
|
||||||
TerminationCompleted, TerminationBudgetExceeded, TerminationDeadline,
|
TerminationCompleted, TerminationBudgetExceeded, TerminationDeadline,
|
||||||
TerminationConfirmationPending, TerminationToolFailure, TerminationRefused,
|
TerminationConfirmationPending, TerminationToolFailure, TerminationRefused,
|
||||||
|
TerminationGatewayFailure,
|
||||||
} {
|
} {
|
||||||
if !ok.Valid() {
|
if !ok.Valid() {
|
||||||
t.Errorf("%q should be a valid termination", ok)
|
t.Errorf("%q should be a valid termination", ok)
|
||||||
@@ -1075,3 +1098,31 @@ func TestTheTrajectoryRecordsWhichChunksGroundedTheAnswer(t *testing.T) {
|
|||||||
t.Error("nothing in the trajectory says a retrieval happened")
|
t.Error("nothing in the trajectory says a retrieval happened")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The whole reason GatewayFailure exists: a provider that will not answer is
|
||||||
|
// not a tool that broke, and for two weeks the trajectory said it was.
|
||||||
|
func TestTerminationForSeparatesGatewayFromTool(t *testing.T) {
|
||||||
|
gw := func(code string) error { return &gateway.Error{Code: code, Message: "x"} }
|
||||||
|
cases := []struct {
|
||||||
|
name string
|
||||||
|
err error
|
||||||
|
want Termination
|
||||||
|
}{
|
||||||
|
{"rate limited is the gateway's", gw(gateway.CodeRateLimited), TerminationGatewayFailure},
|
||||||
|
{"invalid request is the gateway's", gw(gateway.CodeInvalidRequest), TerminationGatewayFailure},
|
||||||
|
{"bad credential is the gateway's", gw(gateway.CodeUnauthorized), TerminationGatewayFailure},
|
||||||
|
{"upstream is the gateway's", gw(gateway.CodeUpstream), TerminationGatewayFailure},
|
||||||
|
{"not configured is the gateway's", gw(gateway.CodeNotConfigured), TerminationGatewayFailure},
|
||||||
|
{"refused keeps its own reason", gw(gateway.CodeRefused), TerminationRefused},
|
||||||
|
{"timeout keeps its own reason", gw(gateway.CodeTimeout), TerminationDeadline},
|
||||||
|
{"a non-gateway error is still the tool layer's", errors.New("tool exploded"), TerminationToolFailure},
|
||||||
|
{"a wrapped gateway error is still found", fmt.Errorf("delegating: %w", gw(gateway.CodeRateLimited)), TerminationGatewayFailure},
|
||||||
|
}
|
||||||
|
for _, tc := range cases {
|
||||||
|
t.Run(tc.name, func(t *testing.T) {
|
||||||
|
if got := terminationFor(tc.err); got != tc.want {
|
||||||
|
t.Errorf("terminationFor(%v) = %s, want %s", tc.err, got, tc.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -285,7 +285,7 @@ func (r *Recorder) SetModel(model string) {
|
|||||||
|
|
||||||
// Finish closes the trajectory with its termination reason and returns it.
|
// Finish closes the trajectory with its termination reason and returns it.
|
||||||
//
|
//
|
||||||
// A reason that is not one of the six is recorded as ToolFailure rather than
|
// A reason that is not one of the seven is recorded as ToolFailure rather than
|
||||||
// stored as-is: an unrecognised termination is a bug in the loop, and writing
|
// stored as-is: an unrecognised termination is a bug in the loop, and writing
|
||||||
// it verbatim would let that bug propagate into every eval and dashboard that
|
// it verbatim would let that bug propagate into every eval and dashboard that
|
||||||
// groups by this column.
|
// groups by this column.
|
||||||
|
|||||||
@@ -167,6 +167,18 @@ type ExecutionResult struct {
|
|||||||
// the token of whichever the person approves as ExecutionInput.Confirmation
|
// the token of whichever the person approves as ExecutionInput.Confirmation
|
||||||
// on the next call.
|
// on the next call.
|
||||||
Confirmations []*tools.Confirmation `json:"confirmations,omitempty"`
|
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.
|
// RuntimeError is a structured error containing context for execution failures.
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"sort"
|
||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
@@ -352,6 +353,11 @@ func (s *DefinitionsService) UpdateAgent(ctx context.Context, ident authctx.Iden
|
|||||||
if err := s.refuseSubagentCycle(ctx, ident, agent.ID, agent.Subagents); err != nil {
|
if err := s.refuseSubagentCycle(ctx, ident, agent.ID, agent.Subagents); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
if agent.Status == "archived" {
|
||||||
|
if err := s.refuseArchivingDependency(ctx, ident, agent.ID); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
if err := s.refusePublishedRewrite(ctx, ident, repo.KindAgent,
|
||||||
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
agent.ID, markdown, agent.Status, agent.Version); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -361,6 +367,12 @@ func (s *DefinitionsService) UpdateAgent(ctx context.Context, ident authctx.Iden
|
|||||||
if !isStr || (status != "draft" && status != "published" && status != "archived") {
|
if !isStr || (status != "draft" && status != "published" && status != "archived") {
|
||||||
return nil, domain.Validation("status must be one of: draft, published, archived", map[string]string{"status": "invalid"})
|
return nil, domain.Validation("status must be one of: draft, published, archived", map[string]string{"status": "invalid"})
|
||||||
}
|
}
|
||||||
|
if status == "archived" {
|
||||||
|
did, _ := existing["definition_id"].(string)
|
||||||
|
if err := s.refuseArchivingDependency(ctx, ident, did); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
input.Status = &status
|
input.Status = &status
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -671,6 +683,70 @@ func (s *DefinitionsService) refuseSubagentCycle(ctx context.Context,
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// refuseArchivingDependency fails an archive while a published agent in the
|
||||||
|
// organization still delegates to the definition.
|
||||||
|
//
|
||||||
|
// §3 says an unknown subagent key fails at publish, not at run time, and
|
||||||
|
// refuseSubagentCycle is half of that. This is the other half. Publish
|
||||||
|
// validation proves the edge exists when the PARENT is written; nothing
|
||||||
|
// re-checked it when the CHILD was later archived, so a spec could pass
|
||||||
|
// validation on Monday and be delegating into nothing by Friday. That is what
|
||||||
|
// happened to krow-workforce-agent on 2026-09-15: activity-agent was archived
|
||||||
|
// under it, and every run since logged runtime.unknown_subagent and answered
|
||||||
|
// activity questions without its activity capability — quietly, because the
|
||||||
|
// parent still Completed.
|
||||||
|
//
|
||||||
|
// Only published parents count. A draft that names this agent is the author's
|
||||||
|
// problem at their next publish, where rejectUnknownTools-style validation
|
||||||
|
// will tell them; refusing an archive on the strength of a draft would let an
|
||||||
|
// abandoned experiment pin a production agent in place forever.
|
||||||
|
//
|
||||||
|
// Unlike refuseSubagentCycle this FAILS CLOSED when the graph cannot be read.
|
||||||
|
// The cycle check can afford to fail open because the runtime depth cap holds
|
||||||
|
// regardless; the only backstop here is a parent that keeps answering with a
|
||||||
|
// capability missing, which is the failure this exists to prevent.
|
||||||
|
func (s *DefinitionsService) refuseArchivingDependency(ctx context.Context,
|
||||||
|
ident authctx.Identity, definitionID string) error {
|
||||||
|
|
||||||
|
if definitionID == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
rows, _, err := s.repo.ListAgents(ctx, ident, repo.DefinitionListParams{
|
||||||
|
Visibility: "organization", Limit: 500,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return domain.Internal(fmt.Errorf("could not check whether any published agent delegates to %q: %w", definitionID, err))
|
||||||
|
}
|
||||||
|
|
||||||
|
var dependents []string
|
||||||
|
for _, rec := range rows {
|
||||||
|
id, _ := rec["definition_id"].(string)
|
||||||
|
markdown, _ := rec["markdown"].(string)
|
||||||
|
status, _ := rec["status"].(string)
|
||||||
|
if id == "" || id == definitionID || markdown == "" || status != "published" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
parsed, err := definition.ParseAgent(markdown, definition.Options{})
|
||||||
|
if err != nil || parsed == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, sub := range parsed.Subagents {
|
||||||
|
if sub == definitionID {
|
||||||
|
dependents = append(dependents, id)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(dependents) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
sort.Strings(dependents)
|
||||||
|
return domain.Conflict(fmt.Sprintf(
|
||||||
|
"%s cannot be archived while a published agent delegates to it: %s. "+
|
||||||
|
"Publish a version of each without it in `subagents` first, or archive them too.",
|
||||||
|
definitionID, strings.Join(dependents, ", ")))
|
||||||
|
}
|
||||||
|
|
||||||
// refusePublishedRewrite fails a publish that would change a version already
|
// refusePublishedRewrite fails a publish that would change a version already
|
||||||
// published, BEFORE anything is written.
|
// published, BEFORE anything is written.
|
||||||
//
|
//
|
||||||
|
|||||||
10
migrations/000016_gateway_failure_termination.down.sql
Normal file
10
migrations/000016_gateway_failure_termination.down.sql
Normal file
@@ -0,0 +1,10 @@
|
|||||||
|
-- Reverting narrows the vocabulary, so rows that used the seventh value are
|
||||||
|
-- folded back into ToolFailure FIRST — the classification they would have had
|
||||||
|
-- before 000016 — or the narrower CHECK cannot be re-added at all. This loses
|
||||||
|
-- the distinction the up migration introduced, which is what reverting means.
|
||||||
|
UPDATE agent_runs SET termination = 'ToolFailure' WHERE termination = 'GatewayFailure';
|
||||||
|
ALTER TABLE agent_runs DROP CONSTRAINT agent_runs_termination_check;
|
||||||
|
ALTER TABLE agent_runs ADD CONSTRAINT agent_runs_termination_check CHECK (termination IN (
|
||||||
|
'Completed', 'BudgetExceeded', 'Deadline',
|
||||||
|
'ConfirmationPending', 'ToolFailure', 'Refused'
|
||||||
|
));
|
||||||
22
migrations/000016_gateway_failure_termination.up.sql
Normal file
22
migrations/000016_gateway_failure_termination.up.sql
Normal file
@@ -0,0 +1,22 @@
|
|||||||
|
-- A seventh termination reason: GatewayFailure.
|
||||||
|
--
|
||||||
|
-- 000006 chose a CHECK over an enum type so that "a new termination reason
|
||||||
|
-- should be a migration, but not one that requires ALTER TYPE". This is that
|
||||||
|
-- migration.
|
||||||
|
--
|
||||||
|
-- Until now the model provider failing — rate limited, request rejected, key
|
||||||
|
-- refused, unreachable — was recorded as ToolFailure, because the enum had no
|
||||||
|
-- other place for it. On 2026-09-22 that was 131 of 318 production runs, of
|
||||||
|
-- which not one was a tool failing. The distinction is the difference between
|
||||||
|
-- "what did we break" and "what are we not paying for", and the column that
|
||||||
|
-- every dashboard and eval groups by could not make it.
|
||||||
|
--
|
||||||
|
-- Existing rows are left as they are. Rewriting history from the trajectory
|
||||||
|
-- text would be guesswork against a message format that has changed twice;
|
||||||
|
-- the trajectory entries still carry the gateway.* error code for anyone who
|
||||||
|
-- needs to reclassify the past.
|
||||||
|
ALTER TABLE agent_runs DROP CONSTRAINT agent_runs_termination_check;
|
||||||
|
ALTER TABLE agent_runs ADD CONSTRAINT agent_runs_termination_check CHECK (termination IN (
|
||||||
|
'Completed', 'BudgetExceeded', 'Deadline',
|
||||||
|
'ConfirmationPending', 'ToolFailure', 'Refused', 'GatewayFailure'
|
||||||
|
));
|
||||||
Reference in New Issue
Block a user