From 6849363a37893f5d2e1c589607e5345f765b7705 Mon Sep 17 00:00:00 2001 From: Suriyakumarvijayanayagam Date: Sat, 29 Aug 2026 11:59:20 +0530 Subject: [PATCH] Implement delegation, so an agent's subagents are more than decoration MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit krow-workforce-agent has declared five subagents since it was written and answered every question by itself. Everything for §6 existed except the delegation: the parser read `subagents:`, runtime.Agent 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 already described sharing a budget with subagents. Nothing called any of it. A subagent is offered to the parent's model as a tool, because §6 says that is what delegation is from the parent's side. Three rules are enforced rather than assumed, each with a test that fails if it stops holding: I1 The subagent runs as the ORIGINAL caller. It cannot read anything the person could not read directly. §6 It SHARES the parent's budget. The test sets MaxSteps to 1, spends it in the parent, and asserts the child terminates BudgetExceeded — an assertion that only passes when the budget is shared, and that a fresh budget would quietly turn green. §3 Depth is capped at 2. At the cap no subagent is loaded or offered, so a cycle reaching run time is bounded rather than unbounded. I4 survives too: a write a SUBAGENT wants approved still stops the whole run and asks a person, rather than being performed because it happened one level down. Two bugs found by running it rather than by reading it: - delegate() read the error before the result. finish returns a non-nil error for every termination that is not Completed, INCLUDING ConfirmationPending — which is not a failure but a run that stopped to ask a question. Reading the error first discarded the result and with it the confirmation, so a subagent's write silently never happened and nobody was asked. - Delegated trajectories were never persisted at all. parent_run_id is a foreign key and a subagent finishes BEFORE the run that delegated to it, so every child insert named a parent row that did not exist yet. The database refused it; finish deliberately does not fail a run over a sink error; and the entry recording that the trajectory could not be saved was itself in the trajectory that was not saved. Children are now buffered and written by finish after the parent's own row, each arriving with its descendants already ordered behind it, so one pass writes a whole tree parent-first. The regression test asserts on save ORDER, because a MemorySink has no foreign key and will pass either way. Verified end to end against a live model: an agent with no tools of its own and one subagent produced delegation-probe run=run_16622d7de6 parent=(root) talent-pool-agent run=run_64160684b1 parent=run_16622d7de6 with the subagent's answer reaching the parent's model. Full suite green, only TestLive* skipped. Not addressed: §3's publish-time cycle detection, which needs the whole agent set in hand. The depth cap is what holds without it, and is the half that matters at run time. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g --- go-api/internal/runtime/delegate.go | 267 +++++++++++++++++ go-api/internal/runtime/delegate_test.go | 362 +++++++++++++++++++++++ go-api/internal/runtime/loop.go | 83 +++++- go-api/internal/runtime/trajectory.go | 30 ++ go-api/internal/runtime/wire.go | 8 +- 5 files changed, 747 insertions(+), 3 deletions(-) create mode 100644 go-api/internal/runtime/delegate.go create mode 100644 go-api/internal/runtime/delegate_test.go 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)) }