Skip to content

Commit 3edfe00

Browse files
Make director task a spawn plus wait wrapper (#685)
* Make director task() a spawn_agent plus wait_agents wrapper Closed-director task() was a second full spawn engine. When a session store is present it now starts the worker through spawn_agent and blocks on wait_agents, so one mailbox owns completion. Custom AgentProfile lookup still uses the legacy await-run path. * Preserve task cancel, auth, and abort contracts through the fleet wrapper Director task() now routes through spawn_agent + wait_agents, which was misclassifying AbortError as failed:aborted, dropping Re-authenticate auth wording, and leaving the child running when the parent tool aborted. Map those outcomes back to the legacy fused-task parent contract and prettier the inherited agent-progress tip.
1 parent 09be3f9 commit 3edfe00

5 files changed

Lines changed: 307 additions & 18 deletions

File tree

src/agent/tools.ts

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -282,6 +282,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
282282
const orchestratorTools: AgentTool[] = [];
283283
if (subAgentsEnabled && args.subAgent !== undefined) {
284284
const sa = args.subAgent;
285+
const fleetRecords = sa.sessions !== undefined ? createFleetRecords() : undefined;
285286
orchestratorTools.push(
286287
createTaskTool({
287288
cwd,
@@ -302,6 +303,7 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
302303
...(args.getBlobReader !== undefined ? { getBlobReader: args.getBlobReader } : {}),
303304
...(sa.useWorktree !== undefined ? { useWorktree: sa.useWorktree } : {}),
304305
...(args.telemetry !== undefined ? { telemetry: args.telemetry } : {}),
306+
...(fleetRecords !== undefined ? { fleetRecords } : {}),
305307
}),
306308
);
307309
if (sa.profiles !== undefined) {
@@ -321,9 +323,8 @@ export async function createAgentToolset(args: AgentToolsetArgs): Promise<AgentT
321323
// Mirror nested runSubAgent's orchestrator fleet mount (run.ts), but
322324
// reuse the existing TUI/exec session store — do not allocate a private
323325
// store only for these verbs. spawnAllowlist stays unwired on primary.
324-
if (sa.sessions !== undefined) {
326+
if (sa.sessions !== undefined && fleetRecords !== undefined) {
325327
const fleetSessions = sa.sessions;
326-
const fleetRecords = createFleetRecords();
327328
const fleetDeps = {
328329
permissionGate,
329330
inheritMcpTools: () => inheritedMcpTools,

src/subagent/agent-fleet.ts

Lines changed: 24 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,7 @@ import type { Settings } from "../config/settings.js";
6363
import { resolveEffortForRole } from "../provider/reasoning-effort.js";
6464
import { isCodexProviderName } from "../config/codex-providers.js";
6565
import { buildDispatchBrief, type TaskIntent } from "./report.js";
66-
import type { SubAgentSessionStore } from "./session-store.js";
66+
import { DEFAULT_CANCEL_REASON, type SubAgentSessionStore } from "./session-store.js";
6767
import type {
6868
NestedDispatchDeps,
6969
RunSubAgentParams,
@@ -75,6 +75,8 @@ import { cleanupSubAgentWorktree, createSubAgentWorktree, WorktreeError } from "
7575
import { NOOP_TELEMETRY, type Telemetry } from "../telemetry/index.js";
7676
import { classifyAgentName } from "../telemetry/classify.js";
7777
import type { DirectorPackage } from "../agent/directors/types.js";
78+
import { formatSubAgentTaskAuthFailureMessage } from "./inference-auth-failure.js";
79+
import { isSubAgentCancelError } from "./dispose.js";
7880

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

@@ -363,6 +365,8 @@ export type AgentFleetDeps = SubAgentSandboxDeps & {
363365
useWorktree?: boolean;
364366
/** Optional wall-clock budget (ms) forwarded to runSubAgent. */
365367
deadlineMs?: number;
368+
/** When false, tear the worker down on completion (task wrapper). Default true. */
369+
persist?: boolean;
366370
settings?: Settings | (() => Settings | undefined);
367371
catalog?: readonly ProviderCatalogEntry[] | (() => readonly ProviderCatalogEntry[]);
368372
onEvent?: (event: ReactorEmittedEvent) => void;
@@ -660,7 +664,7 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
660664
// stays alive for followup (agentRetained / interrupt keep-alive) —
661665
// matching run.ts's persisting gate so followup_task does not hit a
662666
// removed cwd.
663-
persist: true,
667+
persist: deps.persist !== false,
664668
onAgentReady: ({ close, interrupt, followup, deliver }) => {
665669
deps.sessions.registerClose(session.id, async (deadlineMs) => {
666670
try {
@@ -716,11 +720,24 @@ export function createSpawnAgentTool(deps: AgentFleetDeps): AgentTool {
716720
})
717721
.catch((err) => {
718722
// Always terminalize fleetRecords — including pre-progress cancel that
719-
// rethrows with no salvage — so wait_agents does not hang. fail()
720-
// no-ops when cancel already flipped the strip status.
721-
const message = err instanceof Error ? err.message : String(err);
722-
deps.fleetRecords.reject(session.id, message);
723-
deps.sessions.fail(session.id, message);
723+
// rethrows with no salvage — so wait_agents does not hang. Prefer
724+
// cancel semantics over fail when the strip already cancelled or the
725+
// throw is an AbortError (legacy task() parent contract).
726+
const alreadyCancelled = deps.sessions.get(session.id)?.status === "cancelled";
727+
if (alreadyCancelled || isSubAgentCancelError(err, childCtl.signal)) {
728+
if (!alreadyCancelled) {
729+
deps.sessions.cancel(session.id, DEFAULT_CANCEL_REASON);
730+
}
731+
const message = err instanceof Error ? err.message : String(err);
732+
deps.fleetRecords.reject(session.id, message);
733+
return;
734+
}
735+
// Auth failures keep the actionable Re-authenticate wording that
736+
// task()'s fused path surfaces via formatSubAgentTaskAuthFailureMessage.
737+
const authMessage = formatSubAgentTaskAuthFailureMessage(description, err);
738+
const failReason = authMessage ?? (err instanceof Error ? err.message : String(err));
739+
deps.fleetRecords.reject(session.id, failReason);
740+
deps.sessions.fail(session.id, failReason);
724741
})
725742
.finally(() => {
726743
telemetry.capture("subagent_end", {

src/subagent/run.ts

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -521,6 +521,8 @@ export async function runSubAgent(params: RunSubAgentParams): Promise<RunSubAgen
521521
);
522522
}
523523
const nd = params.nestedDispatch;
524+
const fleetSessions = nd.sessions ?? createSubAgentSessionStore();
525+
const fleetRecords = createFleetRecords();
524526
tools = [
525527
...tools,
526528
createTaskTool({
@@ -542,7 +544,8 @@ export async function runSubAgent(params: RunSubAgentParams): Promise<RunSubAgen
542544
telemetry: liveTelemetry,
543545
...(nd.onEvent !== undefined ? { onEvent: nd.onEvent } : {}),
544546
...(nd.onProgress !== undefined ? { onProgress: nd.onProgress } : {}),
545-
...(nd.sessions !== undefined ? { sessions: nd.sessions } : {}),
547+
sessions: fleetSessions,
548+
fleetRecords,
546549
...(nd.settings !== undefined ? { settings: nd.settings } : {}),
547550
...(nd.catalog !== undefined ? { catalog: nd.catalog } : {}),
548551
...(nd.profiles !== undefined ? { profiles: nd.profiles } : {}),
@@ -569,16 +572,9 @@ export async function runSubAgent(params: RunSubAgentParams): Promise<RunSubAgen
569572
createReadAgentTraceTool(nd.getWorkdirBase, {
570573
actorId: params.id,
571574
tier,
572-
getNodes: () => nd.sessions?.list() ?? [],
575+
getNodes: () => fleetSessions.list(),
573576
}),
574577
];
575-
// spawn_agent/wait_agents need a session store as their mailbox;
576-
// reuse the orchestrator's if it has one, else give this install its
577-
// own. fleetRecords holds terminal results the session store's
578-
// display cap would otherwise evict before wait_agents collects them
579-
// (see agent-fleet.ts).
580-
const fleetSessions = nd.sessions ?? createSubAgentSessionStore();
581-
const fleetRecords = createFleetRecords();
582578
const lifecycleAuthority = {
583579
actorId: params.id,
584580
tier,

0 commit comments

Comments
 (0)