Files
krow_backend/go-api/internal/httpserver/runs.go
Suriya 5166fde764
Some checks failed
CI / test (push) Failing after 4m37s
CI / fixture (push) Failing after 8s
Add GatewayFailure: the provider not answering is not a tool failing
terminationFor sent every gateway error that was not Refused or Timeout
to ToolFailure, because the enum had nowhere else to put it. On
2026-09-22 that was 131 of 318 production runs, and not one of them was
a tool failing: 56 were retired model ids, 25 an exhausted Anthropic
balance, 45 Groq's free-tier rate limit -- the only one still happening.
An operator reading the termination column saw "a tool is broken" for
two weeks while the actual answer was "we are not paying for capacity".

GatewayFailure is the seventh termination. Rate limited, request
rejected, credential refused and unreachable land there; Refused and
Deadline keep their own reasons; a non-gateway error is still the tool
layer's. A delegation whose subagent died at the gateway now carries
that reason up to the parent instead of reading as a tool call that
failed.

Migration 000016 widens the CHECK that 000006 chose precisely so this
would be a migration rather than an ALTER TYPE. Its down folds any
GatewayFailure rows back to ToolFailure BEFORE narrowing the constraint,
which is the order that works; verified up, down and up again on a
scratch database. Existing rows are left as they are -- the trajectory
entries still carry the gateway.* code for anyone reclassifying history.

The surface wording is the one termination where "try again" is honest
advice, since the dominant cause clears within a minute.

Full suite run against a real database, including the tests that skip
without one.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
2026-09-22 12:48:09 +05:30

390 lines
16 KiB
Go

