package runtime import ( "context" "crypto/rand" "encoding/hex" "encoding/json" "errors" "fmt" "strings" "github.com/krow/krow-backend/go-api/internal/gateway" "github.com/krow/krow-backend/go-api/internal/knowledge" "github.com/krow/krow-backend/go-api/internal/tools" ) // ModelExecutor runs an agent against a model. // // This is the agent loop. It is spec-driven and there is exactly one of it: no // branch anywhere below asks which agent it is running. An agent's identity // reaches this code only as data — its instructions, its tier, its skills — // which is what I6 means in practice and what makes adding an agent a data // change rather than a deploy. // // The loop runs until the model stops asking for tools, or until a bound is // reached. Every exit is one of the six terminations. type ModelExecutor struct { // subagents resolves a spec's `subagents:` into runnable agents. nil means // delegation is off and a spec that declares subagents runs alone — see // delegate.go. subagents SubagentResolver gw gateway.Gateway sink Sink tools *tools.Registry // retriever is the knowledge layer, or nil for an agent platform with no // documents in it. Nil is a supported state rather than a broken one: every // agent built so far answers from the operational tables through tools, and // none of them needs a corpus. retriever Retriever } // Retriever is what the loop needs from the knowledge layer. // // An interface rather than the concrete type so the runtime does not import the // knowledge package's whole surface, and so a test can drive the loop with a // scripted corpus. Deliberately narrow: the loop retrieves, it does not ingest, // and it has no way to ask for anything other than the caller's own rows — // knowledge.Query requires a principal and this signature carries one. type Retriever interface { Retrieve(ctx context.Context, q knowledge.Query) (*knowledge.Results, error) } // WithRetriever attaches a knowledge layer to an executor. func (m *ModelExecutor) WithRetriever(r Retriever) *ModelExecutor { m.retriever = r return m } var _ AgentExecutor = (*ModelExecutor)(nil) // NewModelExecutor builds the loop over a model gateway. // // A nil sink is DiscardSink rather than a panic: a service wired without a // trajectory store should still answer, and losing the record is a worse // outcome than nothing but not one worth refusing a correct answer over. // A nil registry is an empty one: an agent that names no tools does not need // one, and a nil map dereference is a worse way to discover that than an agent // that simply has nothing to call. func NewModelExecutor(gw gateway.Gateway, sink Sink, reg *tools.Registry) *ModelExecutor { if sink == nil { sink = DiscardSink{} } if reg == nil { reg = tools.NewRegistry() } return &ModelExecutor{gw: gw, sink: sink, tools: reg} } // newRunID returns an opaque run identifier. // // Random rather than sequential: a run id appears in logs and in support // conversations, and a sequential one would leak how many runs a deployment // has served. func newRunID() string { var b [16]byte if _, err := rand.Read(b[:]); err != nil { // crypto/rand does not fail in practice; if it ever does, a run // without an id is still better than a run that refuses to start. return "run-unknown" } return "run_" + hex.EncodeToString(b[:]) } // ExecuteAgent runs one agent turn and returns a structured result. // // It never returns a bare error into user-facing text. Every exit is a // termination reason plus a trajectory, because §6 requires exactly one // termination per run and §10 requires user-facing text to be derived at the // surface layer rather than raised from here. func (m *ModelExecutor) ExecuteAgent(ctx context.Context, agent *Agent, input ExecutionInput) (*ExecutionResult, error) { tier, _ := gateway.ParseTier(agent.Reasoning) return m.executeWithLimits(ctx, agent, input, LimitsForTier(string(tier))) } // executeWithLimits is ExecuteAgent with the bounds supplied rather than // derived. // // The seam exists for two reasons and will earn its keep for the second. Today // it lets a test drive a real deadline instead of asserting on a counter. When // the spec gains a `limits:` block, that block resolves here and ExecuteAgent // stays the one-line default — so per-agent limits arrive without the loop // itself changing shape. func (m *ModelExecutor) executeWithLimits( ctx context.Context, agent *Agent, input ExecutionInput, limits Limits, ) (*ExecutionResult, error) { return m.executeRun(ctx, agent, input, limits, delegation{}) } // executeRun is the loop. `del` is what a SUBAGENT inherits from its parent — // a shared budget, a parent run id, and a depth — and its zero value is an // ordinary root run that owns its own budget and has no parent. // // One loop, not two. A delegated run is the same code on the same path; the // only things that differ are where its budget came from and what its // trajectory is linked to. §6 says there is exactly one loop, and a second one // for subagents would be the first place the two drifted apart. func (m *ModelExecutor) executeRun( ctx context.Context, agent *Agent, input ExecutionInput, limits Limits, del delegation, ) (*ExecutionResult, error) { tier, known := gateway.ParseTier(agent.Reasoning) // §6: a subagent SHARES the parent's budget and never gets a fresh one. // Delegating would otherwise be a way to buy more steps. budget := del.budget if budget == nil { budget = NewBudget(limits) } skillIDs := make([]string, len(agent.ResolvedSkills)) for i, s := range agent.ResolvedSkills { skillIDs[i] = s.ID } rec := NewRecorder(&Trajectory{ RunID: newRunID(), ParentRunID: del.parentRunID, OrgID: input.Identity.OrgID, UserID: input.Identity.UserID, AgentID: agent.ID, AgentVersion: agent.Version, Tier: string(tier), }) // A spec naming a tier the vocabulary does not have still runs, at the // default — but it is recorded, so a definition that has drifted is // visible in the trajectory rather than silently reinterpreted. if !known { rec.Error("runtime.unknown_tier", fmt.Sprintf("%q is not a reasoning mode; running at %s", agent.Reasoning, tier)) } // The deadline is the budget's, so an in-flight model call is torn down // rather than returning into a run that has already ended. runCtx, cancel := budget.Context(ctx) defer cancel() // Anything the caller wants on the record, before the run does anything. // A run that silently could not do what was asked of it is the failure // worth preventing here. for _, note := range input.Notes { rec.Error("runtime.note", note) } question := strings.TrimSpace(input.Input) if question == "" { return m.finish(ctx, rec, budget, TerminationToolFailure, agent, skillIDs, "", &RuntimeError{Code: "runtime.empty_input", Message: "a run needs a question"}) } rec.Message("user", question) // The tools this agent may use. Unknown names are recorded and dropped // rather than failing the run: §3 says an unknown tool fails validation at // *publish*, so one reaching run time means a tool was withdrawn under a // live spec — degrading is better than an outage, provided someone is told. toolDefs, unknown := m.toolsFor(agent) for _, name := range unknown { rec.Error("runtime.unknown_tool", fmt.Sprintf("%q is not a registered tool; it was not offered", name)) } // Subagents are offered as tools, because from this agent's side that is // exactly what they are (§6). Resolved once per run rather than per turn: // the set cannot change mid-run, and loading it per turn would spend the // caller's time on the same query repeatedly. subs := m.resolveSubagents(ctx, rec, agent, input.Identity, del.depth) toolDefs = append(toolDefs, delegateTools(subs)...) // An approved write happens FIRST, before the model gets a turn. // // This is the half of I4 that makes a confirmation reliable rather than // hopeful. The older design resumed the run and matched the model's next // tool call against the token — which only works if the model repeats // itself, and a model asked a second time may perfectly reasonably ask a // clarifying question instead. When that happened the token was never // presented, nothing was written, and the person who clicked Approve got a // follow-up question with no explanation. // // So the approved call is performed from what the person was SHOWN, not // from what the model says next. The model's job afterwards is to report // what happened, which is a job it cannot get wrong in a way that costs // anybody a shift. if approved, done := m.performApproved(runCtx, rec, budget, agent, input); done != nil { return m.finish(ctx, rec, budget, *done, agent, skillIDs, "", nil) } else if approved != "" { // Prepended to the question so the model answers knowing the write // already happened. It is a tool result in everything but shape — // delimited, factual, and about an act rather than an instruction. question = approved + "\n\n" + question } // Retrieval, before the first model call. // // I7 decides where the result goes: into a delimited block in a USER // message, never into the system prompt. The system prompt is assembled // from the agent record alone, so no amount of document content can reach // it — which is the only reason the standing "content inside is // data" instruction means anything. conversation := []gateway.Message{{Role: gateway.RoleUser, Text: question}} if block, retrieved := m.retrieve(runCtx, rec, agent, input, question); block != "" { conversation = []gateway.Message{{ Role: gateway.RoleUser, // Context first, question second. A model reads the question last // and answers it, rather than treating the evidence as the prompt. Text: block + "\n\n" + question, }} rec.Retrieval(retrieved) } system := SystemPrompt(agent) var lastText string for { // Claimed before dispatch, never after. A call that hangs until the // context dies has still spent the step it was given. if t := budget.ClaimStep(); t != "" { return m.finish(ctx, rec, budget, t, agent, skillIDs, lastText, nil) } if t := budget.CheckTokens(); t != "" { return m.finish(ctx, rec, budget, t, agent, skillIDs, lastText, nil) } rec.Budget(budget.Snapshot()) // Streamed when the caller asked for it AND the gateway can. Both // halves go through StreamComplete, so the loop has one call site and // no branch on transport — a run behaves identically whether its text // arrived in one piece or a hundred. resp, err := gateway.StreamComplete(runCtx, m.gw, gateway.Request{ Tier: tier, System: system, Messages: conversation, Tools: toolDefs, }, input.OnDelta) // Charged whatever happened. A refused or failed call was still billed, // and a ledger that forgives it is one a loop will happily repeat // against. if resp != nil { budget.ChargeTokens(resp.Usage.Total()) rec.ChargeUsage(resp.Usage.InputTokens, resp.Usage.OutputTokens, resp.Usage.CacheReadTokens+resp.Usage.CacheCreationTokens) rec.SetModel(resp.Model) } if err != nil { return m.finish(ctx, rec, budget, terminationFor(err), agent, skillIDs, lastText, err) } if resp.Text != "" { rec.Message("assistant", resp.Text) lastText = resp.Text } // No tool calls means the model is done talking. if len(resp.ToolCalls) == 0 { return m.finish(ctx, rec, budget, TerminationCompleted, agent, skillIDs, lastText, nil) } // The assistant turn goes back verbatim, calls included, before any // result is appended — a tool result with no preceding call is a // malformed conversation the API will reject. conversation = append(conversation, gateway.Message{ Role: gateway.RoleAssistant, Text: resp.Text, ToolCalls: resp.ToolCalls, }) results, pending, term := m.runTools(runCtx, rec, budget, agent, input, resp.ToolCalls, subs, del.depth) if term != "" { return m.finish(ctx, rec, budget, term, agent, skillIDs, lastText, nil) } // I4. A run that wants to write stops here and asks. It does not // continue with the reads it also made, does not summarise, and does // not get another turn to reconsider — the next thing that happens is a // person deciding, and the run resumes only if they say yes. if len(pending) > 0 { res, err := m.finish(ctx, rec, budget, TerminationConfirmationPending, agent, skillIDs, lastText, nil) res.Confirmations = pending return res, err } // Every result in ONE user turn. Splitting them is accepted and quietly // teaches the model to stop calling tools in parallel. conversation = append(conversation, gateway.Message{ Role: gateway.RoleUser, ToolResults: results, }) } } // performApproved carries out a write a person approved. // // Returns the sentence describing what happened, for the model to report from. // A run with no confirmation token does nothing here and returns "". // // A token that authorises nothing — unknown, expired, already spent, somebody // else's — is NOT an error and does not end the run. It is recorded and the run // continues, because the most common cause is a person clicking Approve twice, // and the honest response to that is to answer the question again rather than // to fail. func (m *ModelExecutor) performApproved( ctx context.Context, rec *Recorder, budget *Budget, agent *Agent, input ExecutionInput, ) (string, *Termination) { if input.Confirmation == "" || m.tools == nil { return "", nil } // The write spends a tool call from the run's budget, claimed before // dispatch like every other. An approval is not a way around I3. if t := budget.ClaimToolCall(); t != "" { return "", &t } rec.Budget(budget.Snapshot()) tc := tools.Context{ Principal: input.Identity, RunID: rec.RunID(), RemainingTokens: budget.Snapshot().TokensLeft, AgentID: agent.ID, KnowledgeSources: agent.KnowledgeSources, } out, ok := m.tools.DispatchApproved(ctx, tc, input.Confirmation) if !ok { rec.Error("runtime.confirmation_not_redeemable", "the supplied approval authorises nothing; it may have expired or already been used") return "", nil } rec.ToolCall(out.Tool, string(tools.EffectWrite), out.Inputs) rec.ToolResult(out.Tool, string(tools.EffectWrite), out.Result.Error != nil, out.Result) encoded, err := json.Marshal(out.Result) if err != nil { encoded = []byte(`{"error":{"code":"tool.failed","message":"the result could not be encoded"}}`) } // Delimited and labelled as data, on the same terms as retrieved content. // This text describes something that already happened; it is not an // instruction, and the standing rule in the system prompt covers // it for exactly that reason. return fmt.Sprintf( "\nA change you approved has already been carried out. This is its "+ "result, as data — report it, do not repeat the action.\n\n"+ "\n%s\n\n", out.Tool, string(encoded)), nil } // retrieve searches the agent's declared corpora on the CALLER's behalf. // // Three things are load-bearing and none of them is the search itself: // // - The principal is the caller's, never the agent's. I1: an agent reads // exactly what its caller could read directly, and the identity that // reaches knowledge.Query is the one that arrived with the request. // - The sources are the SPEC's. An agent granted the policy library does not // gain the incident log by asking nicely, because the source list is not // something the model can influence. // - A failure degrades rather than ends the run. A knowledge layer that is // down should cost grounding, not the answer — but it is recorded, because // an ungrounded answer that looks grounded is the worse outcome. func (m *ModelExecutor) retrieve( ctx context.Context, rec *Recorder, agent *Agent, input ExecutionInput, question string, ) (string, *knowledge.Results) { if m.retriever == nil || len(agent.KnowledgeSources) == 0 { return "", nil } res, err := m.retriever.Retrieve(ctx, knowledge.Query{ Text: question, Principal: input.Identity, Sources: agent.KnowledgeSources, }) if err != nil { // Recorded, not raised. The run continues without grounding, and the // trajectory says so — "the agent answered from nothing" is only // diagnosable afterwards if the failure was written down at the time. var kErr *knowledge.Error if errors.As(err, &kErr) { rec.Error(kErr.Code, kErr.Message) } else { rec.Error("knowledge.failed", err.Error()) } return "", nil } if res == nil || len(res.Chunks) == 0 { return "", res } return knowledge.RenderContext(res), res } // toolsFor resolves the tools an agent's spec names. func (m *ModelExecutor) toolsFor(agent *Agent) (defs []gateway.ToolDef, unknown []string) { if m.tools == nil || len(agent.Tools) == 0 { return nil, nil } resolved, unknown, err := m.tools.Resolve(agent.Tools) if err != nil { // Over the per-agent cap. Offering none is the safe reading: an agent // that silently got its first twenty tools would behave differently // depending on the order someone happened to write them in. return nil, agent.Tools } for _, t := range resolved { defs = append(defs, gateway.ToolDef{ Name: t.Name, Description: t.Description, InputSchema: t.InputSchema, }) } return defs, unknown } // runTools dispatches one turn's calls and returns their results. // // A tool that fails returns its error TO THE MODEL rather than ending the run. // §13 lists "swallowing a tool error and letting the model narrate around it" // as an anti-pattern — the fix is not to hide the failure but to hand it over // as a failure, so the model can say it could not look rather than inventing // what it would have found. // // The tool-call budget is claimed per call, before dispatch. Running out ends // the run: a model that has exhausted its calls cannot make progress, and // letting it continue would spend the remaining step budget on turns that can // only apologise. // // A write that needs approving comes back as a pending confirmation rather than // a result. Those are collected across the whole turn rather than returned at // the first one, so a person is asked about every write the model wanted in one // go instead of being walked through them one dialog at a time — and so that // the reads in the same turn, which are safe, still run and are still recorded. func (m *ModelExecutor) runTools( ctx context.Context, rec *Recorder, budget *Budget, agent *Agent, input ExecutionInput, calls []gateway.ToolCall, subs map[string]*Agent, depth int, ) (results []gateway.ToolResult, pending []*tools.Confirmation, term Termination) { results = make([]gateway.ToolResult, 0, len(calls)) for _, call := range calls { if t := budget.ClaimToolCall(); t != "" { return nil, nil, t } rec.Budget(budget.Snapshot()) // A delegation spends the same tool-call budget as any other call, // claimed above, and then spends the parent's remaining budget inside // the subagent. It is charged twice on purpose: once for asking, and // then for whatever the asking cost. if sub, isDelegation := subs[call.Name]; isDelegation { rec.ToolCall(call.Name, "delegate", json.RawMessage(call.Input)) answer, subPending := m.delegate(ctx, rec, budget, sub, input, json.RawMessage(call.Input), depth) encoded, err := json.Marshal(answer) if err != nil { encoded = []byte(`{"error":"the subagent's answer could not be encoded"}`) } rec.ToolResult(call.Name, "delegate", answer.Error != "", tools.Result{Data: answer}) // I4 survives delegation. A write a SUBAGENT wants approved is // still a write, and it stops this run the same way one from a // direct tool call does — see the pending check in the loop. if len(subPending) > 0 { for _, c := range subPending { rec.Confirmation(call.Name, c) } pending = append(pending, subPending...) continue } results = append(results, gateway.ToolResult{ CallID: call.ID, Content: string(encoded), IsError: answer.Error != "", }) continue } // The declared effect travels with the record. An eval asking "did this // run change anything" reads it from here rather than keeping its own // list of which tools write — a list that goes stale on the first tool // anybody adds. var effect string if t, ok := m.tools.Get(call.Name); ok { effect = string(t.Effect) } rec.ToolCall(call.Name, effect, json.RawMessage(call.Input)) res := m.tools.Dispatch(ctx, tools.Context{ Principal: input.Identity, RunID: rec.RunID(), RemainingTokens: budget.Snapshot().TokensLeft, Confirmation: input.Confirmation, // From the spec, never from the call. A model that asked to search // a corpus its agent was not granted is asking for a source list it // has no way to set. KnowledgeSources: agent.KnowledgeSources, }, call.Name, call.Input) // A pending confirmation never reaches the model. It is a question for // a person, and handing it back as a tool result would invite the model // to reason about it — to explain why it should be approved, or to try // a different tool that might not ask. Neither is its business. if res.Confirmation != nil { rec.Confirmation(call.Name, res.Confirmation) pending = append(pending, res.Confirmation) continue } rec.ToolResult(call.Name, effect, res.Error != nil, res) encoded, err := json.Marshal(res) if err != nil { encoded = []byte(`{"error":{"code":"tool.failed","message":"the result could not be encoded"}}`) } results = append(results, gateway.ToolResult{ CallID: call.ID, Content: string(encoded), IsError: res.Error != nil, }) } return results, pending, "" } // terminationFor maps a failure to the reason a run ends with. // // The mapping matters more than it looks: Deadline and BudgetExceeded are // different questions to an operator ("too slow" versus "too expensive"), and // 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 // the trajectory. func terminationFor(err error) Termination { var gwErr *gateway.Error if !errors.As(err, &gwErr) { return TerminationToolFailure } switch gwErr.Code { case gateway.CodeRefused: return TerminationRefused case gateway.CodeTimeout: return TerminationDeadline default: return TerminationToolFailure } } // finish closes the trajectory, persists it, and builds the caller's result. // // Persistence uses the *caller's* context, not the run's: the run context is // cancelled at the deadline, and a run that ended by running out of time is // exactly the one whose record is most worth keeping. func (m *ModelExecutor) finish( ctx context.Context, rec *Recorder, budget *Budget, term Termination, agent *Agent, skillIDs []string, output string, cause error, ) (*ExecutionResult, error) { if cause != nil { var gwErr *gateway.Error if errors.As(cause, &gwErr) { rec.Error(gwErr.Code, gwErr.Message) } else { rec.Error("runtime.failed", cause.Error()) } } rec.Budget(budget.Snapshot()) traj := rec.Finish(term) // A sink that fails must not fail the run — the answer was already // produced. It is recorded in the trajectory we could not save, which is // the best available place for it. if err := m.sink.Save(ctx, traj); err != nil { rec.Error("runtime.trajectory_unsaved", err.Error()) } // Delegated runs are written AFTER this one, because parent_run_id is a // foreign key and a subagent finishes first. Each child arrives with its // own descendants already ordered behind it, so one pass here writes a // whole tree parent-first. for _, child := range rec.Children() { if err := m.sink.Save(ctx, child); err != nil { rec.Error("runtime.subrun_unsaved", fmt.Sprintf("%s: %s", child.RunID, err.Error())) } } res := &ExecutionResult{ Success: term == TerminationCompleted, Output: output, AgentID: agent.ID, AgentVersion: agent.Version, ResolvedSkills: skillIDs, RunID: traj.RunID, Termination: term, Usage: traj.Usage, } if term == TerminationCompleted { return res, nil } // A bounded run is not an exception. The caller gets a result carrying the // reason; the error exists so a Go caller that ignores the result still // notices, and it is structured so the surface layer derives the wording. rtErr := &RuntimeError{ Code: "runtime." + strings.ToLower(string(term)), Message: terminationMessage(term), Target: agent.ID, Cause: cause, } res.Error = rtErr return res, rtErr } // terminationMessage is the internal explanation for a termination. Not // user-facing copy — §10 puts that at the surface layer, which is free to say // something kinder using the code. func terminationMessage(t Termination) string { switch t { case TerminationBudgetExceeded: return "the run reached its budget before finishing" case TerminationDeadline: return "the run reached its deadline before finishing" case TerminationRefused: return "the model declined to answer" case TerminationConfirmationPending: return "the run is waiting on a confirmation" case TerminationToolFailure: return "the run failed" default: return string(t) } } // SystemPrompt assembles an agent's system prompt from its spec. // // I7 is the whole design of this function. Retrieved document text, tool // results and user messages are all untrusted, and none of them are reachable // from here: it reads the agent record and nothing else. When Phase 2 adds // retrieval, the retrieved chunks go into a delimited block in a *user* // message — not into this string — and the standing instruction below is what // makes that delimiter mean something. func SystemPrompt(agent *Agent) string { var b strings.Builder b.WriteString("You are ") b.WriteString(agent.Name) if agent.Description != "" { b.WriteString(", ") b.WriteString(agent.Description) } b.WriteString(".\n\n") if instructions := strings.TrimSpace(agent.Instructions); instructions != "" { b.WriteString(instructions) b.WriteString("\n\n") } if len(agent.Pages) > 0 { b.WriteString("You answer on: ") b.WriteString(strings.Join(agent.Pages, ", ")) b.WriteString(". Anywhere else, say plainly that you do not cover it.\n\n") } // Stated even when nothing was retrieved, because the boundary has to be // established before content arrives rather than alongside it. // // The sentence comes from the knowledge package, beside the renderer that // emits the fence. A prompt promising while the renderer wrote // would be a defence that had quietly stopped existing, and two // copies of a string in two packages is exactly how that happens. b.WriteString(knowledge.ContextInstruction) b.WriteString("\n\n") b.WriteString("State a figure only where the records you were given show it. " + "When you cannot answer from them, say so rather than estimating.") return b.String() }