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 } s.logUnsaved(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) } } // 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 } s.logUnsaved(ident, res) 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 } s.logUnsaved(ident, res) send(map[string]any{"run": buildRunResponse(res)}) fmt.Fprint(w, "data: [DONE]\n\n") flusher.Flush() }