Skip to content

Commit afe6c80

Browse files
committed
Add parent steering for running workers
1 parent fdcfafb commit afe6c80

10 files changed

Lines changed: 423 additions & 33 deletions

File tree

src/agent/fleet-verbs-mount.test.ts

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
/**
2-
* Primary createAgentToolset mounts the six fleet verbs beside task /
3-
* search_agents / read_agent_trace when subAgent (with the shared TUI
4-
* sessions store) is wired. Leaves / no-subAgent toolsets stay without them.
2+
* Primary createAgentToolset mounts the fleet verbs beside task / search_agents /
3+
* read_agent_trace when subAgent (with the shared TUI sessions store) is wired.
4+
* Leaves / no-subAgent toolsets stay without them.
55
*/
66
import { mkdtempSync } from "node:fs";
77
import { tmpdir } from "node:os";
@@ -17,10 +17,11 @@ const FLEET_VERBS = [
1717
"resume_agent",
1818
"interrupt_agent",
1919
"followup_task",
20+
"send_input",
2021
] as const;
2122

2223
describe("primary fleet verb mount", () => {
23-
test("createAgentToolset registers the six fleet verbs when subAgent + sessions are set", async () => {
24+
test("createAgentToolset registers the fleet verbs when subAgent + sessions are set", async () => {
2425
const cwd = mkdtempSync(join(tmpdir(), "corbits-fleet-mount-"));
2526
const { createAgentToolset } = await import("./tools.js");
2627
const permissionGate = {

src/agent/tool-search.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ export const CORE_TOOL_NAMES: readonly string[] = [
4949
"resume_agent",
5050
"interrupt_agent",
5151
"followup_task",
52+
"send_input",
5253
];
5354

5455
const ORCHESTRATOR_ONLY_TOOL_NAMES: readonly string[] = [
@@ -60,6 +61,7 @@ const ORCHESTRATOR_ONLY_TOOL_NAMES: readonly string[] = [
6061
"resume_agent",
6162
"interrupt_agent",
6263
"followup_task",
64+
"send_input",
6365
];
6466

6567
// Session-start facts that gate a core tool's advertisement. Each must be

src/agent/tools.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ import {
4949
createResumeAgentTool,
5050
createInterruptAgentTool,
5151
createFollowupTaskTool,
52+
createSendInputTool,
5253
} from "../subagent/lifecycle-tools.js";
5354
import { parseManageTasksArgs } from "./tasks.js";
5455
import { createListDirTool } from "../util/list-dir.js";
@@ -348,6 +349,8 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
348349
createResumeAgentTool({ sessions: fleetSessions }),
349350
createInterruptAgentTool({ sessions: fleetSessions }),
350351
createFollowupTaskTool({ sessions: fleetSessions }),
352+
// Tier 1: primary may target any worker — omit authority (unrestricted).
353+
createSendInputTool({ sessions: fleetSessions }),
351354
);
352355
}
353356
}

src/subagent/agent-fleet.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -449,10 +449,11 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
449449
// Keep the session open after a clean completion, and hand the
450450
// store a bounded close for close_agent to call later.
451451
persist: true,
452-
onAgentReady: ({ close, interrupt, followup }) => {
452+
onAgentReady: ({ close, interrupt, followup, deliver }) => {
453453
deps.sessions.registerClose(session.id, close);
454454
deps.sessions.registerInterrupt(session.id, interrupt);
455455
deps.sessions.registerFollowup(session.id, followup);
456+
deps.sessions.registerDeliver(session.id, deliver);
456457
deps.sessions.markRunning(session.id);
457458
},
458459
};

src/subagent/authority.ts

Lines changed: 5 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -11,10 +11,11 @@
1111
* this same gate).
1212
* - assertCanTargetAgent: a Tier 2 nested orchestrator may act only on its
1313
* own descendants, never a sibling or anything above it in the tree.
14-
* Tier 1 (the primary orchestrator) may target anyone. Callers pass the
15-
* live fleet as a flat list of {id, parentSessionId} nodes — the same
16-
* shape SubAgentSessionStore already tracks — so no parallel tree
17-
* structure is needed.
14+
* Tier 1 (the primary orchestrator) may target anyone. Addressing verbs
15+
* such as read_agent_trace and send_input call this at their handler
16+
* boundary. Callers pass the live fleet as a flat list of {id,
17+
* parentSessionId} nodes — the same shape SubAgentSessionStore already
18+
* tracks — so no parallel tree structure is needed.
1819
*/
1920

2021
import type { SubagentTier } from "../agent/directors/types.js";
@@ -89,15 +90,6 @@ function isDescendant(
8990
}
9091

9192
/**
92-
* SEAM, NOT YET A LIVE GATE: this function has no production call site today.
93-
* No verb in this codebase currently lets one live agent target another
94-
* (`task` only spawns; it never addresses an existing session), so the
95-
* subtree rule below is exercised only by authority.test.ts — it is not
96-
* enforced at runtime yet. It exists now so future verbs that make one
97-
* agent addressable by another can call it from day one instead of
98-
* inventing their own check. Until one of those wires a call site here, do
99-
* not describe this rule as enforced; only assertTierMayMountFleetVerb is.
100-
*
10193
* Authority rule (root owns its tree; a child manages only its own
10294
* descendants): throws unless `actor` is Tier 1, or `targetId` is `actor.id`
10395
* itself, or a descendant of `actor.id` in `nodes`. A Tier 3 leaf holds no

src/subagent/lifecycle-tools.test.ts

Lines changed: 172 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import {
55
createResumeAgentTool,
66
createInterruptAgentTool,
77
createFollowupTaskTool,
8+
createSendInputTool,
89
} from "./lifecycle-tools.js";
910
import { createSubAgentSessionStore } from "./session-store.js";
1011

@@ -13,7 +14,8 @@ async function callTool(
1314
| ReturnType<typeof createCloseAgentTool>
1415
| ReturnType<typeof createResumeAgentTool>
1516
| ReturnType<typeof createInterruptAgentTool>
16-
| ReturnType<typeof createFollowupTaskTool>,
17+
| ReturnType<typeof createFollowupTaskTool>
18+
| ReturnType<typeof createSendInputTool>,
1719
args: Record<string, unknown>,
1820
): Promise<Record<string, unknown>> {
1921
if (tool.kind !== "full") throw new Error(`expected full tool, got ${tool.kind}`);
@@ -258,3 +260,172 @@ describe("interrupt_agent / followup_task", () => {
258260
expect(followupErr.isError).toBe(true);
259261
});
260262
});
263+
264+
describe("send_input", () => {
265+
test("soft-delivers a durable message to a running worker without awaiting a reply", async () => {
266+
const sessions = createSubAgentSessionStore();
267+
const worker = sessions.start({
268+
description: "worker",
269+
agentId: "a",
270+
brief: "b",
271+
retained: true,
272+
});
273+
sessions.markRunning(worker.id);
274+
const delivered: string[] = [];
275+
sessions.registerDeliver(worker.id, (message) => {
276+
delivered.push(message);
277+
});
278+
279+
const sendInput = createSendInputTool({ sessions });
280+
const result = await callTool(sendInput, {
281+
target: worker.id,
282+
message: "stop and inspect line 4",
283+
});
284+
285+
expect(result).toEqual({ agent_id: worker.id, status: "running" });
286+
expect(delivered).toEqual(["stop and inspect line 4"]);
287+
expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running");
288+
});
289+
290+
test("interrupts and queues a next-turn message without awaiting the reply", async () => {
291+
const sessions = createSubAgentSessionStore();
292+
const worker = sessions.start({
293+
description: "worker",
294+
agentId: "a",
295+
brief: "b",
296+
retained: true,
297+
});
298+
sessions.markRunning(worker.id);
299+
let interrupted = false;
300+
let followupStarted = false;
301+
sessions.registerInterrupt(worker.id, () => {
302+
interrupted = true;
303+
});
304+
sessions.registerFollowup(worker.id, async (message) => {
305+
followupStarted = true;
306+
expect(message).toBe("drop the broad refactor and patch only the test");
307+
await new Promise((resolve) => setTimeout(resolve, 25));
308+
return "queued turn finished";
309+
});
310+
sessions.registerDeliver(worker.id, () => {
311+
throw new Error("interrupt:true should not soft-deliver");
312+
});
313+
314+
const sendInput = createSendInputTool({ sessions });
315+
const result = await callTool(sendInput, {
316+
target: worker.id,
317+
message: "drop the broad refactor and patch only the test",
318+
interrupt: true,
319+
});
320+
321+
expect(result).toEqual({ agent_id: worker.id, status: "interrupted" });
322+
expect(interrupted).toBe(true);
323+
expect(followupStarted).toBe(true);
324+
expect(sessions.get(worker.id)?.lifecycleStatus).toBe("interrupted");
325+
});
326+
327+
test("fails closed when interrupt:true cannot queue the followup", async () => {
328+
const sessions = createSubAgentSessionStore();
329+
const worker = sessions.start({
330+
description: "worker",
331+
agentId: "a",
332+
brief: "b",
333+
retained: true,
334+
});
335+
sessions.markRunning(worker.id);
336+
let interrupted = false;
337+
sessions.registerInterrupt(worker.id, () => {
338+
interrupted = true;
339+
});
340+
341+
const sendInput = createSendInputTool({ sessions });
342+
if (sendInput.kind !== "full") throw new Error("expected full tool");
343+
const result = await sendInput.handler(
344+
{
345+
id: "missing-followup",
346+
name: "send_input",
347+
arguments: { target: worker.id, message: "steer after interrupt", interrupt: true },
348+
},
349+
new AbortController().signal,
350+
);
351+
352+
expect(result.isError).toBe(true);
353+
expect(interrupted).toBe(false);
354+
expect(sessions.get(worker.id)?.lifecycleStatus).toBe("running");
355+
});
356+
357+
test("rejects empty and oversize messages", async () => {
358+
const sessions = createSubAgentSessionStore();
359+
const worker = sessions.start({
360+
description: "worker",
361+
agentId: "a",
362+
brief: "b",
363+
retained: true,
364+
});
365+
sessions.markRunning(worker.id);
366+
sessions.registerDeliver(worker.id, () => {});
367+
const sendInput = createSendInputTool({ sessions });
368+
369+
if (sendInput.kind !== "full") throw new Error("expected full tool");
370+
const empty = await sendInput.handler(
371+
{ id: "empty", name: "send_input", arguments: { target: worker.id, message: " " } },
372+
new AbortController().signal,
373+
);
374+
expect(empty.isError).toBe(true);
375+
376+
const oversize = await sendInput.handler(
377+
{
378+
id: "big",
379+
name: "send_input",
380+
arguments: { target: worker.id, message: "x".repeat(24_001) },
381+
},
382+
new AbortController().signal,
383+
);
384+
expect(oversize.isError).toBe(true);
385+
});
386+
387+
test("enforces nested orchestrator descendant authority", async () => {
388+
const sessions = createSubAgentSessionStore();
389+
const nested = sessions.start({
390+
id: "nested",
391+
description: "nested",
392+
agentId: "a",
393+
brief: "b",
394+
});
395+
const child = sessions.start({
396+
id: "child",
397+
description: "child",
398+
agentId: "a",
399+
brief: "b",
400+
parentSessionId: nested.id,
401+
});
402+
const sibling = sessions.start({
403+
id: "sibling",
404+
description: "sibling",
405+
agentId: "a",
406+
brief: "b",
407+
});
408+
for (const session of [nested, child, sibling]) {
409+
sessions.markRunning(session.id);
410+
sessions.registerDeliver(session.id, () => {});
411+
}
412+
const sendInput = createSendInputTool({
413+
sessions,
414+
authority: {
415+
actorId: nested.id,
416+
tier: "nested-orchestrator",
417+
getNodes: () => sessions.list(),
418+
},
419+
});
420+
421+
const ok = await callTool(sendInput, { target: child.id, message: "continue" });
422+
expect(ok.status).toBe("running");
423+
424+
if (sendInput.kind !== "full") throw new Error("expected full tool");
425+
const denied = await sendInput.handler(
426+
{ id: "denied", name: "send_input", arguments: { target: sibling.id, message: "continue" } },
427+
new AbortController().signal,
428+
);
429+
expect(denied.isError).toBe(true);
430+
});
431+
});

0 commit comments

Comments
 (0)