From 408a92d91a240f03e456843f88748fc4e7b2be61 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sat, 11 Jul 2026 16:17:03 +0200 Subject: [PATCH] feat(coding-agent): enabled asynchronous background task execution - Enabled granular task execution by allowing batches to interleave blocking items with non-blocking async background spawns. - Updated task orchestration to support simultaneous inline result collection and persistent background job tracking. - Improved agent visibility in the job tool by reporting running subagents even when not explicitly linked to a backing job ID. - Enhanced terminal state handling to prevent premature tool block closures while async background operations remain active. --- docs/tools/task.md | 15 +- packages/coding-agent/CHANGELOG.md | 3 + .../coding-agent/src/async/job-manager.ts | 9 + .../src/modes/controllers/event-controller.ts | 10 +- .../controllers/tan-command-controller.ts | 2 +- .../coding-agent/src/prompts/tools/job.md | 2 +- .../coding-agent/src/prompts/tools/task.md | 5 +- packages/coding-agent/src/task/index.ts | 409 ++++++++++++------ packages/coding-agent/src/task/render.ts | 17 +- packages/coding-agent/src/tools/job.ts | 176 +++++++- packages/coding-agent/src/vibe/runtime.ts | 2 +- .../test/job-renderer-preview.test.ts | 35 ++ .../test/job-tool-agent-roster.test.ts | 159 +++++++ ...vent-controller-task-async-updates.test.ts | 135 ++++++ .../test/task/task-blocking-split.test.ts | 280 ++++++++++++ 15 files changed, 1096 insertions(+), 163 deletions(-) create mode 100644 packages/coding-agent/test/job-tool-agent-roster.test.ts create mode 100644 packages/coding-agent/test/modes/controllers/event-controller-task-async-updates.test.ts create mode 100644 packages/coding-agent/test/task/task-blocking-split.test.ts diff --git a/docs/tools/task.md b/docs/tools/task.md index 3601f1530..4bdf27a99 100644 --- a/docs/tools/task.md +++ b/docs/tools/task.md @@ -1,6 +1,6 @@ # task -> Spawn subagents — one per call, or a `tasks[]` batch per call (`task.batch`, default on). With `async.enabled=true`, spawns run in the background; otherwise the call blocks until they finish. +> Spawn subagents — one per call, or a `tasks[]` batch per call (`task.batch`, default on). With `async.enabled=true`, spawns run in the background; otherwise the call blocks until they finish. Execution mode is per item: an item whose agent type declares `blocking: true` (e.g. `scout`) runs inline and returns its result in the call, while non-blocking items in the same call still spawn as background jobs. ## Source - Entry: `packages/coding-agent/src/task/index.ts` @@ -52,10 +52,10 @@ The tool returns one text block plus `details: TaskToolDetails`. Background response (`async.enabled=true`): - `content`: `` Spawned agent `` (job ``). The result will be delivered when it yields. ... `` plus a coordination hint (`irc` DM when enabled, otherwise `job`). A batch call instead returns `` Spawned N background agents using . ... `` (the deduped per-item agent types, comma-joined) with a per-agent `- `` (job ``)` listing. -- `details`: `{ projectAgentsDir: null, results: [], totalDurationMs: 0, progress: [], async: { state: "running", jobId, type: "task" } }`. A batch call keeps one shared `progress[]` snapshot; `async.jobId` is the first started job and `async.state` aggregates ("running" until every job settles, "failed" if any spawn failed). +- `details`: `{ projectAgentsDir, results, totalDurationMs, progress: [], async: { state, jobId, type: "task" } }`. The call keeps one shared `progress[]` snapshot; `async.jobId` is the first started job and `async.state` aggregates over the async spawns ("running" until every job settles, "failed" if any spawn failed) — jobs that settled before the call returned are already reflected. A mixed call's `results` carries the blocking spawns' inline `SingleResult`s (pure background calls return `results: []`). - Live progress keeps streaming into the same tool block via `onUpdate(...)`; each final result arrives later as an async-result injection into the parent conversation. The delivery text appends a follow-up hint: `` is now idle — message it via `irc` to follow up; transcript at history:// `` (aborted variant points at the transcript only). -Settled response (`async.enabled=false`, no job manager, blocking agent, or async job body): +Settled response (`async.enabled=false`, no job manager, every item's agent `blocking: true`, or async job body): - `content`: summary rendered from `packages/coding-agent/src/prompts/tools/task-summary.md` with a preview capped at 5000 chars; `agent://` holds the full output. A sync batch concatenates the per-spawn summaries. - `details.results`: one `SingleResult` per spawn; `usage`, `outputPaths` populated (aggregated across spawns for a sync batch). @@ -75,12 +75,13 @@ Artifacts and side channels: ## Flow 1. `TaskTool.create(...)` discovers agents once per cwd through a process-level memo (`discoverAgentsForCreate`) to render the dynamic prompt description. 2. `execute(...)` repairs raw params (`repairTaskParams`), then validates: `schema` is always rejected; `tasks`/`context` are rejected unless `task.batch` is on; batch calls need a non-empty `tasks` (a `task` per item, unique provided names), a non-empty shared `context`, and no top-level `task` alongside `tasks`; flat calls need `task`. The call is then normalized into its spawn list (`resolveSpawnItems`). -3. Sync execution runs when `async.enabled=false`, the session has no `AsyncJobManager` (orphaned host), or the selected agent definition declares `blocking: true`; the call then runs every spawn through `#executeSync(...)` inline under the session-scoped semaphore. -4. Background execution runs only when `async.enabled=true` and the session has an `AsyncJobManager`: +3. Per-item execution split: items whose agent type declares `blocking: true` run inline; the rest become background jobs. The whole call runs sync when `async.enabled=false`, the session has no `AsyncJobManager` (orphaned host), or every item is blocking; inline spawns run through `#executeSync(...)` under the session-scoped semaphore. +4. Background execution (any non-blocking item with `async.enabled=true` and an `AsyncJobManager`): - agent ids are allocated up front via `AgentOutputManager.allocate(...)` — each item's `name`, or a generated AdjectiveNoun name — one per spawn; - one `type: "task"` job per spawn is registered with `session.asyncJobManager` (`id` = agent id, `queued: true`, `ownerId` = caller agent id) and the tool returns immediately; - each job body acquires the session-scoped `Semaphore` (one per `TaskTool` instance, sized from `task.maxConcurrency` at first use), marks the job running, runs `#executeSync(...)` with that spawn's params, and reports progress through the shared `buildAsyncDetails`/`onUpdate`; - a failed or aborted run throws `TaskJobError` so the job lands `failed`, but the agent itself stays registered and interrogable. + - a mixed call registers the async jobs first, then runs its blocking items inline and returns once they settle — the text combines the inline summaries with the spawned-job listing, and the block keeps rendering the still-running background rows beside the inline results. 5. `#executeSync(...)` runs the spawn path (`#runSpawn`), which rediscovers agents from disk, so runtime resolution can differ from the create-time description. 6. It resolves each spawn's requested `agent` type, rejects unknown or settings-disabled agents, and enforces parent spawn policy plus `PI_BLOCKED_AGENT` self-recursion prevention. 7. Output schema priority: agent frontmatter `output` → inherited parent session schema (the call itself never carries one). @@ -99,8 +100,8 @@ Artifacts and side channels: ## Modes / Variants - Execution mode - - Background job — `async.enabled=true`; spawns go through `AsyncJobManager`. - - Sync inline — `async.enabled=false`, no job manager, or `blocking: true` agent. + - Background job — `async.enabled=true`; non-blocking spawns go through `AsyncJobManager`. + - Sync inline — `async.enabled=false`, no job manager, or the item's agent declares `blocking: true` (per item: a mixed call runs both modes). - Batch mode (`task.batch`, default on) - on — `{ context, tasks[] }`: one independent spawn per item, required `context` shared across the call's spawns, `agent`/`isolated` per item. Lifecycle, revival, and concurrency semantics match N parallel single calls. - off — single spawn per call; `tasks`/`context` are rejected and removed from the schema. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 049ccb186..f024690aa 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -31,6 +31,9 @@ - Fixed visible per-keystroke lag while searching in the `/resume` session picker. Literal matches now rank synchronously from a cached per-session haystack, fuzzy scoring runs in bounded background chunks that converge to the same ranking (large listings previously rebuilt a fuzzy index per token per session on every keystroke), and the prompt-history SQLite lookup — an FTS query plus a LIKE scan over every stored prompt — is debounced off the keystroke path. - Fixed compiled Linux binary extension loading when bundled web-search header generation cannot read `header-generator` data files from the build-time path. ([#5178](https://github.com/can1357/oh-my-pi/issues/5178)) +- Fixed `job` `list`/empty-poll snapshots returning empty output: a no-job result now says so explicitly, and both snapshots list running subagents that have no backing job (agents woken via `irc`, spawns owned by another agent) with an `irc` coordination hint, so the tool's picture matches the UI's running-agent badge. Polling an agent id that has no job now explains the agent's registry state instead of a bare "no matching jobs". +- Fixed mixed blocking/non-blocking `task` batches degrading every spawn to synchronous execution: `blocking: true` was all-or-nothing per call, so one `scout` item silently stripped its non-blocking siblings of background execution (no job ids despite the tool description promising them, the whole turn blocked on the slowest worker, and every spawn died with the turn's abort signal). Execution mode is now per item — blocking items run inline and return their results in the call while non-blocking items still spawn as background jobs — with one shared progress/result aggregate, and the tool description now marks BLOCKING agents and documents the split. +- Fixed the task block "disappearing" when a background spawn's job settled while the tool call was still executing (e.g. beside a blocking spawn in the same turn): the final job frame replaced the block with an empty-results skeleton, dropped it from live tracking so the call's real result never rendered, and duplicated every completion frame. Final async frames now only finalize a block the call has already parked as background, returned task results report the converged job state instead of a stale "running", job completion frames carry the shared aggregate (inline results included) exactly once, and a mixed call's block renders its inline result rows alongside the still-running background rows. ## [16.4.4] - 2026-07-11 diff --git a/packages/coding-agent/src/async/job-manager.ts b/packages/coding-agent/src/async/job-manager.ts index c0fc5bbba..bdb3cbdb7 100644 --- a/packages/coding-agent/src/async/job-manager.ts +++ b/packages/coding-agent/src/async/job-manager.ts @@ -44,6 +44,12 @@ export interface AsyncJob { * supply an id (e.g. legacy tests, SDK consumers without an agent context). */ ownerId?: string; + /** + * Registry id of the subagent this job runs (task/tan/vibe jobs). Lets + * job-view code link a job row to its AgentRegistry ref even when the job + * id differs from the agent id (vibe turn jobs, tan clones). + */ + agentId?: string; /** * Job is registered but parked behind a caller-managed gate (e.g. a task * batch semaphore). Queued jobs do not count toward the running-job limit @@ -79,6 +85,8 @@ export interface AsyncJobRegisterOptions { id?: string; /** Registry id of the agent that owns this job; used to scope cancelAll. */ ownerId?: string; + /** Registry id of the subagent this job runs; see {@link AsyncJob.agentId}. */ + agentId?: string; onProgress?: (text: string, details?: Record) => void | Promise; /** Register the job in queued state; see {@link AsyncJob.queued}. */ queued?: boolean; @@ -192,6 +200,7 @@ export class AsyncJobManager { abortController, promise: Promise.resolve(), ownerId: options?.ownerId, + agentId: options?.agentId, queued: options?.queued === true, }; diff --git a/packages/coding-agent/src/modes/controllers/event-controller.ts b/packages/coding-agent/src/modes/controllers/event-controller.ts index e2baa993e..ddea4f86a 100644 --- a/packages/coding-agent/src/modes/controllers/event-controller.ts +++ b/packages/coding-agent/src/modes/controllers/event-controller.ts @@ -943,12 +943,18 @@ export class EventController { if (component) { const asyncState = (event.partialResult.details as { async?: { state?: string } } | undefined)?.async?.state; const isFinalAsyncState = asyncState === "completed" || asyncState === "failed"; + // A final async snapshot is terminal only for a parked background + // block (the call already returned and was kept alive for its jobs). + // While the call is still executing — a mixed blocking+async task + // call whose jobs settle before its blocking subset — treat it as a + // partial frame: `tool_execution_end` still owns the terminal result. + const isTerminal = isFinalAsyncState && this.#backgroundToolCallIds.has(event.toolCallId); component.updateResult( { ...event.partialResult, isError: asyncState === "failed" }, - !isFinalAsyncState, + !isTerminal, event.toolCallId, ); - if (isFinalAsyncState) { + if (isTerminal) { this.ctx.pendingTools.delete(event.toolCallId); this.#backgroundToolCallIds.delete(event.toolCallId); } diff --git a/packages/coding-agent/src/modes/controllers/tan-command-controller.ts b/packages/coding-agent/src/modes/controllers/tan-command-controller.ts index 5eb98ba65..0fcec31fc 100644 --- a/packages/coding-agent/src/modes/controllers/tan-command-controller.ts +++ b/packages/coding-agent/src/modes/controllers/tan-command-controller.ts @@ -160,7 +160,7 @@ export class TanCommandController { } } }, - { ownerId }, + { ownerId, agentId: cloneId }, ); } catch (error) { if (cloneFile) await removeCloneSession(cloneFile); diff --git a/packages/coding-agent/src/prompts/tools/job.md b/packages/coding-agent/src/prompts/tools/job.md index 848ee86fc..c98acfa0d 100644 --- a/packages/coding-agent/src/prompts/tools/job.md +++ b/packages/coding-agent/src/prompts/tools/job.md @@ -8,4 +8,4 @@ Background tasks deliver their results automatically the moment they finish. You - To watch EVERY running job, issue a call with NO fields at all (no `poll`, no `cancel`, no `list`). NEVER pass an array of every running ID. - A finished job's output, or the interrupting message and reason, is included in the next turn. - **Stop execution:** Pass `cancel` with job IDs to kill jobs that have hung, stalled, or are no longer needed. A cancel-only call returns immediately. -- **Snapshot:** Pass `list: true` to get the current status of all jobs without waiting. +- **Snapshot:** Pass `list: true` to get the current status of all jobs without waiting. The listing also names running subagents that have no job entry (e.g. agents woken via `irc`, or spawns owned by another agent) — those are coordinated through `irc`, not this tool. diff --git a/packages/coding-agent/src/prompts/tools/task.md b/packages/coding-agent/src/prompts/tools/task.md index 05d288548..acf9ce4c1 100644 --- a/packages/coding-agent/src/prompts/tools/task.md +++ b/packages/coding-agent/src/prompts/tools/task.md @@ -1,5 +1,6 @@ {{#if asyncEnabled}}{{#if batchEnabled}}Delegate work to background subagents by passing multiple items in a single `tasks[]` batch.{{else}}Delegate work to ONE background subagent per call.{{/if}} -Execution does not block your turn: you receive agent and job IDs immediately, and the final results deliver themselves when the subagents finish.{{else}}{{#if batchEnabled}}Run subagents synchronously by passing items in a `tasks[]` batch.{{else}}Run ONE subagent synchronously per call.{{/if}} +Execution does not block your turn: you receive agent and job IDs immediately, and the final results deliver themselves when the subagents finish.{{#if hasBlockingAgents}} +Exception: agents marked BLOCKING below run inline — their results return in this call, while non-blocking items in the same batch still spawn as background jobs.{{/if}}{{else}}{{#if batchEnabled}}Run subagents synchronously by passing items in a `tasks[]` batch.{{else}}Run ONE subagent synchronously per call.{{/if}} Execution blocks your turn: the call only returns once the work is completely finished.{{/if}} # Task Design @@ -54,7 +55,7 @@ Agent spawning is currently disabled. {{else}} Pick the most specific agent for each task. Use the default worker only when no specialist below fits. {{#list agents join="\n"}} -### {{name}}{{#if readOnly}} (READ-ONLY: no edit/write/command tools){{/if}} +### {{name}}{{#if readOnly}} (READ-ONLY: no edit/write/command tools){{/if}}{{#if blocking}} (BLOCKING: runs inline; its result returns in this call){{/if}} {{description}} {{#if readOnly}}Use ONLY for investigation and reporting; do the edits yourself or assign them to a writing agent.{{/if}} {{/list}} diff --git a/packages/coding-agent/src/task/index.ts b/packages/coding-agent/src/task/index.ts index fe46ef080..70deecc41 100644 --- a/packages/coding-agent/src/task/index.ts +++ b/packages/coding-agent/src/task/index.ts @@ -200,6 +200,7 @@ function renderDescription( name: agent.name, description: agent.description, readOnly: isReadOnlyAgent(agent), + blocking: agent.blocking === true, })); return prompt.render(taskDescriptionTemplate, { agents: renderedAgents, @@ -209,6 +210,7 @@ function renderDescription( isolationEnabled, batchEnabled, asyncEnabled, + hasBlockingAgents: renderedAgents.some(agent => agent.blocking), ircEnabled, }); } @@ -326,6 +328,65 @@ function spawnParamsFor(params: TaskParams, item: TaskItem, defaultAgent: string return spawn; } +/** One sync-executed spawn: its item, position in the original call, and (for mixed calls) a pre-claimed agent id. */ +interface SyncSpawnRef { + item: TaskItem; + index: number; + preAllocatedId?: string; +} + +/** Merged view of a sync spawn set's payloads: joined text plus flattened results/usage/paths. */ +interface MergedSyncPayloads { + contentParts: string[]; + results: SingleResult[]; + usage?: Usage; + outputPaths?: string[]; + projectAgentsDir: string | null; +} + +/** + * Merge per-spawn sync payloads into one result view. `index` is each spawn's + * position in the original call so batch rows keep stable ordering; a missing + * payload (cancelled before start) becomes an explanatory content line. + */ +function mergeSyncPayloads( + spawns: SyncSpawnRef[], + payloads: (AgentToolResult | undefined)[], +): MergedSyncPayloads { + const results: SingleResult[] = []; + const contentParts: string[] = []; + const outputPaths: string[] = []; + const usageTotals = createUsageTotals(); + let hasUsage = false; + let projectAgentsDir: string | null = null; + for (let position = 0; position < spawns.length; position++) { + const payload = payloads[position]; + const { item, index } = spawns[position]; + if (!payload) { + contentParts.push(`Task ${item.name?.trim() || `#${index + 1}`}: cancelled before start.`); + continue; + } + projectAgentsDir ??= payload.details?.projectAgentsDir ?? null; + const text = payload.content.find(part => part.type === "text")?.text; + if (text) contentParts.push(text); + for (const result of payload.details?.results ?? []) { + results.push({ ...result, index }); + if (result.usage) { + addUsageTotals(usageTotals, result.usage); + hasUsage = true; + } + if (result.outputPath) outputPaths.push(result.outputPath); + } + } + return { + contentParts, + results, + usage: hasUsage ? usageTotals : undefined, + outputPaths: outputPaths.length > 0 ? outputPaths : undefined, + projectAgentsDir, + }; +} + /** Generic worker agent types; several in one call usually means a more specific type exists. */ const GENERIC_SPAWN_AGENTS: ReadonlySet = new Set(["task", "sonic"]); @@ -368,10 +429,11 @@ export function buildCoordinationAdvisory( /** * Compose the non-blocking advisory appended to a `task` result: the * specialization nudge (from the per-item resolved agent types), plus — only - * when the siblings keep running after this call (`willRunAsync`) — the - * coordination suggestion. Coordination is gated on async because a sync - * fanout's siblings have already finished, so a "coordinate while they run" - * hint would misfire. Returns undefined when neither applies. + * when some spawns keep running after this call (`willRunAsync`) — the + * coordination suggestion over those still-live spawns (`items`). Coordination + * is gated on async because a sync spawn has already finished by the time the + * call returns, so a "coordinate while they run" hint would misfire. Returns + * undefined when neither applies. */ export function composeSpawnAdvisory(args: { agents: string[]; @@ -568,27 +630,30 @@ export class TaskTool implements AgentTool item.agent?.trim() || defaultAgent); - // Blocking is all-or-nothing for the call: one `blocking: true` - // agent type sends the whole fanout down the sync path. - const hasBlockingAgent = resolvedAgents.some( + // Execution mode is per item: an item whose agent type declares + // `blocking: true` runs inline on this turn (the parent waits on its + // result); every other item becomes a background job when async + // execution is available. + const itemBlocking = resolvedAgents.map( name => this.#discoveredAgents.find(agent => agent.name === name)?.blocking === true, ); const asyncEnabled = this.session.settings.get("async.enabled"); const manager = asyncEnabled ? this.session.asyncJobManager : undefined; + const asyncItems = manager ? spawnItems.filter((_, index) => !itemBlocking[index]) : []; const depthCapacity = canSpawnAtDepth( this.session.settings.get("task.maxRecursionDepth") ?? 2, this.session.taskDepth ?? 0, ); const ircEnabled = isIrcEnabled(this.session.settings, this.session.taskDepth ?? 0); - // Coordination only makes sense when the siblings keep running after this - // call returns (async). In the sync fallback they have already completed, - // so a "coordinate while they run" hint would misfire. - const willRunAsync = !!manager && !hasBlockingAgent; + // Coordination only makes sense for spawns that keep running after this + // call returns (the async subset). Blocking items have already completed + // by then, so a "coordinate while they run" hint would misfire. + const willRunAsync = asyncItems.length > 0; const advisory = this.session.suppressSpawnAdvisory ? undefined : composeSpawnAdvisory({ agents: resolvedAgents, - items: spawnItems, + items: asyncItems, depthCapacity, ircEnabled, willRunAsync, @@ -609,12 +674,12 @@ export class TaskTool implements AgentTool null)); - const agentLabel = [...new Set(resolvedAgents)].join(", "); - const spawns: Array<{ agentId: string; item: TaskItem; progress: AgentProgress }> = []; + const callStartedAt = Date.now(); + const spawns: Array<{ + agentId: string; + item: TaskItem; + index: number; + blocking: boolean; + progress: AgentProgress; + }> = []; for (let index = 0; index < spawnItems.length; index++) { const item = spawnItems[index]; const agentType = resolvedAgents[index]; @@ -636,6 +707,8 @@ export class TaskTool implements AgentTool !spawn.blocking); + const syncSpawns = spawns.filter(spawn => spawn.blocking); + const agentLabel = [...new Set(asyncSpawns.map(spawn => spawn.progress.agent))].join(", "); - // Aggregate async state for the one tool call: every spawn's job reports - // into the shared progress snapshot; the call stays "running" until all - // jobs settle, then turns "failed" if any spawn failed. The single-spawn - // case passes the job's own suggestion through (pre-batch behavior). - const single = spawns.length === 1; + // Aggregate state for the one tool call. Async spawns report into the + // shared progress snapshot through their jobs: the async half stays + // "running" until every job settles, then turns "failed" if any spawn + // failed. Blocking spawns run inline below and land in `results` before + // the call returns, so post-return job updates never drop them. let settledCount = 0; let failedCount = 0; - let primaryJobId = spawns[0].agentId; - const buildAsyncDetails = (state: "running" | "completed" | "failed", jobId: string): TaskToolDetails => ({ - projectAgentsDir: null, - results: [], - totalDurationMs: 0, + let primaryJobId = asyncSpawns[0].agentId; + const syncResults: SingleResult[] = []; + let syncUsage: Usage | undefined; + let syncOutputPaths: string[] | undefined; + let syncProjectAgentsDir: string | null = null; + const buildAsyncDetails = (): TaskToolDetails => ({ + projectAgentsDir: syncProjectAgentsDir, + results: [...syncResults], + totalDurationMs: Date.now() - callStartedAt, + usage: syncUsage, + outputPaths: syncOutputPaths, progress: spawns.map(spawn => ({ ...spawn.progress })), async: { - state: single ? state : settledCount < spawns.length ? "running" : failedCount > 0 ? "failed" : "completed", - jobId: single ? jobId : primaryJobId, + state: settledCount < asyncSpawns.length ? "running" : failedCount > 0 ? "failed" : "completed", + jobId: primaryJobId, type: "task", }, }); const started: Array<{ agentId: string; jobId: string }> = []; const failedSchedules: string[] = []; - for (const spawn of spawns) { + for (const spawn of asyncSpawns) { try { const jobId = this.#registerSpawnJob({ manager, @@ -704,58 +786,129 @@ export class TaskTool implements AgentTool 0 + ? ` Failed to schedule ${failedSchedules.length} spawn${failedSchedules.length === 1 ? "" : "s"}: ${failedSchedules.join("; ")}.` + : ""; + const coordinationHint = + started.length === 1 + ? ircEnabled + ? `DM \`${started[0].agentId}\` via \`irc\` to coordinate while it runs; use \`job\` only to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.` + : `Use \`job\` to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.` + : ircEnabled + ? `DM these ids via \`irc\` to coordinate while they run; use \`job\` only to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.` + : `Use \`job\` to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task by id.`; + + if (syncSpawns.length === 0) { + if (spawns.length === 1) { + const { agentId, jobId } = started[0]; + onUpdate?.({ + content: [{ type: "text", text: `Spawned agent \`${agentId}\`...` }], + details: buildAsyncDetails(), + }); + return withAdvisory({ + content: [ + { + type: "text", + text: `Spawned agent \`${agentId}\` (job \`${jobId}\`). The result will be delivered when it yields. ${coordinationHint}`, + }, + ], + details: buildAsyncDetails(), + }); + } + const startedListing = started.map(({ agentId, jobId }) => `- \`${agentId}\` (job \`${jobId}\`)`).join("\n"); onUpdate?.({ - content: [{ type: "text", text: `Spawned agent \`${agentId}\`...` }], - details: buildAsyncDetails("running", jobId), + content: [{ type: "text", text: `Spawned ${started.length} agents...` }], + details: buildAsyncDetails(), }); return withAdvisory({ content: [ { type: "text", - text: `Spawned agent \`${agentId}\` (job \`${jobId}\`). The result will be delivered when it yields. ${coordinationHint}`, + text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result will be delivered when that agent yields.\n${startedListing}\n${coordinationHint}`, }, ], - details: buildAsyncDetails("running", jobId), + details: buildAsyncDetails(), }); } - const coordinationHint = ircEnabled - ? `DM these ids via \`irc\` to coordinate while they run; use \`job\` only to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task.` - : `Use \`job\` to inspect (\`list\`), wait (\`poll\`), or cancel a stuck task by id.`; - const scheduleFailureSummary = - failedSchedules.length > 0 - ? ` Failed to schedule ${failedSchedules.length} spawn${failedSchedules.length === 1 ? "" : "s"}: ${failedSchedules.join("; ")}.` - : ""; - const startedListing = started.map(({ agentId, jobId }) => `- \`${agentId}\` (job \`${jobId}\`)`).join("\n"); + // Mixed call: the async jobs above already run detached; the blocking + // subset runs inline and gates the call's return — exactly what each + // agent type declares (`blocking: true` = the parent waits on it). + const syncLabel = syncSpawns.map(spawn => `\`${spawn.agentId}\``).join(", "); onUpdate?.({ - content: [{ type: "text", text: `Spawned ${started.length} agents...` }], - details: buildAsyncDetails("running", primaryJobId), - }); - return withAdvisory({ content: [ { type: "text", - text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result will be delivered when that agent yields.\n${startedListing}\n${coordinationHint}`, + text: `Running ${syncLabel} inline; ${started.length} background agent${started.length === 1 ? "" : "s"} spawned...`, }, ], - details: buildAsyncDetails("running", primaryJobId), + details: buildAsyncDetails(), + }); + const payloads = await this.#runSyncSpawns({ + toolCallId, + params, + defaultAgent, + signal, + spawns: syncSpawns.map(spawn => ({ item: spawn.item, index: spawn.index, preAllocatedId: spawn.agentId })), + onItemProgress: onUpdate + ? (index, progress) => { + const spawn = spawns[index]; + if (spawn) spawn.progress = { ...progress, index }; + onUpdate({ + content: [{ type: "text", text: `Running ${syncLabel} inline...` }], + details: buildAsyncDetails(), + }); + } + : undefined, + }); + const merged = mergeSyncPayloads( + syncSpawns.map(spawn => ({ item: spawn.item, index: spawn.index })), + payloads, + ); + syncResults.push(...merged.results); + syncUsage = merged.usage; + syncOutputPaths = merged.outputPaths; + syncProjectAgentsDir = merged.projectAgentsDir; + // Settle the inline spawns' progress rows from their merged results so + // post-return job updates carry final statuses, not the last snapshot. + for (let position = 0; position < syncSpawns.length; position++) { + const spawn = syncSpawns[position]; + const result = merged.results.find(r => r.id === spawn.agentId); + if (result) { + spawn.progress.status = result.aborted + ? "aborted" + : result.exitCode === 0 && !result.error + ? "completed" + : "failed"; + spawn.progress.durationMs = result.durationMs; + } else { + spawn.progress.status = payloads[position] ? "failed" : "aborted"; + } + } + + const spawnedSummary = + started.length > 0 + ? `Spawned ${started.length} background agent${started.length === 1 ? "" : "s"}.${scheduleFailureSummary} Each result will be delivered when that agent yields.\n${started.map(({ agentId, jobId }) => `- \`${agentId}\` (job \`${jobId}\`)`).join("\n")}\n${coordinationHint}` + : scheduleFailureSummary.trim(); + const text = [merged.contentParts.join("\n\n"), spawnedSummary] + .filter(section => section.trim().length > 0) + .join("\n\n"); + return withAdvisory({ + content: [{ type: "text", text: text.length > 0 ? text : "No results." }], + details: buildAsyncDetails(), }); } @@ -772,7 +925,7 @@ export class TaskTool implements AgentTool TaskToolDetails; + buildDetails: () => TaskToolDetails; onUpdate?: AgentToolUpdateCallback; onSettled?: (failed: boolean) => void; }): string { @@ -788,7 +941,7 @@ export class TaskTool implements AgentTool { + async ({ signal: runSignal, reportProgress, markRunning }) => { const startedAt = Date.now(); const semaphore = this.#getSpawnSemaphore(); let semaphoreHeld = false; @@ -822,10 +975,7 @@ export class TaskTool implements AgentTool, - ); + await reportProgress(`Running background task ${agentId}...`); const result = await this.#executeSync( toolCallId, spawnParams, @@ -855,14 +1005,7 @@ export class TaskTool implements AgentTool, - ); - onUpdate?.({ - content: [{ type: "text", text: statusText }], - details: buildDetails(resultFailed ? "failed" : "completed", ownJobId), - }); + await reportProgress(statusText); const deliveryText = `${finalText}${buildFollowUpHint(singleResult?.aborted === true)}`; if (resultFailed) { // Mark the job itself failed; the failed agent stays interrogable. @@ -877,11 +1020,7 @@ export class TaskTool implements AgentTool); - onUpdate?.({ - content: [{ type: "text", text: statusText }], - details: buildDetails("failed", ownJobId), - }); + await reportProgress(statusText); const message = error instanceof Error ? error.message : String(error); const hint = AgentRegistry.global().get(agentId) ? buildFollowUpHint(false) : ""; throw new TaskJobError(`${message}${hint}`); @@ -891,21 +1030,21 @@ export class TaskTool implements AgentTool { - const progressDetails = (details as TaskToolDetails | undefined) ?? buildDetails("running", agentId); - onUpdate?.({ content: [{ type: "text", text }], details: progressDetails }); + onProgress: text => { + onUpdate?.({ content: [{ type: "text", text }], details: buildDetails() }); }, }, ); } /** - * Sync fallback fan-out (no job manager, or a `blocking: true` agent): run - * every spawn to completion inline and merge the per-spawn payloads into a - * single tool result. The session-scoped semaphore still bounds concurrency - * across parallel task calls. + * Sync fan-out (async unavailable, or every item's agent type is + * `blocking: true`): run every spawn to completion inline and merge the + * per-spawn payloads into a single tool result. The session-scoped + * semaphore still bounds concurrency across parallel task calls. */ async #executeSyncFanout( toolCallId: string, @@ -915,8 +1054,8 @@ export class TaskTool implements AgentTool, ): Promise> { - const semaphore = this.#getSpawnSemaphore(); if (spawnItems.length === 1) { + const semaphore = this.#getSpawnSemaphore(); const invokedAt = Date.now(); await semaphore.acquire(signal); const acquiredAt = Date.now(); @@ -952,30 +1091,75 @@ export class TaskTool implements AgentTool { + const payloads = await this.#runSyncSpawns({ + toolCallId, + params, + defaultAgent, + signal, + spawns: spawnItems.map((item, index) => ({ item, index })), + onItemProgress: onUpdate + ? (index, progress) => { + latestProgress.set(index, { ...progress, index }); + emitCombined(); + } + : undefined, + }); + + const merged = mergeSyncPayloads( + spawnItems.map((item, index) => ({ item, index })), + payloads, + ); + return { + content: [{ type: "text", text: merged.contentParts.join("\n\n") }], + details: { + projectAgentsDir: merged.projectAgentsDir, + results: merged.results, + totalDurationMs: Date.now() - startTime, + usage: merged.usage, + outputPaths: merged.outputPaths, + }, + }; + } + + /** + * Run a set of spawns to completion inline, bounded by the session spawn + * semaphore. `preAllocatedId` reuses an id claimed up front (mixed calls); + * `index` is each item's position in the original call so progress rows and + * merged results keep stable ordering. Per-item progress snapshots flow + * through `onItemProgress`. Returns per-spawn payloads in input order; + * `undefined` marks a spawn cancelled before it started. + */ + async #runSyncSpawns(args: { + toolCallId: string; + params: TaskParams; + defaultAgent: string; + spawns: SyncSpawnRef[]; + signal?: AbortSignal; + onItemProgress?: (index: number, progress: AgentProgress) => void; + }): Promise<(AgentToolResult | undefined)[]> { + const { toolCallId, params, defaultAgent, spawns, signal, onItemProgress } = args; + const semaphore = this.#getSpawnSemaphore(); + const { results } = await mapWithConcurrencyLimit( + spawns, + spawns.length, + async (spawn, _position, workerSignal) => { const invokedAt = Date.now(); await semaphore.acquire(workerSignal); const acquiredAt = Date.now(); try { - const itemOnUpdate: AgentToolUpdateCallback | undefined = onUpdate + const itemOnUpdate: AgentToolUpdateCallback | undefined = onItemProgress ? update => { const progress = update.details?.progress?.[0]; - if (progress) { - latestProgress.set(index, { ...progress, index }); - emitCombined(); - } + if (progress) onItemProgress(spawn.index, progress); } : undefined; return await this.#executeSync( toolCallId, - spawnParamsFor(params, item, defaultAgent), + spawnParamsFor(params, spawn.item, defaultAgent), workerSignal, itemOnUpdate, - undefined, - index, + spawn.preAllocatedId, + spawn.index, false, { invokedAt, acquiredAt }, ); @@ -985,42 +1169,7 @@ export class TaskTool implements AgentTool part.type === "text")?.text; - if (text) contentParts.push(text); - for (const result of payload.details?.results ?? []) { - results.push({ ...result, index }); - if (result.usage) { - addUsageTotals(usageTotals, result.usage); - hasUsage = true; - } - if (result.outputPath) outputPaths.push(result.outputPath); - } - } - - return { - content: [{ type: "text", text: contentParts.join("\n\n") }], - details: { - projectAgentsDir, - results, - totalDurationMs: Date.now() - startTime, - usage: hasUsage ? usageTotals : undefined, - outputPaths: outputPaths.length > 0 ? outputPaths : undefined, - }, - }; + return results; } /** diff --git a/packages/coding-agent/src/task/render.ts b/packages/coding-agent/src/task/render.ts index 562ebc94c..78005670e 100644 --- a/packages/coding-agent/src/task/render.ts +++ b/packages/coding-agent/src/task/render.ts @@ -1579,8 +1579,10 @@ export function renderResult( const frozen = options.renderContext?.frozen === true; const lines: string[] = []; + // Result rows win once any exist; progress rows for spawns without a + // result (a mixed call's async subset) render as a supplement below. const shouldRenderProgress = - Boolean(details.progress && details.progress.length > 0) && (isPartial || details.results.length === 0); + Boolean(details.progress && details.progress.length > 0) && details.results.length === 0; if (shouldRenderProgress && details.progress) { const ordered = orderProgressForDisplay(details.progress); // Collapsed view keeps the live edge: finished rows sort to the top of @@ -1607,6 +1609,19 @@ export function renderResult( ); } + // Mixed blocking+async call: async spawns never land in `results` + // (their payloads deliver through jobs) — keep their rows visible + // beside the finalized inline results, live while running and + // settled once their jobs finish. + const supplementalProgress = details.progress + ? orderProgressForDisplay( + details.progress.filter(progress => !details.results.some(res => res.id === progress.id)), + ) + : []; + for (const progress of supplementalProgress) { + lines.push(...renderAgentProgress(progress, "", " ", expanded, theme, spinnerFrame, frozen)); + } + const summaryParts: string[] = []; if (abortedCount > 0) summaryParts.push(theme.fg("error", `${abortedCount} aborted`)); if (successCount > 0) summaryParts.push(theme.fg("success", `${successCount} succeeded`)); diff --git a/packages/coding-agent/src/tools/job.ts b/packages/coding-agent/src/tools/job.ts index 2e6911429..360cbf908 100644 --- a/packages/coding-agent/src/tools/job.ts +++ b/packages/coding-agent/src/tools/job.ts @@ -61,9 +61,26 @@ interface CancelOutcome { message: string; } +/** + * A live subagent from the AgentRegistry that has no backing job in the + * AsyncJobManager — e.g. an idle agent woken (or a parked agent revived) via + * `irc`, or a spawn owned by another agent. Surfaced by `list` and empty-poll + * snapshots so the job tool's picture matches the UI's running-agent count. + */ +interface AgentActivitySnapshot { + id: string; + parentId?: string; + /** Latest activity gist recorded by the registry (display-only). */ + activity?: string; + /** Time since the agent was registered. */ + ageMs: number; +} + export interface JobToolDetails { jobs: JobSnapshot[]; cancelled?: { id: string; status: CancelStatus }[]; + /** Running subagents not represented by a job row in this result. */ + agents?: AgentActivitySnapshot[]; } /** @@ -118,7 +135,9 @@ export class JobTool implements AgentTool { if (params.cancel?.length || params.poll?.length) { throw new ToolError("`list` cannot be combined with `poll` or `cancel`."); } - return this.#buildResult(manager, manager.getAllJobs(ownerFilter), []); + const jobs = manager.getAllJobs(ownerFilter); + const agents = this.#runningAgentsOutsideJobs(); + return this.#buildResult(manager, jobs, [], agents); } const cancelIds = params.cancel ?? []; @@ -166,15 +185,37 @@ export class JobTool implements AgentTool { const cancelledJobs = this.#visibleJobs(manager, cancelIds, ownerId); return this.#buildResult(manager, cancelledJobs, cancelOutcomes); } - const message = requestedPollIds?.length - ? `No matching jobs found for IDs: ${requestedPollIds.join(", ")}` - : "No running background jobs to wait for."; + // Zero pollable jobs is not necessarily "nothing running": agents + // woken via irc or owned by another agent run with no job entry. + // Report them so the snapshot matches the UI's running-agent count + // (task job ids are agent ids, so a stale poll id often names one). + const agents = this.#runningAgentsOutsideJobs(); + const lines: string[] = []; + if (requestedPollIds?.length) { + lines.push(`No matching jobs found for IDs: ${requestedPollIds.join(", ")}`); + const registry = this.session.agentRegistry; + for (const id of requestedPollIds) { + const ref = registry?.get(id); + if (!ref) continue; + lines.push( + ref.status === "running" + ? `- \`${id}\` is a running agent with no job entry — coordinate via \`irc\`; transcript at history://${id}` + : `- \`${id}\` is a ${ref.status} agent (its job is gone) — transcript at history://${id}`, + ); + } + } else { + lines.push("No running background jobs to wait for."); + } + if (agents.length > 0) { + lines.push("", ...this.#describeAgents(agents)); + } return { - content: [{ type: "text", text: message }], - details: { jobs: [] }, + content: [{ type: "text", text: lines.join("\n") }], + details: { jobs: [], ...(agents.length ? { agents } : {}) }, // Nothing found / nothing to wait for is noise once consumed — - // the follow-up call has already corrected course. - useless: true, + // the follow-up call has already corrected course. Running agents + // are real state the model may act on, so keep those results. + ...(agents.length === 0 ? { useless: true } : {}), }; } @@ -266,6 +307,57 @@ export class JobTool implements AgentTool { return out; } + /** + * Running subagents from the registry that are not covered by one of the + * caller's running jobs. Agents woken via `irc` (idle wake / park revival) + * and spawns owned by another agent run with no AsyncJobManager entry, yet + * the UI's agent badge counts them — a snapshot must account for that + * activity instead of implying the system is quiet. Existence is already + * public via the `irc` roster, so listing ids here leaks nothing new; job + * *control* stays owner-scoped. + */ + #runningAgentsOutsideJobs(): AgentActivitySnapshot[] { + const registry = this.session.agentRegistry; + if (!registry) return []; + const selfId = this.session.getAgentId?.() ?? undefined; + // Cover = the caller's RUNNING jobs only. A settled job still sitting in + // delivery retention must not hide its agent if that agent was re-woken + // (e.g. via irc) and is running again without a job. + const covered = new Set(); + const manager = this.session.asyncJobManager; + if (manager) { + for (const job of manager.getRunningJobs(selfId ? { ownerId: selfId } : undefined)) { + covered.add(job.id); + if (job.agentId) covered.add(job.agentId); + } + } + const now = Date.now(); + const out: AgentActivitySnapshot[] = []; + for (const ref of registry.list()) { + if (ref.kind !== "sub" || ref.status !== "running") continue; + if (ref.id === selfId || covered.has(ref.id)) continue; + out.push({ + id: ref.id, + ...(ref.parentId ? { parentId: ref.parentId } : {}), + ...(ref.activity ? { activity: ref.activity } : {}), + ageMs: Math.max(0, now - ref.createdAt), + }); + } + return out; + } + + /** Model-facing lines for the running-agents section shared by `list` and empty-poll results. */ + #describeAgents(agents: AgentActivitySnapshot[]): string[] { + const lines = [`## Running Agents (${agents.length}) — not job-backed\n`]; + for (const agent of agents) { + const parent = agent.parentId ? ` (spawned by \`${agent.parentId}\`)` : ""; + const activity = agent.activity ? ` — ${agent.activity}` : ""; + lines.push(`- \`${agent.id}\`${parent} — up ${formatDuration(agent.ageMs)}${activity}`); + } + lines.push("", "These agents have no job entry; coordinate via `irc`, transcripts at `history://`."); + return lines; + } + #snapshotJobs( jobs: { id: string; @@ -305,6 +397,7 @@ export class JobTool implements AgentTool { errorText?: string; }[], cancelOutcomes: CancelOutcome[], + agents: AgentActivitySnapshot[] = [], ): AgentToolResult { // Deduplicate by id (cancelled jobs may also appear in the watched set). const seen = new Set(); @@ -350,9 +443,21 @@ export class JobTool implements AgentTool { } } + if (agents.length > 0) { + if (lines.length > 0) lines.push(""); + lines.push(...this.#describeAgents(agents)); + } + + // A tool result must never be empty text — the model cannot tell "no + // jobs" from a malfunction (reported exactly that way in QA). + if (lines.length === 0) { + lines.push("No background jobs."); + } + const details: JobToolDetails = { jobs: jobResults, ...(cancelOutcomes.length ? { cancelled: cancelOutcomes.map(({ id, status }) => ({ id, status })) } : {}), + ...(agents.length ? { agents } : {}), }; return { content: [{ type: "text", text: lines.join("\n").trimEnd() }], @@ -463,8 +568,9 @@ export const jobToolRenderer = { args?: JobRenderArgs, ): Component { let jobs = result.details?.jobs ?? []; + const agents = result.details?.agents ?? []; - if (jobs.length === 0) { + if (jobs.length === 0 && agents.length === 0) { const fallback = result.content?.find(c => c.type === "text")?.text || "No jobs to process"; const header = renderStatusLine({ icon: "warning", title: describeTarget(args) || "Job" }, uiTheme); return new Text([header, formatEmptyMessage(fallback, uiTheme)].join("\n"), 0, 0); @@ -474,7 +580,10 @@ export const jobToolRenderer = { ? !args.list && (!args.cancel || args.cancel.length === 0 || args.poll !== undefined) : true; - if (!options.isPartial && isPollCall) { + // Agent-carrying results (list / empty-poll roster) are real snapshots, + // not displaceable waiting frames — only agentless polls collapse their + // still-running rows once sealed. + if (!options.isPartial && isPollCall && agents.length === 0) { jobs = jobs.filter(job => job.status !== "running"); if (jobs.length === 0) { return new Text("", 0, 0); @@ -490,20 +599,26 @@ export const jobToolRenderer = { if (counts.completed > 0) meta.push(uiTheme.fg("success", `${counts.completed} done`)); if (counts.failed > 0) meta.push(uiTheme.fg("error", `${counts.failed} failed`)); if (counts.cancelled > 0) meta.push(uiTheme.fg("warning", `${counts.cancelled} cancelled`)); + if (agents.length > 0 && jobs.length > 0) { + meta.push(uiTheme.fg("accent", `${agents.length} agent${agents.length === 1 ? "" : "s"}`)); + } - const headerIcon: ToolUIStatus = counts.failed > 0 ? "warning" : counts.running > 0 ? "info" : "success"; + const headerIcon: ToolUIStatus = + counts.failed > 0 ? "warning" : counts.running > 0 || agents.length > 0 ? "info" : "success"; const jobsNoun = jobs.length === 1 ? "job" : "jobs"; const description = - counts.running > 0 - ? counts.running === jobs.length - ? `waiting on ${jobs.length} ${jobsNoun}` - : `waiting on ${counts.running} of ${jobs.length} ${jobsNoun}` - : `${jobs.length} ${jobsNoun} settled`; + jobs.length === 0 + ? `${agents.length} running agent${agents.length === 1 ? "" : "s"} — no jobs` + : counts.running > 0 + ? counts.running === jobs.length + ? `waiting on ${jobs.length} ${jobsNoun}` + : `waiting on ${counts.running} of ${jobs.length} ${jobsNoun}` + : `${jobs.length} ${jobsNoun} settled`; const header = renderStatusLine( { icon: headerIcon, - spinnerFrame: counts.running > 0 ? options.spinnerFrame : undefined, + spinnerFrame: counts.running > 0 || agents.length > 0 ? options.spinnerFrame : undefined, title: description, meta, }, @@ -598,7 +713,32 @@ export const jobToolRenderer = { uiTheme, ); - const all = [header, ...itemLines].map(l => truncateToWidth(l, width, Ellipsis.Unicode)); + // Agents run outside job control; render them as their own tree so + // they never skew the job counts or the "waiting on N jobs" title. + const agentLines = + agents.length === 0 + ? [] + : renderTreeList( + { + items: agents, + expanded, + maxCollapsed: COLLAPSED_LIST_LIMIT, + itemType: "agent", + renderItem: agent => { + const icon = formatStatusIcon("running", uiTheme, options.spinnerFrame); + const badge = formatBadge("agent", "accent", uiTheme); + const gist = agent.activity + ? ` ${uiTheme.fg("toolOutput", truncateToWidth(replaceTabs(agent.activity), LABEL_MAX_WIDTH, Ellipsis.Unicode))}` + : ""; + const parent = agent.parentId ? uiTheme.fg("dim", ` ← ${agent.parentId}`) : ""; + const age = uiTheme.fg("dim", formatDuration(agent.ageMs)); + return [`${icon} ${uiTheme.fg("muted", agent.id)} ${badge}${gist} ${age}${parent}`]; + }, + }, + uiTheme, + ); + + const all = [header, ...itemLines, ...agentLines].map(l => truncateToWidth(l, width, Ellipsis.Unicode)); cached = { key, lines: all }; return all; }, diff --git a/packages/coding-agent/src/vibe/runtime.ts b/packages/coding-agent/src/vibe/runtime.ts index 200ed2c32..d7e1215aa 100644 --- a/packages/coding-agent/src/vibe/runtime.ts +++ b/packages/coding-agent/src/vibe/runtime.ts @@ -603,7 +603,7 @@ export class VibeSessionRegistry { ); } }, - { id: `${record.id}-t${turnIndex}`, ownerId: record.ownerId }, + { id: `${record.id}-t${turnIndex}`, agentId: record.id, ownerId: record.ownerId }, ); turn.jobId = jobId; record.turn = turn; diff --git a/packages/coding-agent/test/job-renderer-preview.test.ts b/packages/coding-agent/test/job-renderer-preview.test.ts index 74976e1d2..e261a9673 100644 --- a/packages/coding-agent/test/job-renderer-preview.test.ts +++ b/packages/coding-agent/test/job-renderer-preview.test.ts @@ -232,5 +232,40 @@ describe("job renderer task-result preview", () => { expect(output).toContain("Job3 running"); expect(output).toContain("waiting on 2 of 3 jobs"); }); + + it("renders agent rows for running agents outside job control", () => { + const result = { + content: [{ type: "text" as const, text: "" }], + details: { + jobs: [], + agents: [{ id: "Worker", parentId: "Main", activity: "grepping the tree", ageMs: 65_000 }], + }, + }; + const component = jobToolRenderer.renderResult( + result, + { expanded: true, isPartial: false } as Parameters[1], + theme, + { list: true }, + ); + const output = Bun.stripANSI((component.render(120) as readonly string[]).join("\n")); + expect(output).toContain("1 running agent — no jobs"); + expect(output).toContain("Worker"); + expect(output).toContain("grepping the tree"); + }); + + it("keeps a sealed bare-poll result visible when it carries an agent roster", () => { + const result = { + content: [{ type: "text" as const, text: "No running background jobs to wait for." }], + details: { jobs: [], agents: [{ id: "Worker", ageMs: 1_000 }] }, + }; + const component = jobToolRenderer.renderResult( + result, + { expanded: true, isPartial: false } as Parameters[1], + theme, + { poll: [] }, + ); + const output = Bun.stripANSI((component.render(120) as readonly string[]).join("\n")); + expect(output).toContain("Worker"); + }); }); }); diff --git a/packages/coding-agent/test/job-tool-agent-roster.test.ts b/packages/coding-agent/test/job-tool-agent-roster.test.ts new file mode 100644 index 000000000..ce689f1f6 --- /dev/null +++ b/packages/coding-agent/test/job-tool-agent-roster.test.ts @@ -0,0 +1,159 @@ +/** + * The `job` tool's snapshot contract: `list` and empty-poll results must never + * come back as empty text, and they must surface running subagents that have + * no backing job (irc-woken/revived agents, spawns owned by another agent) so + * the tool's picture matches the UI's running-agent count. Regression for the + * QA report "job list returned no status output despite known running + * background jobs and subagents". + */ +import { afterEach, describe, expect, test } from "bun:test"; +import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async"; +import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; +import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools"; +import { JobTool } from "@oh-my-pi/pi-coding-agent/tools/job"; + +const managers: AsyncJobManager[] = []; + +function createManager(): AsyncJobManager { + const manager = new AsyncJobManager({ onJobComplete: () => {} }); + managers.push(manager); + return manager; +} + +function createToolSession(options: { + manager?: AsyncJobManager; + registry?: AgentRegistry; + agentId?: string; +}): ToolSession { + return { + cwd: process.cwd(), + hasUI: false, + settings: { + get: (key: string) => (key === "async.pollWaitDuration" ? "5s" : undefined), + }, + getSessionFile: () => null, + getSessionSpawns: () => null, + getAgentId: () => options.agentId ?? null, + asyncJobManager: options.manager, + agentRegistry: options.registry, + } as unknown as ToolSession; +} + +function registerRunningSub(registry: AgentRegistry, id: string, parentId = "Main"): void { + registry.register({ id, displayName: id, kind: "sub", parentId, session: null }); +} + +function resultText(result: { content: Array<{ type: string; text?: string }> }): string { + return result.content.find(part => part.type === "text")?.text ?? ""; +} + +const neverResolves = () => new Promise(() => {}); + +afterEach(async () => { + for (const manager of managers.splice(0)) { + await manager.dispose({ timeoutMs: 200 }); + } +}); + +describe("job list snapshot", () => { + test("empty list reports 'no jobs' instead of empty output", async () => { + const tool = new JobTool(createToolSession({ manager: createManager(), agentId: "Main" })); + + const result = await tool.execute("call", { list: true }); + + expect(resultText(result)).toBe("No background jobs."); + expect(result.details?.jobs).toEqual([]); + }); + + test("list surfaces running subagents that have no backing job", async () => { + const registry = new AgentRegistry(); + registerRunningSub(registry, "Worker"); + registerRunningSub(registry, "Idler"); + registry.setStatus("Idler", "idle"); + registry.register({ id: "advisor", displayName: "advisor", kind: "advisor", session: null }); + registry.register({ id: "Main", displayName: "Main", kind: "main", session: null }); + const tool = new JobTool(createToolSession({ manager: createManager(), registry, agentId: "Main" })); + + const result = await tool.execute("call", { list: true }); + + expect(result.details?.agents?.map(agent => agent.id)).toEqual(["Worker"]); + const text = resultText(result); + expect(text).toContain("Running Agents (1)"); + expect(text).toContain("Worker"); + expect(result.useless).toBeUndefined(); + }); + + test("agents covered by the caller's running jobs are not double-listed", async () => { + const manager = createManager(); + const registry = new AgentRegistry(); + // Task-style spawn: job id == agent id. + manager.register("task", "AgentA", neverResolves, { id: "AgentA", agentId: "AgentA", ownerId: "Main" }); + registerRunningSub(registry, "AgentA"); + // Vibe-style turn job: job id differs from the agent id; linkage via agentId. + manager.register("task", "vibe turn", neverResolves, { id: "vibe-1-t1", agentId: "vibe-1", ownerId: "Main" }); + registerRunningSub(registry, "vibe-1"); + // Woken via irc: running agent with no job at all. + registerRunningSub(registry, "Loner"); + const tool = new JobTool(createToolSession({ manager, registry, agentId: "Main" })); + + const result = await tool.execute("call", { list: true }); + + expect(result.details?.jobs.map(job => job.id).sort()).toEqual(["AgentA", "vibe-1-t1"]); + expect(result.details?.agents?.map(agent => agent.id)).toEqual(["Loner"]); + manager.cancel("AgentA"); + manager.cancel("vibe-1-t1"); + }); + + test("a settled job in retention does not hide its re-woken agent", async () => { + const manager = createManager(); + const registry = new AgentRegistry(); + manager.register("task", "AgentB", async () => "done", { id: "AgentB", agentId: "AgentB", ownerId: "Main" }); + await manager.waitForAll(); + // The agent was re-woken (e.g. via irc) after its job completed. + registerRunningSub(registry, "AgentB"); + const tool = new JobTool(createToolSession({ manager, registry, agentId: "Main" })); + + const result = await tool.execute("call", { list: true }); + + expect(result.details?.jobs.find(job => job.id === "AgentB")?.status).toBe("completed"); + expect(result.details?.agents?.map(agent => agent.id)).toEqual(["AgentB"]); + }); +}); + +describe("job poll with no matching jobs", () => { + test("bare poll with nothing running stays a useless no-op message", async () => { + const tool = new JobTool(createToolSession({ manager: createManager(), agentId: "Main" })); + + const result = await tool.execute("call", {}); + + expect(resultText(result)).toBe("No running background jobs to wait for."); + expect(result.useless).toBe(true); + }); + + test("bare poll reports running agents outside job control", async () => { + const registry = new AgentRegistry(); + registerRunningSub(registry, "Worker"); + const tool = new JobTool(createToolSession({ manager: createManager(), registry, agentId: "Main" })); + + const result = await tool.execute("call", {}); + + const text = resultText(result); + expect(text).toContain("No running background jobs to wait for."); + expect(text).toContain("Worker"); + expect(result.details?.agents?.map(agent => agent.id)).toEqual(["Worker"]); + expect(result.useless).toBeUndefined(); + }); + + test("polling an agent id that has no job explains the agent's state", async () => { + const registry = new AgentRegistry(); + registerRunningSub(registry, "Worker"); + const tool = new JobTool(createToolSession({ manager: createManager(), registry, agentId: "Main" })); + + const result = await tool.execute("call", { poll: ["Worker"] }); + + const text = resultText(result); + expect(text).toContain("No matching jobs found for IDs: Worker"); + expect(text).toContain("running agent with no job entry"); + expect(text).toContain("history://Worker"); + }); +}); diff --git a/packages/coding-agent/test/modes/controllers/event-controller-task-async-updates.test.ts b/packages/coding-agent/test/modes/controllers/event-controller-task-async-updates.test.ts new file mode 100644 index 000000000..185831e4d --- /dev/null +++ b/packages/coding-agent/test/modes/controllers/event-controller-task-async-updates.test.ts @@ -0,0 +1,135 @@ +/** + * Contracts: final async `task` snapshots vs. the tool call's own lifecycle. + * + * A `task` call with background jobs streams `tool_execution_update` frames + * whose `details.async.state` can settle ("completed"/"failed") at any time + * relative to the call's `tool_execution_end` (mixed blocking+async calls run + * their jobs while the call is still executing). + * + * 1. A final async frame arriving BEFORE the call's end is a partial frame: + * the block stays tracked so `tool_execution_end` still delivers the + * terminal result (previously the block was dropped from tracking and the + * real result never rendered — the "disappearing task call"). + * 2. A final async frame arriving AFTER an end that parked the block as + * background ("running") finalizes and untracks it. + */ +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; +import { resetSettingsForTest, Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import type { ToolExecutionComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tool-execution"; +import { TranscriptContainer } from "@oh-my-pi/pi-coding-agent/modes/components/transcript-container"; +import { EventController } from "@oh-my-pi/pi-coding-agent/modes/controllers/event-controller"; +import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme"; +import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types"; +import type { TaskToolDetails } from "@oh-my-pi/pi-coding-agent/task/types"; + +function taskResult(asyncState: "running" | "completed" | "failed" | undefined, text: string) { + const details: TaskToolDetails = { + projectAgentsDir: null, + results: [], + totalDurationMs: 5, + ...(asyncState ? { async: { state: asyncState, jobId: "Job1", type: "task" as const } } : {}), + }; + return { content: [{ type: "text" as const, text }], details }; +} + +describe("EventController task async update finalization", () => { + const sealed: ToolExecutionComponent[] = []; + + beforeEach(async () => { + resetSettingsForTest(); + await Settings.init({ inMemory: true }); + await initTheme(); + }); + + afterEach(() => { + for (const component of sealed.splice(0)) component.seal(); + vi.restoreAllMocks(); + resetSettingsForTest(); + }); + + function createFixture() { + const chatContainer = new TranscriptContainer(); + const pendingTools = new Map(); + const ctx = { + isInitialized: true, + init: vi.fn(async () => {}), + ui: { requestRender: vi.fn(), requestComponentRender: vi.fn() }, + statusLine: { invalidate: vi.fn() }, + updateEditorTopBorder: vi.fn(), + toolOutputExpanded: false, + pendingTools, + chatContainer, + session: { getToolByName: () => undefined, isStreaming: true }, + showWarning: vi.fn(), + viewSession: { getToolByName: () => undefined }, + sessionManager: { getCwd: () => process.cwd() }, + setTodos: vi.fn(), + } as unknown as InteractiveModeContext; + return { controller: new EventController(ctx), pendingTools }; + } + + async function startTask(controller: EventController, pendingTools: Map) { + await controller.handleEvent({ + type: "tool_execution_start", + toolCallId: "tc-task", + toolName: "task", + args: { context: "ctx", tasks: [{ agent: "task", task: "work" }] }, + }); + const component = pendingTools.get("tc-task")!; + sealed.push(component); + return component; + } + + it("keeps the block tracked when a final async frame precedes tool_execution_end", async () => { + const { controller, pendingTools } = createFixture(); + const component = await startTask(controller, pendingTools); + + // The job settled while the call is still executing (mixed call). + await controller.handleEvent({ + type: "tool_execution_update", + toolCallId: "tc-task", + toolName: "task", + args: {}, + partialResult: taskResult("completed", "Background task Job1 complete."), + }); + expect(pendingTools.get("tc-task")).toBe(component); + expect(component.isTranscriptBlockFinalized()).toBe(false); + + // The call's own result still lands and finalizes the block. + await controller.handleEvent({ + type: "tool_execution_end", + toolCallId: "tc-task", + toolName: "task", + result: taskResult("completed", "Inline results + spawned listing."), + isError: false, + }); + expect(pendingTools.has("tc-task")).toBe(false); + expect(component.isTranscriptBlockFinalized()).toBe(true); + }); + + it("finalizes a parked background block when its jobs settle after the end", async () => { + const { controller, pendingTools } = createFixture(); + const component = await startTask(controller, pendingTools); + + await controller.handleEvent({ + type: "tool_execution_end", + toolCallId: "tc-task", + toolName: "task", + result: taskResult("running", "Spawned agent `Job1` (job `Job1`)."), + isError: false, + }); + // Background: kept tracked so later job frames can update it. + expect(pendingTools.get("tc-task")).toBe(component); + expect(component.isTranscriptBlockFinalized()).toBe(true); + + await controller.handleEvent({ + type: "tool_execution_update", + toolCallId: "tc-task", + toolName: "task", + args: {}, + partialResult: taskResult("completed", "Background task Job1 complete."), + }); + expect(pendingTools.has("tc-task")).toBe(false); + expect(component.isTranscriptBlockFinalized()).toBe(true); + }); +}); diff --git a/packages/coding-agent/test/task/task-blocking-split.test.ts b/packages/coding-agent/test/task/task-blocking-split.test.ts new file mode 100644 index 000000000..76a9da175 --- /dev/null +++ b/packages/coding-agent/test/task/task-blocking-split.test.ts @@ -0,0 +1,280 @@ +/** + * Contracts: per-item blocking split in `task` spawns. + * + * An item whose agent type declares `blocking: true` runs inline (the call + * waits on its result); non-blocking items in the same call still spawn as + * background jobs. Previously blocking was all-or-nothing per call: one scout + * in a batch silently dragged every sibling down the sync path — no job ids + * despite the tool description promising them, the whole turn blocked on the + * slowest worker, and every spawn died with the turn signal. + * + * 1. A mixed batch splits: the blocking item's result returns inline, the + * non-blocking item registers a job that keeps running past the return. + * 2. Returned details reflect settled jobs (async.state converges), and + * post-return job updates keep the inline results — they never regress to + * an empty-results skeleton. + * 3. An all-blocking batch stays fully synchronous (no jobs). + * 4. Async schedule failure in a mixed call still returns the inline results + * and reports the failed spawn. + */ +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; +import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async/job-manager"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { AgentLifecycleManager } from "@oh-my-pi/pi-coding-agent/registry/agent-lifecycle"; +import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; +import { TaskTool } from "@oh-my-pi/pi-coding-agent/task"; +import * as discoveryModule from "@oh-my-pi/pi-coding-agent/task/discovery"; +import * as executorModule from "@oh-my-pi/pi-coding-agent/task/executor"; +import type { AgentDefinition, SingleResult, TaskParams, TaskToolDetails } from "@oh-my-pi/pi-coding-agent/task/types"; +import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools"; + +const taskAgent: AgentDefinition = { + name: "task", + description: "General-purpose task agent", + systemPrompt: "You are a task agent.", + source: "bundled", +}; + +const scoutAgent: AgentDefinition = { + name: "scout", + description: "Read-only scout", + systemPrompt: "You are a scout.", + source: "bundled", + blocking: true, +}; + +function createSession(options: { manager?: AsyncJobManager; settings?: Record } = {}): ToolSession { + return { + cwd: "/tmp", + hasUI: false, + settings: Settings.isolated(options.settings ?? { "async.enabled": true, "task.batch": true }), + getSessionFile: () => null, + getSessionSpawns: () => "*", + getAgentId: () => null, + asyncJobManager: options.manager, + } as unknown as ToolSession; +} + +function makeResult(id: string, agent: string, overrides: Partial = {}): SingleResult { + return { + index: 0, + id, + agent, + agentSource: "bundled", + task: "task prompt", + assignment: "Do the thing.", + exitCode: 0, + output: `${id} output.`, + stderr: "", + truncated: false, + durationMs: 5, + tokens: 0, + requests: 1, + ...overrides, + }; +} + +function mockDiscovery(): void { + vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({ + agents: [taskAgent, scoutAgent], + projectAgentsDir: null, + }); +} + +function firstText(result: { content: Array<{ type: string; text?: string }> }): string { + const content = result.content.find(part => part.type === "text"); + return content?.type === "text" ? (content.text ?? "") : ""; +} + +describe("task per-item blocking split", () => { + const managers: AsyncJobManager[] = []; + + function createManager(): AsyncJobManager { + const manager = new AsyncJobManager({ onJobComplete: () => {} }); + managers.push(manager); + return manager; + } + + beforeEach(() => { + AgentRegistry.resetGlobalForTests(); + AgentLifecycleManager.resetGlobalForTests(); + }); + + afterEach(async () => { + vi.restoreAllMocks(); + for (const manager of managers.splice(0)) { + await manager.dispose({ timeoutMs: 1000 }); + } + AgentLifecycleManager.resetGlobalForTests(); + AgentRegistry.resetGlobalForTests(); + }); + + it("runs blocking items inline while non-blocking siblings spawn as jobs", async () => { + mockDiscovery(); + const gates = new Map>(); + const started: string[] = []; + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { + const id = options.id ?? "?"; + started.push(id); + const gate = Promise.withResolvers(); + gates.set(id, gate); + await gate.promise; + return makeResult(id, options.agent.name); + }); + + const manager = createManager(); + const tool = await TaskTool.create(createSession({ manager })); + + const executePromise = tool.execute("tc-mixed", { + context: "ctx", + tasks: [ + { name: "ScoutOne", agent: "scout", task: "Research A." }, + { name: "WorkerOne", agent: "task", task: "Build B." }, + ], + } as TaskParams); + + // Both spawns run concurrently: the async job is registered immediately, + // the blocking scout gates the call's return. + const deadline = Date.now() + 2_000; + while (started.length < 2) { + if (Date.now() > deadline) throw new Error(`spawns never started: ${JSON.stringify(started)}`); + await Bun.sleep(5); + } + const workerJob = manager.getJob("WorkerOne"); + expect(workerJob).toBeDefined(); + expect(workerJob!.status).toBe("running"); + + // Releasing only the scout settles the call; the worker keeps running. + gates.get("ScoutOne")!.resolve(); + const result = await executePromise; + + const text = firstText(result); + expect(text).toContain('id="ScoutOne"'); + expect(text).toContain("ScoutOne output."); + expect(text).toContain("Spawned 1 background agent"); + expect(text).toContain("- `WorkerOne` (job `WorkerOne`)"); + expect(text).not.toContain("WorkerOne output."); + + expect(result.details?.results.map(r => r.id)).toEqual(["ScoutOne"]); + expect(result.details?.async?.state).toBe("running"); + const progressById = new Map(result.details?.progress?.map(p => [p.id, p.status])); + expect(progressById.get("ScoutOne")).toBe("completed"); + expect(progressById.get("WorkerOne")).toBe("running"); + expect(manager.getJob("WorkerOne")!.status).toBe("running"); + + gates.get("WorkerOne")!.resolve(); + await manager.getJob("WorkerOne")!.promise; + expect(manager.getJob("WorkerOne")!.status).toBe("completed"); + }); + + it("keeps inline results in post-return job updates and converges async state", async () => { + mockDiscovery(); + const gates = new Map>(); + const started: string[] = []; + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { + const id = options.id ?? "?"; + started.push(id); + const gate = Promise.withResolvers(); + gates.set(id, gate); + await gate.promise; + return makeResult(id, options.agent.name); + }); + + const manager = createManager(); + const tool = await TaskTool.create(createSession({ manager })); + + const updates: Array<{ text: string; details: TaskToolDetails }> = []; + const executePromise = tool.execute( + "tc-mixed-settle", + { + context: "ctx", + tasks: [ + { name: "ScoutTwo", agent: "scout", task: "Research." }, + { name: "WorkerTwo", agent: "task", task: "Build." }, + ], + } as TaskParams, + undefined, + update => { + if (update.details) updates.push({ text: firstText(update), details: update.details }); + }, + ); + + const deadline = Date.now() + 2_000; + while (started.length < 2) { + if (Date.now() > deadline) throw new Error(`spawns never started: ${JSON.stringify(started)}`); + await Bun.sleep(5); + } + + // The async job settles BEFORE the blocking subset: the returned result + // must already report the converged async state, not a stale "running". + gates.get("WorkerTwo")!.resolve(); + await manager.getJob("WorkerTwo")!.promise; + gates.get("ScoutTwo")!.resolve(); + const result = await executePromise; + + expect(result.details?.async?.state).toBe("completed"); + expect(result.details?.results.map(r => r.id)).toEqual(["ScoutTwo"]); + + // The job's completion update carries the shared aggregate — inline + // results included — never an empty-results skeleton, and exactly once. + const completionUpdates = updates.filter(u => u.text.includes("Background task WorkerTwo complete.")); + expect(completionUpdates).toHaveLength(1); + }); + + it("keeps an all-blocking batch fully synchronous", async () => { + mockDiscovery(); + const executed: string[] = []; + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { + executed.push(options.id ?? "?"); + return makeResult(options.id ?? "?", options.agent.name); + }); + + const manager = createManager(); + const tool = await TaskTool.create(createSession({ manager })); + + const result = await tool.execute("tc-all-blocking", { + context: "ctx", + tasks: [ + { name: "ScoutA", agent: "scout", task: "Research A." }, + { name: "ScoutB", agent: "scout", task: "Research B." }, + ], + } as TaskParams); + + expect(executed.sort()).toEqual(["ScoutA", "ScoutB"]); + expect(result.details?.async).toBeUndefined(); + expect(result.details?.results.map(r => r.id).sort()).toEqual(["ScoutA", "ScoutB"]); + expect(manager.getAllJobs()).toHaveLength(0); + }); + + it("returns inline results and reports the failure when async scheduling fails", async () => { + mockDiscovery(); + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => + makeResult(options.id ?? "?", options.agent.name), + ); + + // A running non-queued filler job exhausts maxRunningJobs, so the + // worker spawn's registration throws while the scout still runs inline. + const filler = Promise.withResolvers(); + const manager = new AsyncJobManager({ maxRunningJobs: 1, onJobComplete: () => {} }); + managers.push(manager); + manager.register("bash", "filler", () => filler.promise, { id: "filler" }); + + const tool = await TaskTool.create(createSession({ manager })); + const result = await tool.execute("tc-mixed-schedfail", { + context: "ctx", + tasks: [ + { name: "ScoutThree", agent: "scout", task: "Research." }, + { name: "WorkerThree", agent: "task", task: "Build." }, + ], + } as TaskParams); + filler.resolve("done"); + + const text = firstText(result); + expect(text).toContain('id="ScoutThree"'); + expect(text).toContain("Failed to schedule 1 spawn"); + expect(text).toContain("WorkerThree"); + expect(result.details?.results.map(r => r.id)).toEqual(["ScoutThree"]); + expect(result.details?.async?.state).toBe("failed"); + expect(manager.getJob("WorkerThree")).toBeUndefined(); + }); +});