Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/agent/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
run: runSubAgent,
sessions: fleetSessions,
fleetRecords,
...(sa.useWorktree !== undefined ? { useWorktree: sa.useWorktree } : {}),
...(sa.onEvent !== undefined ? { onEvent: sa.onEvent } : {}),
...(sa.onProgress !== undefined ? { onProgress: sa.onProgress } : {}),
...(sa.settings !== undefined ? { settings: sa.settings } : {}),
Expand Down
83 changes: 83 additions & 0 deletions src/subagent/agent-fleet.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -922,3 +922,86 @@ describe("list_agents", () => {
gate.resolve({ report: "done" });
});
});

describe("spawn_agent parity with task", () => {
test("uses the parent tool call id as the session id", async () => {
const deps = makeDeps(async () => ({ report: "done" }));
const spawn = createSpawnAgentTool(deps);
if (spawn.kind !== "full") throw new Error("expected full tool");
const result = await spawn.handler(
{
id: "call-fixed-id",
name: "spawn_agent",
arguments: { description: "job", prompt: "do it", intent: "explore" },
},
new AbortController().signal,
);
const content = typeof result.content === "string" ? result.content : "";
expect(JSON.parse(content).agent_id).toBe("call-fixed-id");
expect(deps.sessions.get("call-fixed-id")).toBeDefined();
});

test("refuses skywalker as a spawned worker", async () => {
const deps = makeDeps(async () => ({ report: "no" }));
const spawn = createSpawnAgentTool(deps);
const raw = await callToolRaw(spawn, {
description: "nope",
prompt: "do it",
agent: "skywalker",
});
expect(raw.isError).toBe(true);
expect(raw.content).toContain("skywalker is the primary session identity");
});

test("rejects a child outside this director allowlist", async () => {
const deps = makeDeps(async () => ({ report: "no" }));
deps.spawnAllowlist = ["intern", "explorer", "critic"];
const spawn = createSpawnAgentTool(deps);
const raw = await callToolRaw(spawn, {
description: "build",
prompt: "ship it",
agent: "builder",
});
expect(raw.isError).toBe(true);
expect(raw.content).toContain("allowlist");
expect(raw.content).toContain("builder");
});

test("a maySpawn director is launched as an orchestrator with nestedDispatch", async () => {
const captured: RunSubAgentParams[] = [];
const deps = makeDeps(async (params) => {
captured.push(params);
return { report: "ok" };
});
const spawn = createSpawnAgentTool(deps);
await callTool(spawn, {
description: "arch",
prompt: "judge this",
agent: "greybeard",
});
await new Promise((resolve) => setTimeout(resolve, 20));
expect(captured).toHaveLength(1);
expect(captured[0]!.orchestrator).toBe(true);
expect(captured[0]!.orchestratorTier).toBe("nested-orchestrator");
expect(captured[0]!.nestedDispatch).toBeDefined();
expect(captured[0]!.nestedDispatch?.spawnAllowlist).toEqual(["intern", "explorer", "critic"]);
});

test("allowOrchestrator false strips nested spawn even for maySpawn directors", async () => {
const captured: RunSubAgentParams[] = [];
const deps = makeDeps(async (params) => {
captured.push(params);
return { report: "ok" };
});
deps.allowOrchestrator = false;
const spawn = createSpawnAgentTool(deps);
await callTool(spawn, {
description: "arch",
prompt: "judge this",
agent: "greybeard",
});
await new Promise((resolve) => setTimeout(resolve, 20));
expect(captured[0]!.orchestrator).toBeUndefined();
expect(captured[0]!.nestedDispatch).toBeUndefined();
});
});
153 changes: 144 additions & 9 deletions src/subagent/agent-fleet.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,21 +31,25 @@
*
* Argument shape intentionally mirrors `task()`'s (description/prompt/
* context/goals/intent/success_criteria/do_not/report_focus) so a
* caller can swap one for the other. Scope is deliberately narrower than
* `task()` for this first cut: only closed-director dispatch (`agent=` a
* director id, or `intent=`) is supported — no custom AgentProfile lookup,
* no nested orchestration, no re-dispatch ledger. Those remain `task()`-only
* for now; nothing here stops adding them later.
* caller can swap one for the other. Closed-director dispatch also carries
* task()'s isolation and spawn-matrix: worktree cwd, parent allowlist,
* maySpawn nestedDispatch, and deadline. Custom AgentProfile lookup and the
* re-dispatch ledger remain task()-only until task becomes a thin wrapper.
*
*/

