From 890dd885976dc2540ea485839559a8a5097927cc Mon Sep 17 00:00:00 2001 From: TechDufus Date: Thu, 23 Jul 2026 16:41:00 -0500 Subject: [PATCH] fix(coding-agent): make mixed-agent task batches atomic and inspectable --- packages/coding-agent/CHANGELOG.md | 1 + .../src/config/settings-schema.ts | 2 +- packages/coding-agent/src/task/index.ts | 202 +++++++++--------- .../coding-agent/test/task/task-batch.test.ts | 86 +++++++- .../test/task/task-preflight.test.ts | 50 +++-- .../test/task/wire-schema.test.ts | 30 ++- 6 files changed, 242 insertions(+), 129 deletions(-) diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 7677e9bb7..7dfeda5fe 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -39,6 +39,7 @@ - Fixed bash internal-URL expansion skipping unquoted `skill://` (and other supported schemes) inside a legacy backtick command substitution nested directly in double quotes (e.g. ``echo "`cat skill://valid-skill/SKILL.md`"``); `isInsideShellQuote` now treats `` ` `` as an expansion-context boundary like `$()`, including `$()`/backtick nesting in either order, while single-quoted and escaped-backtick text stay literal ([#5645](https://github.com/can1357/oh-my-pi/issues/5645)). - Fixed `omp say` playing no audio for a short single-segment clip on hosts where the first streaming backend (the bundled ffmpeg built without pulse/alsa output) spawns then exits nonzero: the pipe write succeeds before that death and `player.end()` has already closed the input, so neither the broken-pipe replay nor the early-exit handler advanced to `paplay`/`aplay`. `StreamingAudioPlayer` now retains the utterance PCM and, when the streaming backend exits nonzero, replays it through per-file playback so short clips still reach the speakers ([#5875](https://github.com/can1357/oh-my-pi/issues/5875)). - Fixed mid-session `memory.backend` changes leaving runtime state, tools, listeners, and prompt context on different backends; Mnemopi clear/enqueue now rehydrate listeners, and legacy `memories.enabled` no longer activates the local pipeline after migration ([#5638](https://github.com/can1357/oh-my-pi/issues/5638)). +- Fixed mixed-agent `task` batches launching valid siblings when any item failed agent preflight: all effective item policies now validate before job registration or synchronous execution, approval details show effective-agent counts, and the batch setting copy reflects the current per-item agent shape. - Fixed `error.notify` raising a "Stopped with error" toast for provider failures while an auto-retry or async-delivery continuation was pending; the toast now waits for the true terminal settle. - Fixed concurrent MCP config mutations losing updates and racing on a shared temp path: every `mcp.json` read-modify-write (add/update/remove server, disabled/force-enabled lists) is now serialized under a per-file lock, and each atomic write uses a unique temp file so overlapping writers no longer rename each other's `.tmp` out from under them (ENOENT or clobbered config) — reachable in-process via the fire-and-forget extensions-dashboard toggle and across processes on a shared `~/.omp/mcp.json` ([#4104](https://github.com/can1357/oh-my-pi/issues/4104)). - Fixed transient provider stream stalls after tool calls failing to auto-retry even when every call already had a tool result, including synthetic `executed:false` results from OpenAI-completions stalls ([#6414](https://github.com/can1357/oh-my-pi/issues/6414)). diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index 37c465faa..a975c8171 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -4284,7 +4284,7 @@ export const SETTINGS_SCHEMA = { group: "Subagents", label: "Batch Task Calls", description: - "Switch the task tool to its batch shape: one call carries { agent, context, tasks[] } — one subagent per item (with per-item isolation) and a required shared context prepended to every assignment. With async.enabled=true, each spawn runs as an independent background agent with the normal idle/parked lifecycle; otherwise the call blocks for merged results. Disable to restore the flat single-spawn schema.", + "Switch the task tool to its batch shape: one call carries { context, tasks[] } — one subagent per item, with an optional per-item agent (defaulting to the session spawn-policy agent), per-item isolation, and a required shared context prepended to every assignment. With async.enabled=true, each spawn runs as an independent background agent with the normal idle/parked lifecycle; otherwise the call blocks for merged results. Disable to restore the flat single-spawn schema.", }, }, diff --git a/packages/coding-agent/src/task/index.ts b/packages/coding-agent/src/task/index.ts index e3046ba92..c0d33da5e 100644 --- a/packages/coding-agent/src/task/index.ts +++ b/packages/coding-agent/src/task/index.ts @@ -530,19 +530,35 @@ export class TaskTool implements AgentTool 0) { + const defaultAgent = resolveSpawnPolicy(this.session.getSessionSpawns()).defaultAgent; + const effectiveAgent = (item: unknown): string => { + if (item && typeof item === "object" && "agent" in item) { + const agent = item.agent; + if (typeof agent === "string" && agent.trim()) return agent.trim(); + } + return defaultAgent; + }; + const agentCounts = new Map(); + for (const item of tasks) { + const agent = effectiveAgent(item); + agentCounts.set(agent, (agentCounts.get(agent) ?? 0) + 1); } - if (typeof firstTask.agent === "string" && firstTask.agent.trim()) { - lines.push(`Agent: ${truncateForPrompt(firstTask.agent)}`); - } - const itemModel = formatModelForApproval(firstTask.model); - if (itemModel) lines.push(`Model: ${itemModel}`); - if (typeof firstTask.task === "string") { - lines.push(`Task:\n${truncateForPrompt(firstTask.task)}`); + const agentSummary = [...agentCounts].map(([agent, count]) => `${agent} ×${count}`).join(", "); + lines.push(`Batch agents: ${truncateForPrompt(agentSummary)}`); + + const firstTask = tasks[0]; + if (firstTask && typeof firstTask === "object") { + if ("name" in firstTask && typeof firstTask.name === "string" && firstTask.name.trim()) { + lines.push(`Name: ${truncateForPrompt(firstTask.name)}`); + } + lines.push(`Agent: ${truncateForPrompt(effectiveAgent(firstTask))}`); + const itemModel = formatModelForApproval("model" in firstTask ? firstTask.model : undefined); + if (itemModel) lines.push(`Model: ${itemModel}`); + if ("task" in firstTask && typeof firstTask.task === "string") { + lines.push(`Task:\n${truncateForPrompt(firstTask.task)}`); + } } if (tasks.length > 1) { lines.push(`+${tasks.length - 1} more task${tasks.length === 2 ? "" : "s"}`); @@ -686,27 +702,53 @@ export class TaskTool implements AgentTool spawnParamsFor(params, item, defaultAgent)); const resolvedAgents = normalizedSpawnParams.map(spawn => spawn.agent ?? defaultAgent); + // Resolve every item before choosing an execution path. No executor or + // job manager may observe a batch unless every effective policy is valid. + const preflights = await Promise.all( + normalizedSpawnParams.map(async spawn => { + try { + return { policy: await this.#resolveSpawnPreflight(spawn) }; + } catch (error) { + return { error: error instanceof StructuredSubagentError ? error.message : String(error) }; + } + }), + ); + const preflightFailures = preflights + .map((preflight, index) => ("error" in preflight ? { index, error: preflight.error } : undefined)) + .filter((failure): failure is { index: number; error: string } => failure !== undefined); + if (preflightFailures.length > 0) { + if (!batchEnabled) { + return createTaskModeError(`Task execution failed: ${preflightFailures[0]!.error}`); + } + return createTaskModeError( + preflightFailures + .map(({ index, error }) => { + const item = spawnItems[index]!; + return `Task ${item.name?.trim() || `#${index + 1}`} failed preflight: ${error}`; + }) + .join("\n"), + ); + } + const policies = preflights.map(preflight => preflight.policy!); + const itemBlocking = policies.map(policy => policy.effectiveAgent.blocking === true); + // 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 provisionalBlocking = 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 provisionalAsyncItems = manager ? spawnItems.filter((_, index) => !provisionalBlocking[index]) : []; + 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); - if (!manager || provisionalAsyncItems.length === 0) { + if (!manager || asyncItems.length === 0) { // Sync fallback: async execution disabled, orphaned host that never // wired a job manager, or every item's agent type declares - // `blocking: true`. `runStructuredSubagent` performs its own shared - // preflight before reserving an id in these inline paths. + // `blocking: true`. if (asyncEnabled && !this.session.asyncJobManager) { logger.warn("task: no AsyncJobManager registered; falling back to sync execution"); } @@ -714,7 +756,7 @@ export class TaskTool implements AgentTool { - try { - return { policy: await this.#resolveSpawnPreflight(spawn) }; - } catch (error) { - return { error: error instanceof StructuredSubagentError ? error.message : String(error) }; - } - }), - ); - const preflightFailures = preflights - .map((preflight, index) => ("error" in preflight ? { index, error: preflight.error } : undefined)) - .filter((failure): failure is { index: number; error: string } => failure !== undefined); - const renderPreflightFailures = () => - preflightFailures - .map(({ index, error }) => { - const item = spawnItems[index]!; - return `Task ${item.name?.trim() || `#${index + 1}`} failed preflight: ${error}`; - }) - .join("\n"); - if (preflightFailures.length === spawnItems.length) { - return createTaskModeError(renderPreflightFailures()); - } - - const validIndices = preflights.flatMap((preflight, index) => (preflight.policy ? [index] : [])); - const validSpawns = validIndices.map(index => ({ item: spawnItems[index]!, index })); - const itemBlocking = preflights.map(preflight => preflight.policy?.effectiveAgent.blocking === true); - const asyncItems = validIndices.filter(index => !itemBlocking[index]).map(index => spawnItems[index]!); // 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 advisory = this.session.suppressSpawnAdvisory ? undefined : composeSpawnAdvisory({ - agents: validIndices.map(index => resolvedAgents[index]!), + agents: resolvedAgents, items: asyncItems, depthCapacity, ircEnabled, @@ -798,24 +810,15 @@ export class TaskTool implements AgentTool): AgentToolResult => { - if (preflightFailures.length === 0) return result; - const failures = renderPreflightFailures(); - let prepended = false; - const content = result.content.map(part => { - if (!prepended && part.type === "text" && typeof part.text === "string") { - prepended = true; - return { ...part, text: `${failures}\n\n${part.text}` }; - } - return part; - }); - if (!prepended) content.unshift({ type: "text", text: failures }); - return { ...result, content }; - }; if (asyncItems.length === 0) { - return withPreflightFailures( - withAdvisory( - await this.#executeSyncFanout(toolCallId, params, validSpawns, defaultAgent, signal, onUpdate), + return withAdvisory( + await this.#executeSyncFanout( + toolCallId, + params, + spawnItems.map((item, index) => ({ item, index })), + defaultAgent, + signal, + onUpdate, ), ); } @@ -835,12 +838,9 @@ export class TaskTool implements AgentTool = []; - for (const index of validIndices) { - const item = spawnItems[index]!; + for (const [index, item] of spawnItems.entries()) { const agentType = resolvedAgents[index]!; - const preflight = preflights[index]!; - const policy = preflight.policy; - if (!policy) continue; + const policy = policies[index]!; const agentSource = policy.agent.source; const agentId = await outputManager.allocate(item.name?.trim() || generateTaskName()); const assignment = (item.task ?? "").trim(); @@ -877,7 +877,7 @@ export class TaskTool implements AgentTool `- \`${agentId}\` (job \`${jobId}\`)`).join("\n"); onUpdate?.({ content: [{ type: "text", text: `Spawned ${started.length} agents...` }], details: buildAsyncDetails(), }); - return withPreflightFailures( - withAdvisory({ - content: [ - { - type: "text", - text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result auto-delivers on yield unless a settled \`hub jobs\`/\`wait\` snapshot consumes it first.\n${startedListing}\n${coordinationHint}`, - }, - ], - details: buildAsyncDetails(), - }), - ); + return withAdvisory({ + content: [ + { + type: "text", + text: `Spawned ${started.length} background agents using ${agentLabel}.${scheduleFailureSummary} Each result auto-delivers on yield unless a settled \`hub jobs\`/\`wait\` snapshot consumes it first.\n${startedListing}\n${coordinationHint}`, + }, + ], + details: buildAsyncDetails(), + }); } // Mixed call: the async jobs above already run detached; the blocking @@ -1052,12 +1048,10 @@ export class TaskTool implements AgentTool section.trim().length > 0) .join("\n\n"); - return withPreflightFailures( - withAdvisory({ - content: [{ type: "text", text: text.length > 0 ? text : "No results." }], - details: buildAsyncDetails(), - }), - ); + return withAdvisory({ + content: [{ type: "text", text: text.length > 0 ? text : "No results." }], + details: buildAsyncDetails(), + }); } /** diff --git a/packages/coding-agent/test/task/task-batch.test.ts b/packages/coding-agent/test/task/task-batch.test.ts index b56666738..6c6b07110 100644 --- a/packages/coding-agent/test/task/task-batch.test.ts +++ b/packages/coding-agent/test/task/task-batch.test.ts @@ -81,9 +81,9 @@ function makeResult(id: string, overrides: Partial = {}): SingleRe }; } -function mockDiscovery(agent: AgentDefinition = taskAgent): void { +function mockDiscovery(agent: AgentDefinition | AgentDefinition[] = taskAgent): void { vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({ - agents: [agent], + agents: Array.isArray(agent) ? agent : [agent], projectAgentsDir: null, }); } @@ -395,6 +395,88 @@ describe("task.batch spawning", () => { for (const spawn of seen) expect(spawn.parentAgentId).toBe("ParentA"); }); + it("routes each mixed-agent item through its selected definition while preserving caller overrides", async () => { + const scoutSchema = { type: "object", properties: { findings: { type: "array" } } }; + const reviewerSchema = { type: "object", properties: { verdict: { type: "string" } } }; + const callerSchema = { type: "object", properties: { approved: { type: "boolean" } } }; + const scoutAgent: AgentDefinition = { + ...taskAgent, + name: "scout", + description: "Read-only scout", + systemPrompt: "Investigate the assigned target.", + tools: ["read"], + model: ["anthropic/claude-haiku-4-5:low"], + output: scoutSchema, + }; + const reviewerAgent: AgentDefinition = { + ...taskAgent, + name: "reviewer", + description: "Code review specialist", + systemPrompt: "Review the assigned target.", + tools: ["read", "bash"], + model: ["anthropic/claude-sonnet-4-6:medium"], + output: reviewerSchema, + }; + mockDiscovery([scoutAgent, reviewerAgent]); + + const seen: Array<{ + id?: string; + agent: AgentDefinition; + modelOverride?: string | string[]; + outputSchema?: unknown; + outputSchemaSource?: "caller" | "agent" | "session" | "none"; + outputSchemaOverridesAgent?: boolean; + }> = []; + vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { + seen.push({ + id: options.id, + agent: options.agent, + modelOverride: options.modelOverride, + outputSchema: options.outputSchema, + outputSchemaSource: options.outputSchemaSource, + outputSchemaOverridesAgent: options.outputSchemaOverridesAgent, + }); + return makeResult(options.id ?? "?", { agent: options.agent.name }); + }); + + const manager = createManager(); + const tool = await TaskTool.create( + createSession({ manager, settings: { "async.enabled": true, "task.batch": true } }), + ); + const result = await tool.execute("tc-mixed-agents", { + context: "Shared routing context.", + tasks: [ + { name: "Scout", agent: "scout", task: "Investigate." }, + { + name: "Review", + agent: "reviewer", + task: "Review.", + model: "openai-codex/gpt-5.6-sol:high", + outputSchema: callerSchema, + }, + ], + } as TaskParams); + + expect(getFirstText(result)).toContain("Spawned 2 background agents"); + await Promise.all([manager.getJob("Scout")!.promise, manager.getJob("Review")!.promise]); + + const byId = new Map(seen.map(spawn => [spawn.id, spawn])); + const scoutSpawn = byId.get("Scout"); + const reviewerSpawn = byId.get("Review"); + expect(scoutSpawn?.agent).toBe(scoutAgent); + expect(scoutSpawn?.agent.tools).toEqual(["read"]); + expect(scoutSpawn?.modelOverride).toEqual(["anthropic/claude-haiku-4-5:low"]); + expect(scoutSpawn?.outputSchema).toBe(scoutSchema); + expect(scoutSpawn?.outputSchemaSource).toBe("agent"); + expect(scoutSpawn?.outputSchemaOverridesAgent).toBe(false); + expect(reviewerSpawn?.agent).toBe(reviewerAgent); + expect(reviewerSpawn?.agent.tools).toEqual(["read", "bash"]); + expect(reviewerSpawn?.modelOverride).toEqual(["openai-codex/gpt-5.6-sol:high"]); + expect(reviewerSpawn?.outputSchema).toBe(callerSchema); + expect(reviewerSpawn?.outputSchemaSource).toBe("caller"); + expect(reviewerSpawn?.outputSchemaOverridesAgent).toBe(true); + }); + it("treats a one-item batch as a single spawn and forwards context", async () => { mockDiscovery(); let capturedContext: string | undefined; diff --git a/packages/coding-agent/test/task/task-preflight.test.ts b/packages/coding-agent/test/task/task-preflight.test.ts index 08dd7c7a2..97c70880b 100644 --- a/packages/coding-agent/test/task/task-preflight.test.ts +++ b/packages/coding-agent/test/task/task-preflight.test.ts @@ -111,17 +111,42 @@ describe("task async preflight", () => { }, ); - it("reports an invalid batch item synchronously while launching its valid sibling", async () => { + it("rejects an invalid async batch atomically before dispatching any item", async () => { mockDiscovery(); - const seen: string[] = []; - vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => { - seen.push(options.id ?? ""); - return resultFor(options.id ?? ""); - }); + const runSubprocess = vi.spyOn(executorModule, "runSubprocess").mockResolvedValue(resultFor("unexpected")); const jobs = manager(); + const register = vi.spyOn(jobs, "register"); const tool = await TaskTool.create(createSession({ manager: jobs, settings: { "task.batch": true } })); const result = await tool.execute("mixed-preflight", { + context: "Shared context.", + tasks: [ + { name: "Invalid", agent: "missing", task: "Do invalid work." }, + { name: "AlsoInvalid", agent: "also-missing", task: "Do more invalid work." }, + { name: "Valid", agent: "task", task: "Do valid work." }, + ], + } as TaskParams); + + const text = textOf(result); + expect(text).toContain('Task Invalid failed preflight: Unknown agent "missing"'); + expect(text).toContain('Task AlsoInvalid failed preflight: Unknown agent "also-missing"'); + expect(register).not.toHaveBeenCalled(); + expect(runSubprocess).not.toHaveBeenCalled(); + expect(jobs.getJob("Invalid")).toBeUndefined(); + expect(jobs.getJob("AlsoInvalid")).toBeUndefined(); + expect(jobs.getJob("Valid")).toBeUndefined(); + }); + + it("rejects an invalid synchronous batch before running any item", async () => { + mockDiscovery(); + const runSubprocess = vi.spyOn(executorModule, "runSubprocess").mockResolvedValue(resultFor("unexpected")); + const jobs = manager(); + const register = vi.spyOn(jobs, "register"); + const tool = await TaskTool.create( + createSession({ manager: jobs, settings: { "async.enabled": false, "task.batch": true } }), + ); + + const result = await tool.execute("sync-preflight", { context: "Shared context.", tasks: [ { name: "Invalid", agent: "missing", task: "Do invalid work." }, @@ -129,15 +154,10 @@ describe("task async preflight", () => { ], } as TaskParams); - const text = textOf(result); - expect(text).toContain('Task Invalid failed preflight: Unknown agent "missing"'); - expect(text).toContain("Spawned agent `Valid`"); - expect(text.indexOf("Task Invalid failed preflight")).toBeLessThan(text.indexOf("Spawned agent `Valid`")); + expect(textOf(result)).toContain('Task Invalid failed preflight: Unknown agent "missing"'); + expect(register).not.toHaveBeenCalled(); + expect(runSubprocess).not.toHaveBeenCalled(); expect(jobs.getJob("Invalid")).toBeUndefined(); - const valid = jobs.getJob("Valid"); - expect(valid).toBeDefined(); - await valid!.promise; - expect(valid!.status).toBe("completed"); - expect(seen).toEqual(["Valid"]); + expect(jobs.getJob("Valid")).toBeUndefined(); }); }); diff --git a/packages/coding-agent/test/task/wire-schema.test.ts b/packages/coding-agent/test/task/wire-schema.test.ts index fa2fd3603..f5c2d668b 100644 --- a/packages/coding-agent/test/task/wire-schema.test.ts +++ b/packages/coding-agent/test/task/wire-schema.test.ts @@ -107,14 +107,14 @@ describe("task approval details surface the dispatch", () => { vi.restoreAllMocks(); }); - async function makeTool(): Promise { + async function makeTool(spawns = "*"): Promise { vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({ agents: [], projectAgentsDir: null }); return TaskTool.create({ cwd: "/tmp", hasUI: false, - settings: Settings.isolated({ "task.isolation.mode": "none", "task.batch": false }), + settings: Settings.isolated({ "task.isolation.mode": "none", "task.batch": true }), getSessionFile: () => null, - getSessionSpawns: () => "*", + getSessionSpawns: () => spawns, } as unknown as ToolSession); } @@ -132,14 +132,13 @@ describe("task approval details surface the dispatch", () => { expect(lines).toContain("Task:\naudit the auth module"); }); - it("surfaces the first batch item and the remainder count", async () => { - const tool = await makeTool(); + it("summarizes a homogeneous batch whose agents use the session default", async () => { + const tool = await makeTool("scout,reviewer"); const lines = tool.formatApprovalDetails({ context: "shared background", tasks: [ { name: "DbMigrator", - agent: "sonic", model: ["anthropic/claude-sonnet-4", "openai/gpt-5"], task: "migrate the schema", }, @@ -147,10 +146,27 @@ describe("task approval details surface the dispatch", () => { ], }); expect(lines).toContain("Context:\nshared background"); + expect(lines).toContain("Batch agents: scout ×2"); expect(lines).toContain("Name: DbMigrator"); - expect(lines).toContain("Agent: sonic"); + expect(lines).toContain("Agent: scout"); expect(lines).toContain("Model: anthropic/claude-sonnet-4 → openai/gpt-5"); expect(lines).toContain("Task:\nmigrate the schema"); expect(lines).toContain("+1 more task"); }); + + it("summarizes mixed effective agents and safely renders partial batch items", async () => { + const tool = await makeTool("scout,reviewer"); + const lines = tool.formatApprovalDetails({ + tasks: [ + { name: "DefaultScout", task: "map the flow" }, + { agent: " reviewer ", task: "review it" }, + ], + }); + expect(lines).toContain("Batch agents: scout ×1, reviewer ×1"); + expect(lines).toContain("Name: DefaultScout"); + expect(lines).toContain("Agent: scout"); + expect(lines.join("\n")).not.toContain("undefined"); + + expect(() => tool.formatApprovalDetails({ tasks: [undefined, { agent: "reviewer" }] })).not.toThrow(); + }); });