diff --git a/go-api/internal/runtime/delegate.go b/go-api/internal/runtime/delegate.go new file mode 100644 index 0000000..1d1d9c8 --- /dev/null +++ b/go-api/internal/runtime/delegate.go @@ -0,0 +1,267 @@ +package runtime + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/gateway" + "github.com/krow/krow-backend/go-api/internal/tools" +) + +/* ── Delegation ───────────────────────────────────────────────────────────── + +§6: "Delegation is a tool call from the parent's perspective. Subagent runs get +their own trajectory, linked by parent_run_id." + +Everything for this existed except the delegation. The parser read `subagents:`, +the runtime type carried them, the loader populated them, agent_runs had a +parent_run_id column with a self-reference and a no-self-parent constraint, and +budget.go's comments described sharing a budget with subagents. Nothing called +any of it, so krow-workforce-agent declared five subagents and answered every +question alone. + +Three rules from §3 and §6 are load-bearing here, and each is enforced below +rather than assumed: + + I1 A subagent runs as the ORIGINAL caller. It never receives a widened + principal, so it can read exactly what the person could read directly. + §6 A subagent SHARES the parent's budget. It never gets a fresh one, or a + run could buy itself unlimited steps by delegating in a loop. + §3 Depth is capped. The cap is what makes a cycle in the spec graph a + bounded waste rather than an unbounded one. +──────────────────────────────────────────────────────────────────────────── */ + +// MaxDelegationDepth is §3's cap. +// +// The root run is depth 0, so a run at depth 2 is offered no subagents at all +// and delegation stops there. Publish-time cycle detection is the other half +// of §3 and is not here — this is what holds when a cycle reaches run time +// anyway, which is the case worth defending against. +const MaxDelegationDepth = 2 + +// SubagentResolver loads a subagent ready to run, with its own skills and +// dependencies already resolved. +// +// An interface, satisfied by *Loader, so the loop keeps depending on behaviour +// rather than on the repository — the same reason Retriever is one. +type SubagentResolver interface { + LoadExecutableAgent(ctx context.Context, ident authctx.Identity, idOrDefID string) (*Agent, error) +} + +// WithSubagents attaches the resolver that turns a spec's `subagents:` into +// agents this loop can actually run. Without it, delegation is silently off — +// which is the behaviour every deployment had until now. +func (m *ModelExecutor) WithSubagents(r SubagentResolver) *ModelExecutor { + m.subagents = r + return m +} + +// delegation is what a subagent run inherits from its parent. +// +// The zero value is a root run: its own budget, no parent, depth 0. +type delegation struct { + budget *Budget + parentRunID string + depth int +} + +// delegationToolPrefix marks a tool call as a delegation rather than a tool. +const delegationToolPrefix = "ask_" + +// delegationToolName is the name a subagent is offered to the model under. +// +// Hyphens become underscores because agent ids are kebab-case and tool names +// across this registry are snake_case; a model offered both conventions at +// once picks badly. +func delegationToolName(agentID string) string { + return delegationToolPrefix + strings.ReplaceAll(agentID, "-", "_") +} + +// resolveSubagents loads the agents this one may delegate to, keyed by the +// tool name each is offered under. +// +// Every reason a subagent cannot be offered is RECORDED rather than silently +// dropped. A spec that names five subagents and gets three is a spec somebody +// needs to fix, and the trajectory is where they will look. +func (m *ModelExecutor) resolveSubagents( + ctx context.Context, rec *Recorder, agent *Agent, ident authctx.Identity, depth int, +) map[string]*Agent { + + if m.subagents == nil || len(agent.Subagents) == 0 { + return nil + } + if depth >= MaxDelegationDepth { + rec.Error("runtime.delegation_depth", fmt.Sprintf( + "at depth %d of %d; %d subagent(s) were not offered", + depth, MaxDelegationDepth, len(agent.Subagents))) + return nil + } + + out := make(map[string]*Agent, len(agent.Subagents)) + for _, id := range agent.Subagents { + if id == agent.ID { + // The database forbids a run being its own parent; refusing it + // here means the model is never offered the call in the first + // place. + rec.Error("runtime.subagent_self", fmt.Sprintf("%q names itself as a subagent", id)) + continue + } + name := delegationToolName(id) + if m.tools != nil { + if _, taken := m.tools.Get(name); taken { + rec.Error("runtime.subagent_shadowed", fmt.Sprintf( + "subagent %q would be offered as %q, which is a registered tool", id, name)) + continue + } + } + sub, err := m.subagents.LoadExecutableAgent(ctx, ident, id) + if err != nil || sub == nil { + // Same reasoning as an unknown tool: §3 says this fails at + // publish, so reaching run time means the spec changed underneath + // a live agent. Degrade and say so. + rec.Error("runtime.unknown_subagent", fmt.Sprintf( + "%q could not be loaded and was not offered", id)) + continue + } + out[name] = sub + } + return out +} + +// delegateTools describes each subagent to the model as a tool it can call. +// +// The description is the subagent's own, because that is what was written for +// a model to read. A parent choosing between five subagents is doing the same +// job the router does when choosing between five tools, and it needs the same +// quality of description to do it. +func delegateTools(subs map[string]*Agent) []gateway.ToolDef { + if len(subs) == 0 { + return nil + } + defs := make([]gateway.ToolDef, 0, len(subs)) + for name, sub := range subs { + desc := strings.TrimSpace(sub.Description) + if t := strings.TrimSpace(sub.Trigger); t != "" { + desc = strings.TrimSpace(desc + " " + t) + } + if desc == "" { + desc = fmt.Sprintf("Ask the %s agent.", sub.Name) + } + defs = append(defs, gateway.ToolDef{ + Name: name, + Description: fmt.Sprintf( + "Delegate to %s and return its answer. %s "+ + "Ask one self-contained question: this agent cannot see your conversation.", + sub.Name, desc), + InputSchema: map[string]any{ + "type": "object", + "properties": map[string]any{ + "question": map[string]any{ + "type": "string", + "description": "The question to ask, complete on its own. " + + "Include any names, dates or ids it needs.", + }, + }, + "required": []string{"question"}, + "additionalProperties": false, + }, + }) + } + return defs +} + +// delegationRequest is the one argument a delegation takes. +type delegationRequest struct { + Question string `json:"question"` +} + +// delegationAnswer is what the parent's model receives back. +// +// Structured, not prose: §4 says handlers return data and formatting is the +// model's job, and a delegation is a tool call from the parent's side. RunID +// travels with it so a bad answer inside a delegated branch can be found. +type delegationAnswer struct { + Agent string `json:"agent"` + RunID string `json:"runId"` + Termination string `json:"termination"` + Answer string `json:"answer,omitempty"` + Error string `json:"error,omitempty"` +} + +// delegate runs one subagent and returns its answer, plus anything it wants a +// person to approve. +// +// The subagent gets the caller's identity and the PARENT'S budget object, so +// its steps, tool calls, tokens and deadline all come out of the same +// allowance. It does not get the parent's confirmation token: a token +// authorises one specific write that one person was shown, and handing it down +// would let a different tool spend it. +func (m *ModelExecutor) delegate( + ctx context.Context, rec *Recorder, budget *Budget, sub *Agent, + input ExecutionInput, raw json.RawMessage, depth int, +) (answer delegationAnswer, pending []*tools.Confirmation) { + + var req delegationRequest + if err := json.Unmarshal(raw, &req); err != nil || strings.TrimSpace(req.Question) == "" { + return delegationAnswer{ + Agent: sub.ID, Termination: string(TerminationToolFailure), + Error: "a delegation needs a question", + }, nil + } + + // The subagent writes into a buffer rather than straight to the sink: its + // row cannot be inserted until this run's row exists (parent_run_id is a + // foreign key). The buffer is handed to the recorder and flushed by finish + // once this run has been written. + buffered := *m + collected := &MemorySink{} + buffered.sink = collected + m = &buffered + + res, err := m.executeRun(ctx, sub, ExecutionInput{ + Identity: input.Identity, // I1 — the caller, never widened + Input: req.Question, + }, LimitsForTier(sub.Reasoning), delegation{ + budget: budget, // §6 — shared, never fresh + parentRunID: rec.RunID(), + depth: depth + 1, + }) + + // The RESULT is what matters, not the error beside it. finish returns a + // non-nil error for every termination that is not Completed — including + // ConfirmationPending, which is not a failure at all but a run that + // stopped to ask a person something. Reading the error first and + // discarding the result loses that question, and the write it was + // guarding silently never happens. + // Whatever happened, keep what the subagent recorded. A run that failed is + // the one somebody will want to read. + rec.AddChildren(collected.Runs) + + if res == nil { + msg := "the subagent returned nothing" + if err != nil { + msg = err.Error() + } + return delegationAnswer{ + Agent: sub.ID, Termination: string(TerminationToolFailure), Error: msg, + }, nil + } + + out := delegationAnswer{ + Agent: sub.ID, + RunID: res.RunID, + Termination: string(res.Termination), + Answer: res.Output, + } + switch { + case res.Error != nil: + out.Error = res.Error.Error() + case err != nil && res.Termination != TerminationCompleted && + res.Termination != TerminationConfirmationPending: + out.Error = err.Error() + } + return out, res.Confirmations +} diff --git a/go-api/internal/runtime/delegate_test.go b/go-api/internal/runtime/delegate_test.go new file mode 100644 index 0000000..5f8caba --- /dev/null +++ b/go-api/internal/runtime/delegate_test.go @@ -0,0 +1,362 @@ +package runtime + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + "github.com/krow/krow-backend/go-api/internal/authctx" + "github.com/krow/krow-backend/go-api/internal/gateway" + "github.com/krow/krow-backend/go-api/internal/tools" +) + +// fakeSubagents resolves subagents from a map, and remembers who asked — which +// is how I1 is checked below. +type fakeSubagents struct { + agents map[string]*Agent + asked []string + ident authctx.Identity +} + +func (f *fakeSubagents) LoadExecutableAgent(_ context.Context, ident authctx.Identity, id string) (*Agent, error) { + f.asked = append(f.asked, id) + f.ident = ident + if a, ok := f.agents[id]; ok { + return a, nil + } + return nil, fmt.Errorf("no agent %q", id) +} + +func childAgent() *Agent { + return &Agent{ + ID: "talent-pool-agent", Name: "Talent Pool Agent", Version: 1, + Description: "Who is available in the pool.", + Reasoning: "balanced", + Pages: []string{"talent-pool"}, + } +} + +// parentWith returns an agent declaring the given subagents, and a resolver +// that can supply the child. +func parentWith(subs ...string) (*Agent, *fakeSubagents) { + p := testAgent() + p.Subagents = subs + return p, &fakeSubagents{agents: map[string]*Agent{"talent-pool-agent": childAgent()}} +} + +func delegationCall(id, tool, question string) gateway.ToolCall { + return gateway.ToolCall{ + ID: id, Name: tool, + Input: json.RawMessage(fmt.Sprintf(`{"question":%q}`, question)), + } +} + +// A subagent is offered to the model as a tool. Before this, `subagents:` +// parsed, loaded, and was dropped — the model was never told the agent existed. +func TestDelegationOffersSubagentsAsTools(t *testing.T) { + parent, res := parentWith("talent-pool-agent") + gw := &scriptedGateway{} + exec := NewModelExecutor(gw, &MemorySink{}, tools.NewRegistry()).WithSubagents(res) + + if _, err := exec.ExecuteAgent(context.Background(), parent, testInput("who is free?")); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(gw.seen) == 0 { + t.Fatal("the model was never called") + } + var names []string + for _, td := range gw.seen[0].Tools { + names = append(names, td.Name) + } + want := delegationToolName("talent-pool-agent") + if len(names) != 1 || names[0] != want { + t.Fatalf("offered tools = %v, want exactly [%s]", names, want) + } + if len(res.asked) != 1 { + t.Errorf("the resolver was asked %d times, want 1 — subagents resolve once per run", len(res.asked)) + } +} + +// Without a resolver, delegation is off and the agent runs alone. This is the +// behaviour every deployment had, and it must stay a quiet degrade rather than +// a failure. +func TestDelegationIsOffWithoutAResolver(t *testing.T) { + parent, _ := parentWith("talent-pool-agent") + gw := &scriptedGateway{} + exec := NewModelExecutor(gw, &MemorySink{}, tools.NewRegistry()) + + res, err := exec.ExecuteAgent(context.Background(), parent, testInput("who is free?")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Termination != TerminationCompleted { + t.Errorf("Termination = %q, want Completed", res.Termination) + } + if len(gw.seen[0].Tools) != 0 { + t.Errorf("offered %d tools with no resolver, want 0", len(gw.seen[0].Tools)) + } +} + +// The subagent runs, its answer reaches the parent's model, and it writes its +// OWN trajectory linked by parent_run_id (§6). +func TestDelegationRunsTheSubagentAndLinksItsTrajectory(t *testing.T) { + parent, resolver := parentWith("talent-pool-agent") + tool := delegationToolName("talent-pool-agent") + + gw := &scriptedGateway{steps: []*gateway.Response{ + // parent asks + {ToolCalls: []gateway.ToolCall{delegationCall("c1", tool, "who is available?")}, + StopReason: "tool_use", Model: "fake-model"}, + // child answers + {Text: "Five workers are available.", StopReason: "end_turn", Model: "fake-model"}, + // parent answers from it + {Text: "Five are free this week.", StopReason: "end_turn", Model: "fake-model"}, + }} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, tools.NewRegistry()).WithSubagents(resolver) + + res, err := exec.ExecuteAgent(context.Background(), parent, testInput("who is free?")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Termination != TerminationCompleted || res.Output != "Five are free this week." { + t.Fatalf("Termination=%q Output=%q", res.Termination, res.Output) + } + + if len(sink.Runs) != 2 { + t.Fatalf("%d trajectories saved, want 2 — the subagent gets its own", len(sink.Runs)) + } + var child, root *Trajectory + for _, tr := range sink.Runs { + if tr.AgentID == "talent-pool-agent" { + child = tr + } else { + root = tr + } + } + if child == nil || root == nil { + t.Fatal("expected one trajectory per agent") + } + // ORDER MATTERS, and a MemorySink will not tell you so on its own. + // agent_runs.parent_run_id is a foreign key, and a subagent finishes + // before the run that delegated to it — so saving in completion order + // makes every child insert name a parent row that does not exist yet. The + // database refuses it, finish does not fail a run over a sink error, and + // every delegated trajectory disappears without trace. Parent first. + if sink.Runs[0].AgentID != root.AgentID { + t.Errorf("saved %q first, want the parent %q — a child written before its "+ + "parent violates the parent_run_id foreign key and is silently dropped", + sink.Runs[0].AgentID, root.AgentID) + } + if child.ParentRunID != root.RunID { + t.Errorf("child.ParentRunID = %q, want the parent's run id %q", child.ParentRunID, root.RunID) + } + if root.ParentRunID != "" { + t.Errorf("the root run has ParentRunID %q, want empty", root.ParentRunID) + } + + // The parent's model must actually have received the child's answer. + last := gw.seen[len(gw.seen)-1] + var sawAnswer bool + for _, msg := range last.Messages { + for _, tr := range msg.ToolResults { + if strings.Contains(tr.Content, "Five workers are available.") { + sawAnswer = true + } + } + } + if !sawAnswer { + t.Error("the subagent's answer never reached the parent's model") + } +} + +// §6, and the reason delegation cannot be a way to buy more budget. +// +// MaxSteps is 1. The parent spends it, then delegates. If the subagent got a +// FRESH budget it would have a step of its own and answer; sharing the +// parent's means it has nothing left and terminates BudgetExceeded. The +// assertion is on the child's termination, which differs between the two +// designs and cannot be produced by accident. +func TestDelegationSharesTheParentsBudget(t *testing.T) { + parent, resolver := parentWith("talent-pool-agent") + tool := delegationToolName("talent-pool-agent") + + gw := &scriptedGateway{steps: []*gateway.Response{ + {ToolCalls: []gateway.ToolCall{delegationCall("c1", tool, "who is available?")}, + StopReason: "tool_use", Model: "fake-model"}, + {Text: "the child should never get this far", StopReason: "end_turn", Model: "fake-model"}, + }} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, tools.NewRegistry()).WithSubagents(resolver) + + // The parent exhausts the budget too and finish returns an error with its + // result; that is expected here and not what this test is about. + _, _ = exec.executeRun(context.Background(), parent, testInput("who is free?"), + Limits{MaxSteps: 1, MaxToolCalls: 5, MaxTokens: 100_000, Deadline: 30 * time.Second}, + delegation{}) + + var child *Trajectory + for _, tr := range sink.Runs { + if tr.AgentID == "talent-pool-agent" { + child = tr + } + } + if child == nil { + t.Fatal("the subagent never ran") + } + if child.Termination != TerminationBudgetExceeded { + t.Errorf("child Termination = %q, want BudgetExceeded — it was given a fresh budget "+ + "instead of sharing its parent's, which §6 forbids", child.Termination) + } +} + +// I1: a subagent executes as the ORIGINAL caller and never a widened one. +func TestDelegationRunsAsTheOriginalCaller(t *testing.T) { + parent, resolver := parentWith("talent-pool-agent") + gw := &scriptedGateway{} + exec := NewModelExecutor(gw, &MemorySink{}, tools.NewRegistry()).WithSubagents(resolver) + + in := testInput("who is free?") + if _, err := exec.ExecuteAgent(context.Background(), parent, in); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if resolver.ident.UserID != in.Identity.UserID || resolver.ident.OrgID != in.Identity.OrgID { + t.Errorf("subagent resolved as %+v, want the caller %+v", resolver.ident, in.Identity) + } +} + +// §3's depth cap. At the cap, no subagent is offered at all, so a cycle in the +// spec graph is bounded rather than unbounded. +func TestDelegationCapsDepth(t *testing.T) { + parent, resolver := parentWith("talent-pool-agent") + gw := &scriptedGateway{} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, tools.NewRegistry()).WithSubagents(resolver) + + _, err := exec.executeRun(context.Background(), parent, testInput("who is free?"), + LimitsForTier("balanced"), delegation{depth: MaxDelegationDepth}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(gw.seen[0].Tools) != 0 { + t.Errorf("offered %d tools at depth %d, want 0", len(gw.seen[0].Tools), MaxDelegationDepth) + } + if len(resolver.asked) != 0 { + t.Errorf("resolved %d subagents at the cap, want 0 — the cap should short-circuit "+ + "before loading anything", len(resolver.asked)) + } + // And it must be visible, not silent. + if !hasError(sink.Last(), "runtime.delegation_depth") { + t.Error("hitting the depth cap was not recorded in the trajectory") + } +} + +// An agent naming itself is refused before the model is offered the call: the +// database has a no-self-parent constraint, and a run that tried would fail on +// insert rather than on anything legible. +func TestDelegationRefusesSelfReference(t *testing.T) { + parent := testAgent() + parent.Subagents = []string{parent.ID} + resolver := &fakeSubagents{agents: map[string]*Agent{}} + + gw := &scriptedGateway{} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, tools.NewRegistry()).WithSubagents(resolver) + + if _, err := exec.ExecuteAgent(context.Background(), parent, testInput("hello")); err != nil { + t.Fatalf("unexpected error: %v", err) + } + if len(gw.seen[0].Tools) != 0 { + t.Errorf("a self-referencing subagent was offered as a tool") + } + if !hasError(sink.Last(), "runtime.subagent_self") { + t.Error("the self-reference was not recorded") + } +} + +// A subagent that cannot be loaded degrades rather than failing the run — the +// same rule as an unknown tool — but is recorded so somebody can fix the spec. +func TestDelegationRecordsAnUnloadableSubagent(t *testing.T) { + parent := testAgent() + parent.Subagents = []string{"no-such-agent"} + resolver := &fakeSubagents{agents: map[string]*Agent{}} + + gw := &scriptedGateway{} + sink := &MemorySink{} + exec := NewModelExecutor(gw, sink, tools.NewRegistry()).WithSubagents(resolver) + + res, err := exec.ExecuteAgent(context.Background(), parent, testInput("hello")) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if res.Termination != TerminationCompleted { + t.Errorf("Termination = %q, want Completed — an unloadable subagent degrades", res.Termination) + } + if !hasError(sink.Last(), "runtime.unknown_subagent") { + t.Error("the unloadable subagent was not recorded") + } +} + +// I4 survives delegation. A write a SUBAGENT wants approved still stops +// everything and asks a person — it does not get performed because it happened +// one level down. +func TestSubagentConfirmationStopsTheParentRun(t *testing.T) { + writeTool := tools.Tool{ + Name: "assign_worker", + Description: "Assigns a worker to a shift, which is a real change to a real rota.", + InputSchema: map[string]any{"type": "object"}, + Effect: tools.EffectWrite, + Confirm: func(context.Context, tools.Context, json.RawMessage) (*tools.Confirmation, *tools.Result) { + return &tools.Confirmation{Token: "tok_1", Tool: "assign_worker", Title: "Assign Maya to Bar"}, nil + }, + Handler: func(context.Context, tools.Context, json.RawMessage) tools.Result { + t.Error("the write ran without approval") + return tools.OK(map[string]any{}) + }, + } + reg := tools.NewRegistry() + reg.MustRegister(writeTool) + + child := childAgent() + child.Tools = []string{"assign_worker"} + parent := testAgent() + parent.Subagents = []string{child.ID} + resolver := &fakeSubagents{agents: map[string]*Agent{child.ID: child}} + + tool := delegationToolName(child.ID) + gw := &scriptedGateway{steps: []*gateway.Response{ + {ToolCalls: []gateway.ToolCall{delegationCall("c1", tool, "assign Maya")}, + StopReason: "tool_use", Model: "fake-model"}, + {ToolCalls: []gateway.ToolCall{{ID: "c2", Name: "assign_worker", Input: json.RawMessage(`{}`)}}, + StopReason: "tool_use", Model: "fake-model"}, + }} + exec := NewModelExecutor(gw, &MemorySink{}, reg).WithSubagents(resolver) + + // The error beside the result is how the loop reports every termination + // that is not Completed; ConfirmationPending is not a failure and the + // existing confirmation tests ignore it the same way. + res, _ := exec.ExecuteAgent(context.Background(), parent, testInput("assign someone")) + if res.Termination != TerminationConfirmationPending { + t.Fatalf("Termination = %q, want ConfirmationPending — a subagent's write must "+ + "still stop and ask", res.Termination) + } + if len(res.Confirmations) != 1 || res.Confirmations[0].Tool != "assign_worker" { + t.Fatalf("Confirmations = %+v, want the subagent's pending write", res.Confirmations) + } +} + +// hasError reports whether a trajectory recorded an entry with the given code. +func hasError(tr *Trajectory, code string) bool { + if tr == nil { + return false + } + for _, e := range tr.Entries { + if strings.Contains(fmt.Sprint(e), code) { + return true + } + } + return false +} diff --git a/go-api/internal/runtime/loop.go b/go-api/internal/runtime/loop.go index dc5eee1..91e47d3 100644 --- a/go-api/internal/runtime/loop.go +++ b/go-api/internal/runtime/loop.go @@ -25,6 +25,11 @@ import ( // 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 @@ -109,9 +114,29 @@ func (m *ModelExecutor) ExecuteAgent(ctx context.Context, agent *Agent, input Ex // 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) - budget := NewBudget(limits) + + // §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 { @@ -120,6 +145,7 @@ func (m *ModelExecutor) executeWithLimits( rec := NewRecorder(&Trajectory{ RunID: newRunID(), + ParentRunID: del.parentRunID, OrgID: input.Identity.OrgID, UserID: input.Identity.UserID, AgentID: agent.ID, @@ -163,6 +189,13 @@ func (m *ModelExecutor) executeWithLimits( 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 @@ -259,7 +292,7 @@ func (m *ModelExecutor) executeWithLimits( Role: gateway.RoleAssistant, Text: resp.Text, ToolCalls: resp.ToolCalls, }) - results, pending, term := m.runTools(runCtx, rec, budget, agent, input, 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) } @@ -425,6 +458,7 @@ func (m *ModelExecutor) toolsFor(agent *Agent) (defs []gateway.ToolDef, unknown 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)) @@ -434,6 +468,40 @@ func (m *ModelExecutor) runTools( } 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 @@ -535,6 +603,17 @@ func (m *ModelExecutor) finish( 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, diff --git a/go-api/internal/runtime/trajectory.go b/go-api/internal/runtime/trajectory.go index bee854f..ec44f71 100644 --- a/go-api/internal/runtime/trajectory.go +++ b/go-api/internal/runtime/trajectory.go @@ -147,6 +147,36 @@ func (m *MemorySink) Last() *Trajectory { type Recorder struct { mu sync.Mutex t *Trajectory + + // children are trajectories of runs delegated beneath this one, already + // in parent-before-child order. + // + // They are held rather than saved as they finish because agent_runs. + // parent_run_id is a FOREIGN KEY, and a subagent always finishes before + // the run that delegated to it. Saving in completion order means every + // child insert names a parent row that does not exist yet, which the + // database refuses — and finish deliberately does not fail a run over a + // sink error, so the whole thing vanished silently. The trajectory + // recording that it could not be saved was itself the one not saved. + children []*Trajectory +} + +// AddChildren queues trajectories delegated beneath this run, to be written +// after this run's own row exists. +func (r *Recorder) AddChildren(ts []*Trajectory) { + if len(ts) == 0 { + return + } + r.mu.Lock() + defer r.mu.Unlock() + r.children = append(r.children, ts...) +} + +// Children returns the queued delegated trajectories. +func (r *Recorder) Children() []*Trajectory { + r.mu.Lock() + defer r.mu.Unlock() + return append([]*Trajectory(nil), r.children...) } // NewRecorder begins recording a run. diff --git a/go-api/internal/runtime/wire.go b/go-api/internal/runtime/wire.go index dabae32..0379d55 100644 --- a/go-api/internal/runtime/wire.go +++ b/go-api/internal/runtime/wire.go @@ -24,7 +24,13 @@ func NewModelEngine(db repo.Querier, cfg config.Config) *Engine { gw := gateway.NewAnthropic(gateway.FromConfig(cfg.Model)) retriever := knowledge.NewRetriever(db, NewEmbedder(cfg)) exec := NewModelExecutor(gw, NewPostgresSink(db), DefaultTools(db, retriever)). - WithRetriever(retriever) + WithRetriever(retriever). + // Without this, a spec's `subagents:` parses, loads, and is then + // dropped — which is how krow-workforce-agent came to declare five + // subagents and answer every question by itself. The resolver is the + // same Loader the engine uses, so a subagent is loaded exactly the way + // a directly-invoked agent is: same skills, same dependency checks. + WithSubagents(NewLoader(db)) return NewEngine(db, WithAgentExecutor(exec)) }