import { join } from "node:path";

import { tool } from "@intx/agent";
import type { AgentTool } from "@intx/agent";
import { type } from "arktype";
import type { ToolDefinition, ToolResult } from "@intx/types/runtime";
import type { ReactorEmittedEvent } from "@intx/inference";
import { getLogger } from "@intx/log";

import { LOG_NAMESPACE_ROOT } from "../branding.js";
import type { ProviderCatalogEntry } from "../config/index.js";
import { generateSessionId } from "../session/index.js";
import {
isDirectorId,
packageToCapabilities,
Expand All @@ -61,13 +65,18 @@ import { isCodexProviderName } from "../config/codex-providers.js";
import { buildDispatchBrief, type TaskIntent } from "./report.js";
import type { SubAgentSessionStore } from "./session-store.js";
import type {
NestedDispatchDeps,
RunSubAgentParams,
RunSubAgentResult,
SubAgentProvider,
SubAgentSandboxDeps,
} from "./types.js";
import { cleanupSubAgentWorktree, createSubAgentWorktree, WorktreeError } from "./worktree.js";
import { NOOP_TELEMETRY, type Telemetry } from "../telemetry/index.js";
import { classifyAgentName } from "../telemetry/classify.js";
import type { DirectorPackage } from "../agent/directors/types.js";

const log = getLogger([LOG_NAMESPACE_ROOT, "subagent", "agent-fleet"]);

/** Terminal (or running) record for one spawned agent, keyed by agent id. */
interface FleetRecord {
Expand Down Expand Up @@ -346,6 +355,14 @@ export type AgentFleetDeps = SubAgentSandboxDeps & {
* Omit on the primary session — its children are top-level.
*/
parentSessionId?: string;
/** When set, only these director ids may be spawned. */
spawnAllowlist?: readonly string[];
/** When false, maySpawn directors cannot remount fleet verbs. Defaults true. */
allowOrchestrator?: boolean;
/** Isolate each spawn in a git worktree branched from dispatcher HEAD. */
useWorktree?: boolean;
/** Optional wall-clock budget (ms) forwarded to runSubAgent. */
deadlineMs?: number;
settings?: Settings | (() => Settings | undefined);
catalog?: readonly ProviderCatalogEntry[] | (() => readonly ProviderCatalogEntry[]);
onEvent?: (event: ReactorEmittedEvent) => void;
Expand Down Expand Up @@ -373,6 +390,7 @@ export function resolveDirectorDispatch(
systemPromptRole: string;
capabilities: ReturnType<typeof packageToCapabilities>;
roleDefault: ReturnType<typeof defaultEffortForDirector>;
pkg: DirectorPackage;
}
| { ok: false; error: string } {
if (agentId !== undefined && agentId.length > 0) {
Expand All @@ -391,6 +409,7 @@ export function resolveDirectorDispatch(
systemPromptRole: formatDirectorSystemPrompt(pkg),
capabilities: packageToCapabilities(pkg),
roleDefault: defaultEffortForDirector(pkg),
pkg,
};
}
if (intent !== undefined) {
Expand All @@ -403,6 +422,7 @@ export function resolveDirectorDispatch(
systemPromptRole: formatDirectorSystemPrompt(pkg),
capabilities: packageToCapabilities(pkg),
roleDefault: defaultEffortForDirector(pkg),
pkg,
};
}
return {
Expand Down Expand Up @@ -451,12 +471,34 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {

const resolved = resolveDirectorDispatch(agentId, intent);
if (!resolved.ok) return fleetResult(call.id, resolved.error);
if (agentId === "skywalker" || resolved.directorId === "skywalker") {
return fleetResult(
call.id,
"Error: skywalker is the primary session identity, not a spawned worker. Pass spawn_agent(agent=…) for a specialist (builder, explorer, counsel, critic, …).",
);
}
if (deps.spawnAllowlist !== undefined && deps.spawnAllowlist.length > 0) {
if (!deps.spawnAllowlist.includes(resolved.directorId)) {
return fleetResult(
call.id,
`Error: spawn of "${resolved.directorId}" is outside this director's allowlist. Allowed: ${deps.spawnAllowlist.join(", ")}.`,
);
}
}

const settings = deps.settings !== undefined ? resolveDep(deps.settings) : undefined;

const orchestrator = resolved.pkg.spawn.maySpawn === true && deps.allowOrchestrator !== false;
const nestedSpawnAllowlist =
orchestrator &&
resolved.pkg.spawn.allowlist !== undefined &&
resolved.pkg.spawn.allowlist.length > 0
? resolved.pkg.spawn.allowlist
: undefined;

let provider: SubAgentProvider = resolveDep(deps.provider);
const effort = resolveEffortForRole({
orchestrator: false,
orchestrator,
roleDefault: resolved.roleDefault,
...(provider.reasoningEffort !== undefined
? { parentEffort: provider.reasoningEffort }
Expand All @@ -478,6 +520,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
});

const session = deps.sessions.start({
id: call.id,
description,
agentId: resolved.directorId,
brief,
Expand All @@ -503,6 +546,72 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
};

const catalog = deps.catalog !== undefined ? resolveDep(deps.catalog) : undefined;
let worktreeCwd: string | undefined;
let worktreeStashBaseline: readonly string[] | null = [];
let worktreeHeadAtCreate: string | undefined;
if (deps.useWorktree === true) {
const worktreePath = join(deps.getWorkdirBase(), "worktrees", generateSessionId());
try {
const worktree = await createSubAgentWorktree(deps.cwd, worktreePath);
worktreeCwd = worktree.path;
worktreeStashBaseline = worktree.stashBaseline;
worktreeHeadAtCreate = worktree.headAtCreate;
} catch (err) {
const message =
err instanceof WorktreeError
? err.message
: `sub-agent worktree setup failed: ${err instanceof Error ? err.message : String(err)}`;
deps.fleetRecords.reject(session.id, message);
deps.sessions.fail(session.id, message);
return fleetResult(call.id, `Error: ${message}`);
}
}

const nestedDispatch: NestedDispatchDeps | undefined = orchestrator
? {
permissionGate: deps.permissionGate,
...(deps.inheritMcpTools !== undefined
? { inheritMcpTools: deps.inheritMcpTools }
: {}),
...(deps.shellTimeout !== undefined ? { shellTimeout: deps.shellTimeout } : {}),
...(deps.shellEnv !== undefined ? { shellEnv: deps.shellEnv } : {}),
...(deps.extraToolPlugins !== undefined
? { extraToolPlugins: deps.extraToolPlugins }
: {}),
...(deps.getBlobReader !== undefined ? { getBlobReader: deps.getBlobReader } : {}),
getWorkdirBase: deps.getWorkdirBase,
provider: deps.provider,
...(deps.onEvent !== undefined ? { onEvent: deps.onEvent } : {}),
...(deps.onProgress !== undefined ? { onProgress: deps.onProgress } : {}),
sessions: deps.sessions,
...(settings !== undefined ? { settings } : {}),
...(catalog !== undefined ? { catalog } : {}),
parentSessionId: session.id,
...(deps.useWorktree !== undefined ? { useWorktree: deps.useWorktree } : {}),
...(nestedSpawnAllowlist !== undefined ? { spawnAllowlist: nestedSpawnAllowlist } : {}),
}
: undefined;

// Aligns with run.ts: (persist && turnSucceeded) || interruptedKeepAlive.
// When true, the worktree stays until close_agent / eviction calls the
// wrapped close below; otherwise the run's finally reclaims it immediately.
let keepWorktreeAlive = false;
const reclaimWorktree = async (): Promise<void> => {
if (worktreeCwd === undefined) return;
const path = worktreeCwd;
worktreeCwd = undefined;
try {
await cleanupSubAgentWorktree(deps.cwd, path, {
stashBaseline: worktreeStashBaseline,
...(worktreeHeadAtCreate !== undefined ? { headAtCreate: worktreeHeadAtCreate } : {}),
});
} catch (err: unknown) {
log.error("spawn_agent worktree cleanup failed: {error}", {
error: err instanceof Error ? err.message : String(err),
});
}
};

const params: RunSubAgentParams = {
// Name the trace directory after the session-store id so the
// descendant-scoping check behind read_agent_trace can resolve this
Expand All @@ -514,7 +623,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
...(deps.shellEnv !== undefined ? { shellEnv: deps.shellEnv } : {}),
...(deps.extraToolPlugins !== undefined ? { extraToolPlugins: deps.extraToolPlugins } : {}),
...(deps.getBlobReader !== undefined ? { getBlobReader: deps.getBlobReader } : {}),
cwd: deps.cwd,
cwd: worktreeCwd ?? deps.cwd,
workdirBase: deps.getWorkdirBase(),
provider,
...(settings !== undefined ? { settings } : {}),
Expand All @@ -533,11 +642,33 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
...(resolved.capabilities !== undefined ? { capabilities: resolved.capabilities } : {}),
systemPromptRole: resolved.systemPromptRole,
directorId: resolved.directorId,
...(orchestrator
? {
orchestrator: true,
orchestratorTier: resolved.pkg.tier,
nestedDispatch: nestedDispatch!,
}
: {}),
...(deps.deadlineMs !== undefined ? { deadlineMs: deps.deadlineMs } : {}),
tier: resolved.pkg.tier,
...(resolved.pkg.reportContract?.outputType !== undefined
? { reportType: resolved.pkg.reportContract.outputType }
: {}),
// Keep the session open after a clean completion, and hand the
// store a bounded close for close_agent to call later.
// Worktree cleanup is deferred until that close when the session
// stays alive for followup (agentRetained / interrupt keep-alive) —
// matching run.ts's persisting gate so followup_task does not hit a
// removed cwd.
persist: true,
onAgentReady: ({ close, interrupt, followup, deliver }) => {
deps.sessions.registerClose(session.id, close);
deps.sessions.registerClose(session.id, async (deadlineMs) => {
try {
await close(deadlineMs);
} finally {
await reclaimWorktree();
}
});
deps.sessions.registerInterrupt(session.id, interrupt);
deps.sessions.registerFollowup(session.id, followup);
deps.sessions.registerDeliver(session.id, deliver);
Expand All @@ -561,6 +692,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
// "completed" status. Still terminalize fleetRecords so a waiter
// that never saw interrupt_agent (or raced it) cannot hang.
if (result.interrupted === true) {
keepWorktreeAlive = true;
deps.fleetRecords.interrupt(session.id, result.report);
return;
}
Expand All @@ -575,8 +707,10 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
// its agent first, so the store must not treat it as resumable
// just because retained:true was requested at spawn.
// complete() no-ops when status is already cancelled.
const agentRetained = result.agentRetained === true;
if (agentRetained) keepWorktreeAlive = true;
deps.sessions.complete(session.id, result.report, {
agentRetained: result.agentRetained === true,
agentRetained,
...(result.stopReason !== undefined ? { stopReason: result.stopReason } : {}),
});
})
Expand All @@ -594,6 +728,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
status: deps.sessions.get(session.id)?.status ?? "completed",
duration_ms: Date.now() - startedAt,
});
if (!keepWorktreeAlive) void reclaimWorktree();
});

return fleetResult(call.id, JSON.stringify({ agent_id: session.id, status: "running" }));
Expand Down
3 changes: 3 additions & 0 deletions src/subagent/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -598,7 +598,10 @@ export async function runSubAgent(params: RunSubAgentParams): Promise<RunSubAgen
telemetry: liveTelemetry,
sessions: fleetSessions,
fleetRecords,
allowOrchestrator: false,
...(params.id !== undefined ? { parentSessionId: params.id } : {}),
...(nd.useWorktree !== undefined ? { useWorktree: nd.useWorktree } : {}),
...(nd.spawnAllowlist !== undefined ? { spawnAllowlist: nd.spawnAllowlist } : {}),
...(nd.onEvent !== undefined ? { onEvent: nd.onEvent } : {}),
...(nd.onProgress !== undefined ? { onProgress: nd.onProgress } : {}),
...(nd.settings !== undefined ? { settings: nd.settings } : {}),
Expand Down
Loading
Loading