Files
krow_backend/go-api/internal/httpserver/runs.go
Aravind a372340281
Some checks failed
CI / test (push) Failing after 4m41s
CI / fixture (push) Failing after 9s
add the greeting msg
2026-10-05 16:20:17 +05:30

503 lines
21 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/gateway"
"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)
s.logGatewayFailure(ident, res)
writeJSON(w, http.StatusOK, buildRunResponse(res))
}
// logUnsaved is the operator's record of a run whose trajectory did not
// persist. The run itself already answered; §6 says the trajectory is not
// optional telemetry, so losing one is an error even when nothing else went
// wrong, and it carries every field §10 asks a log line to carry.
func (s *Server) logUnsaved(ident authctx.Identity, res *runtime.ExecutionResult) {
for _, detail := range res.Unsaved {
s.log.Error("trajectory unsaved",
"run_id", res.RunID, "tenant_id", ident.OrgID,
"agent_key", res.AgentID, "agent_version", res.AgentVersion,
"detail", detail)
}
}
// logGatewayFailure is the operator's record of the model provider not
// answering, and it exists because the reader of the message cannot report what
// the message does not say.
//
// Four of the five gateway faults need an administrator and will fail
// identically on every retry — a rejected credential, a model id this
// deployment cannot use, no key at all, and a 4xx from the endpoint. The person
// in the chat panel is told, correctly, that retrying will not help; but
// nothing until now told the side that CAN fix it. A deployment failing every
// run produced a stream of 200s and no error line, so the only record of which
// fault it was lived in a trajectory somebody had to know to go and read.
//
// Logged at Error because that is what it is: on a rate limit it is a capacity
// decision worth seeing, and on the other four it is an outage. It carries the
// gateway's code and status and NOT the message — §10 keeps model and document
// text out of the log store, and the code is the part that is actionable
// anyway.
func (s *Server) logGatewayFailure(ident authctx.Identity, res *runtime.ExecutionResult) {
if res.Termination != runtime.TerminationGatewayFailure {
return
}
// Empty rather than invented when the cause did not survive: "which fault
// was it" is the whole point of this line, and a guessed answer to it is
// worse than a visible gap.
code, status := "", 0
var gwErr *gateway.Error
if errors.As(res.Error, &gwErr) {
code, status = gwErr.Code, gwErr.Status
}
s.log.Error("gateway failure",
"run_id", res.RunID, "tenant_id", ident.OrgID,
"agent_key", res.AgentID, "agent_version", res.AgentVersion,
"gateway_code", code, "gateway_status", status)
}
// buildRunResponse turns a runtime result into the client's shape.
//
// Every termination answers 200. That looks wrong at first and is not: the
// 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, res.Error)
}
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.
// `cause` is the run's own error, carried so GatewayFailure can say which
// gateway failure it was. Every other termination ignores it.
func terminationMessage(t runtime.Termination, cause error) string {
switch t {
case runtime.TerminationCompleted:
return ""
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:
return gatewayFailureMessage(cause)
default:
return "The agent did not finish."
}
}
// gatewayFailureMessage tells a GatewayFailure apart from the four others it
// used to be indistinguishable from.
//
// GatewayFailure is everything the gateway can raise except a refusal and a
// timeout, which have terminations of their own. That is five different faults,
// and the one sentence they all produced was "usually it is busy — wait a
// minute and ask again".
//
// For a rate limit that is true and useful. For the other three it is advice
// that CANNOT work: a credential the provider rejected, a model id that is not
// on the configured endpoint, and no key at all will each fail identically on
// every retry, forever. Telling somebody to wait a minute for a misconfigured
// deployment sends them round a loop with no exit, and hides an operator
// problem behind what reads as a transient one — the reader retries instead of
// reporting it, so nobody with access to fix it ever hears.
//
// So each says what it is, and only the two that clear on their own invite a
// retry. The wording stays free of vendor names and status codes: the person
// reading it cannot act on "429 from the model endpoint", and the code is in
// the trajectory for the person who can.
func gatewayFailureMessage(cause error) string {
var gwErr *gateway.Error
if !errors.As(cause, &gwErr) {
// No gateway error to read — either the cause was lost or something
// non-gateway reached this termination. The old sentence is still the
// best guess, so it is what an unknown falls back to.
return "The model behind this agent did not answer — usually it is busy. " +
"Wait a minute and ask again. Nothing was changed."
}
switch gwErr.Code {
case gateway.CodeRateLimited:
return "The model behind this agent is busy right now. " +
"Wait a minute and ask again. Nothing was changed."
case gateway.CodeUpstream:
if gwErr.Status >= 500 {
return "The model behind this agent is having trouble. " +
"Try again in a few minutes. Nothing was changed."
}
return "The model behind this agent could not be reached, and retrying is " +
"unlikely to help. This needs an administrator. Nothing was changed."
case gateway.CodeUnauthorized:
return "This deployment's model credentials were rejected, so the agent " +
"cannot answer. Retrying will not help — it needs an administrator. " +
"Nothing was changed."
case gateway.CodeInvalidRequest:
return "The agent is pointed at a model this deployment cannot use, so it " +
"cannot answer. Retrying will not help — it needs an administrator. " +
"Nothing was changed."
case gateway.CodeNotConfigured:
return "No model is configured for this deployment, so the agent cannot " +
"answer. It needs an administrator. Nothing was changed."
default:
return "The model behind this agent did not answer — usually it is busy. " +
"Wait a minute and ask again. Nothing was changed."
}
}
// runLoadError maps a pre-run failure onto the API's error vocabulary.
//
// These are the errors from LoadExecutableAgent, raised before any run began —
// 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)
s.logGatewayFailure(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)
s.logGatewayFailure(ident, res)
send(map[string]any{"run": buildRunResponse(res)})
fmt.Fprint(w, "data: [DONE]\n\n")
flusher.Flush()
}