package httpserver
import (
"encoding/json"
"errors"
"fmt"
"net/http"
"strings"
"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/runtime"
"github.com/krow/krow-backend/go-api/internal/tools"
)
// The agent run endpoint: the surface layer, and the first thing that can
// actually call the runtime.
//
// Everything under internal/runtime, internal/tools and internal/knowledge has
// been reachable only from tests until now. This file is the seam, and it has
// two jobs that belong nowhere else:
//
// 1. **Deriving user-facing text.** §10 says user-facing wording is produced
// at the surface, not raised from the core. The runtime returns a
// Termination — an enum — and this file decides what a person reads for
// each of the six. A run that hit its budget is not an internal error and
// must not be answered as one.
// 2. **Answering with a shape the client can act on.** A ConfirmationPending
// run is not a failure: it is a question, it comes back 200 with the
// confirmation payload, and the client's job is to ask a person and call
// back with the token. Answering it 500 would make the whole write path
// look broken.
func (s *Server) routeRuns(mux *http.ServeMux) int {
if s.agents == nil {
// No runtime wired — no model credential, or a deployment that does not
// serve agents. The routes are not registered at all rather than
// registered and always failing: a 404 says "this deployment does not
// do that", where a 500 says "this deployment is broken", and only one
// of those is true.
return 0
}
mux.HandleFunc("POST /api/v1/agents/{id}/runs", s.handleAgentRun)
mux.HandleFunc("GET /api/v1/runs/{runId}", s.handleRunGet)
return 2
}
/* ── Request and response ───────────────────────────────────────────────── */
// runRequest is what a client sends to run an agent.
type runRequest struct {
// Input is the caller's question. Required.
Input string `json:"input"`
// AgentVersion pins the run to a published version.
//
// A client resuming a conversation sends the version the FIRST answer came
// back with — every response carries it — so the conversation stays on the
// agent it started with even if somebody publishes an edit mid-thread. Zero
// or absent means whatever is current, which is what a fresh question wants.
//
// It matters most on an approval: a person approved a write while looking
// at one version, and carrying it out under a newer one would perform
// something they were never shown.
AgentVersion int `json:"agentVersion,omitempty"`
// Confirmation is a token a person approved, carried into a resumed run.
//
// It authorises ONE call — the exact tool and arguments it was issued
// against — and supplying it does not put the run into a permissive mode. A
// second write in the same run raises its own confirmation, because a
// person approved one thing. See tools/confirm.go.
Confirmation string `json:"confirmation,omitempty"`
// Context is opaque client state passed to the runtime. Never used for
// authorization: the principal comes from the session, always.
Context map[string]any `json:"context,omitempty"`
}
// runResponse is what comes back.
//
// Deliberately not the ExecutionResult. That struct carries a Go `error` and
// internal wording; this one carries a code and a sentence written for a
// person, which is the §10 boundary made concrete.
type runResponse struct {
RunID string `json:"runId"`
AgentID string `json:"agentId"`
Version int `json:"agentVersion,omitempty"`
Termination string `json:"termination"`
// Output is the assistant's text. Present on a completed run, and also on a
// bounded one — a run that hit its deadline mid-sentence still said
// something, and throwing it away helps nobody.
Output string `json:"output,omitempty"`
// Message is what to show a person when the run did not complete. Derived
// here from the termination, never raised from the core.
Message string `json:"message,omitempty"`
// Confirmations are writes the agent proposed and did not perform. Present
// exactly when termination is ConfirmationPending.
Confirmations []*tools.Confirmation `json:"confirmations,omitempty"`
Usage runUsage `json:"usage"`
}
// runUsage is the token accounting, flattened for the client.
type runUsage struct {
InputTokens int64 `json:"inputTokens"`
OutputTokens int64 `json:"outputTokens"`
CachedTokens int64 `json:"cachedTokens"`
TotalTokens int64 `json:"totalTokens"`
ModelCalls int `json:"modelCalls"`
}
/* ── Running an agent ───────────────────────────────────────────────────── */
// handleAgentRun executes one agent turn.
func (s *Server) handleAgentRun(w http.ResponseWriter, r *http.Request) {
ident, err := authctx.MustFrom(r.Context())
if err != nil {
writeError(w, s.log, domain.Internal(err))
return
}
var req runRequest
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, maxRunRequestBytes)).Decode(&req); err != nil {
writeError(w, s.log, domain.Validation("the request body was not valid JSON", nil))
return
}
if strings.TrimSpace(req.Input) == "" {
writeError(w, s.log, domain.Validation("a run needs an input", map[string]string{
"input": "required",
}))
return
}
// Streamed when the client asks for it, by Accept rather than by a second
// route. It is the same run with the same semantics — the same principal,
// the same budgets, the same confirmation gate — delivered differently. Two
// routes would be two things to keep in step, and the one that drifted
// would be the one nobody tested.
if wantsSSE(r) {
s.streamAgentRun(w, r, ident, req)
return
}
// The principal is the SESSION's, never the body's. I1 begins here: a
// client that could name its own principal could read anything.
res, runErr := s.agents.RunAgent(r.Context(), ident, r.PathValue("id"), runtime.ExecutionInput{
Identity: ident,
Input: req.Input,
AgentVersion: req.AgentVersion,
Confirmation: req.Confirmation,
Context: req.Context,
})
// A load failure — no such agent, not this tenant's, draft, archived — is a
// resource error and answers like one. It is distinguishable from a run
// that started and ended badly, which is the distinction below.
if res == nil || res.Termination == "" {
writeError(w, s.log, runLoadError(runErr))
return
}
writeJSON(w, http.StatusOK, buildRunResponse(res))
}
// buildRunResponse turns a runtime result into the client's shape.
//
// Every termination answers 200. That looks wrong at first and is not: the
// question "did the HTTP request succeed" and the question "did the agent
// finish" are different questions, and collapsing them costs the client the
// second one. A run that hit its budget is a run — it has an id, a trajectory,
// a token cost and often a partial answer — and answering 500 would throw all
// of that away while telling the client to retry something that will fail the
// same way.
func buildRunResponse(res *runtime.ExecutionResult) runResponse {
out := runResponse{
RunID: res.RunID,
AgentID: res.AgentID,
Version: res.AgentVersion,
Termination: string(res.Termination),
Output: res.Output,
Confirmations: res.Confirmations,
Usage: runUsage{
InputTokens: res.Usage.InputTokens,
OutputTokens: res.Usage.OutputTokens,
CachedTokens: res.Usage.CachedTokens,
TotalTokens: res.Usage.TotalTokens,
ModelCalls: res.Usage.ModelCalls,
},
}
if res.Termination != runtime.TerminationCompleted {
out.Message = terminationMessage(res.Termination)
}
return out
}
// terminationMessage is the user-facing wording for each termination.
//
// §10's boundary, and the reason it lives here rather than in the runtime: the
// core's terminationMessage is an internal explanation for a log, and this one
// is a sentence a venue manager reads. They differ on purpose — "the run
// reached its budget before finishing" is accurate and means nothing to
// somebody who has never heard of a token budget.
//
// 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 {
switch t {
case runtime.TerminationCompleted:
return ""
case runtime.TerminationBudgetExceeded:
return "This question needed more work than the agent is allowed to spend in one go. " +
"Try asking for a narrower slice of it."
case runtime.TerminationDeadline:
return "The agent ran out of time before finishing. Anything it had already worked out is above."
case runtime.TerminationConfirmationPending:
return "The agent has proposed a change and is waiting for you to approve it."
case runtime.TerminationToolFailure:
return "The agent could not finish — something it needed did not answer. " +
"Nothing was changed."
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."
default:
return "The agent did not finish."
}
}
// runLoadError maps a pre-run failure onto the API's error vocabulary.
//
// These are the errors from LoadExecutableAgent, raised before any run began —
// so there is no run id, no trajectory and no termination. They are resource
// errors and answer like resource errors.
//
// ErrNotFound and ErrUnauthorized deliberately both become 404. §8's rule about
// denials applies to agents as much as to rows: "this agent exists but is not
// yours" and "there is no such agent" must not be distinguishable, or the
// endpoint becomes a way to enumerate other tenants' agents one id at a time.
func runLoadError(err error) error {
switch {
case err == nil:
return domain.Internal(errors.New("the run produced no result and no error"))
case errors.Is(err, runtime.ErrNotFound), errors.Is(err, runtime.ErrUnauthorized):
return domain.NotFound("agent", "")
case errors.Is(err, runtime.ErrDraftAgent):
return domain.Validation("this agent is still a draft and cannot be run", nil)
case errors.Is(err, runtime.ErrArchivedAgent):
return domain.Validation("this agent is archived and cannot be run", nil)
case errors.Is(err, runtime.ErrNotExecutable),
errors.Is(err, runtime.ErrInvalidDefinition):
return domain.Validation("this agent is not in a runnable state", nil)
case errors.Is(err, runtime.ErrDependencyMissing),
errors.Is(err, runtime.ErrDependencyInactive),
errors.Is(err, runtime.ErrCircularDependency):
return domain.Validation("this agent depends on a skill that is missing or inactive", nil)
default:
return domain.Internal(err)
}
}
// maxRunRequestBytes bounds a run request body.
//
// A question, not a document. Retrieval is how a corpus reaches the model, and
// it goes through the permission layer; a client posting a megabyte of text
// would be routing around that — the text would land in the prompt having been
// read by nobody and authorized by nothing.
const maxRunRequestBytes = 64 << 10
/* ── Reading a trajectory ───────────────────────────────────────────────── */
// handleRunGet returns a recorded run.
//
// §6 requires a full trajectory per run, and this is what makes it worth
// having: "why did the agent say that" is answerable by a support conversation
// pointing at a run id.
//
// Tenant-scoped by the store, not by this handler. I5 — the predicate lives in
// the query, so a run id from another organization is simply absent and answers
// 404, indistinguishable from one that never existed.
func (s *Server) handleRunGet(w http.ResponseWriter, r *http.Request) {
ident, err := authctx.MustFrom(r.Context())
if err != nil {
writeError(w, s.log, domain.Internal(err))
return
}
traj, err := s.runs.Load(r.Context(), ident, r.PathValue("runId"))
if err != nil {
writeError(w, s.log, err)
return
}
writeJSON(w, http.StatusOK, traj)
}
/* ── Streaming ──────────────────────────────────────────────────────────── */
// wantsSSE reports whether the client asked for a streamed response.
func wantsSSE(r *http.Request) bool {
return strings.Contains(r.Header.Get("Accept"), "text/event-stream")
}
// streamAgentRun runs an agent, sending text as it arrives.
//
// The wire format is one JSON object per SSE event, which is the same shape the
// non-streaming response uses for its parts:
//
// {"delta": "…"} assistant text, as the model produces it
// {"run": { … }} the finished run — termination, confirmations, usage
// {"error": { … }} a run that could not start
//
// The final `run` event carries the SAME body the non-streaming path returns.
// That is what keeps the two honest: a client can ignore every delta, read only
// the last event, and be in exactly the state it would have been in without
// streaming.
func (s *Server) streamAgentRun(w http.ResponseWriter, r *http.Request, ident authctx.Identity, req runRequest) {
flusher, ok := w.(http.Flusher)
if !ok {
// Something between here and the client buffers. Streaming into it
// would deliver the whole answer at the end anyway, but silently — so
// the honest move is to answer normally rather than pretend.
res, runErr := s.agents.RunAgent(r.Context(), ident, r.PathValue("id"), runtime.ExecutionInput{
Identity: ident, Input: req.Input,
AgentVersion: req.AgentVersion,
Confirmation: req.Confirmation, Context: req.Context,
})
if res == nil || res.Termination == "" {
writeError(w, s.log, runLoadError(runErr))
return
}
writeJSON(w, http.StatusOK, buildRunResponse(res))
return
}
h := w.Header()
h.Set("Content-Type", "text/event-stream")
h.Set("Cache-Control", "no-store")
// Nginx and friends buffer proxied responses by default, which turns a
// stream into one very late blob. This is the header that turns that off.
h.Set("X-Accel-Buffering", "no")
w.WriteHeader(http.StatusOK)
flusher.Flush()
send := func(payload any) {
encoded, err := json.Marshal(payload)
if err != nil {
return
}
fmt.Fprintf(w, "data: %s\n\n", encoded)
flusher.Flush()
}
res, runErr := s.agents.RunAgent(r.Context(), ident, r.PathValue("id"), runtime.ExecutionInput{
Identity: ident,
Input: req.Input,
AgentVersion: req.AgentVersion,
Confirmation: req.Confirmation,
Context: req.Context,
OnDelta: func(d string) { send(map[string]string{"delta": d}) },
})
// A load failure has no run to report. It is sent as an event rather than a
// status code, because the status was already written when the stream
// opened — an SSE response cannot change its mind about being a 200.
if res == nil || res.Termination == "" {
var de *domain.Error
err := runLoadError(runErr)
if errors.As(err, &de) {
send(map[string]any{"error": map[string]string{"code": de.Code, "message": de.Message}})
} else {
send(map[string]any{"error": map[string]string{"code": "internal", "message": "internal error"}})
}
fmt.Fprint(w, "data: [DONE]\n\n")
flusher.Flush()
return
}
send(map[string]any{"run": buildRunResponse(res)})
fmt.Fprint(w, "data: [DONE]\n\n")
flusher.Flush()
}