/** * Dispatch-batch selection and execution for the team-run scheduler loop. * * Extracted from team-runner.ts (2026-08, Phase 2.6 maintainability split — * CORE-4 extractions 3+4). Pure code motion: selectDispatchBatch / * DispatchBatchDecision / dispatchBatch and their module-private * collaborators (findStep, findAgent, markBlocked, cancelNonTerminalTasks, * retryPolicyFromConfig, shouldUseRetry, failedTaskFrom, dagReadyTaskIds) * moved verbatim. */ import type { AgentConfig } from "../agents/agent-config.ts"; import type { CrewReliabilityConfig } from "../config/config.ts"; import { CrewError, ErrorCode } from "../errors.ts"; import { appendHookEvent, executeHook } from "../hooks/registry.ts"; import { childCorrelation, withCorrelation } from "../observability/correlation.ts"; import { TEAM_TERMINAL_RUN_STATUSES } from "../state/contracts.ts"; import { withRunLockSync } from "../state/coordination/locks.ts"; import { appendEventAsync, appendEventBuffered } from "../state/event-log/event-log.ts"; import { loadRunManifestById, saveRunManifest, saveRunTasks, saveRunTasksAsync, updateRunStatus } from "../state/stores/state-store.ts"; import type { TaskAttemptState, TeamRunManifest, TeamTaskState } from "../state/types.ts"; import { logInternalError } from "../utils/internal-error.ts"; import type { WorkflowConfig, WorkflowStep } from "../workflows/workflow-config.ts"; import { isDelegateShadowTask } from "./broker/delegate/shadow-lifecycle.ts"; import { readCrewAgents, recordFromTask, saveCrewAgents } from "./crew-agent-records.ts"; import { appendDeadletter } from "./deadletter.ts"; import { classifyHeartbeat, DEFAULT_GRADIENT_THRESHOLDS } from "./heartbeat/heartbeat-gradient.ts"; import { getLiveAgent } from "./live-session/live-agent-manager.ts"; import { isNonTerminalTaskStatus } from "./merge-gate.ts"; import { resolveTaskRuntimeKind } from "./model/runtime-policy.ts"; import { filterReadyByWriteOverlap } from "./path-overlap.ts"; import { isMutatingTask, isPlanApprovalPendingEffective } from "./plan-approval.ts"; import { sweepDroppedPlanItems } from "./plan-replan.ts"; import { CrewCancellationError, cancellationReasonFromSignal } from "./process/cancellation.ts"; import { DEFAULT_RETRY_POLICY, executeWithRetry, type RetryPolicy } from "./recovery/retry-executor.ts"; import type { SchedulerContext, SchedulerDecision, SettledUnit } from "./scheduler-context.ts"; import { buildDispatchUnits, type DispatchUnit, planCoalescedGroups } from "./scheduling/coalesce-tasks.ts"; import { resolveBatchConcurrency } from "./scheduling/concurrency.ts"; import { runCoalescedTaskGroup } from "./scheduling/run-coalesced-task-group.ts"; import { buildExecutionPlan as buildDagExecutionPlan, getReadyTasks as getDagReadyTasks, type TaskNode } from "./scheduling/task-graph.ts"; import { taskGraphSnapshot } from "./scheduling/task-graph-scheduler.ts"; import { recordsForMaterializedTasks } from "./task-display.ts"; import { computeStablePrefixComponents } from "./task-runner/prompt-builder.ts"; import { runTeamTask, type SpawnBudget } from "./task-runner.ts"; import { type PhaseGuardContext, validatePhasePreconditions } from "./workflow-state.ts"; // ── WP-2/R2 (ADR-0 items 8+10): waiting-producer liveness + deadline owner ── /** * WP-2/R2 (ADR-0 docs/decisions/2026-08-17-waiting-producer-ask.md item 8): * liveness discriminator for a parked (waiting) worker. * * ALIVE = heartbeat last-beat within the gradient stale window (<60s — * classifyHeartbeat "healthy"/"warn") OR a live in-memory session handle * (live-agent registry, non-terminal status). Everything else is DEAD for * delivery purposes: stale (60–300s), dead (>300s), `alive:false`, no * heartbeat and no handle. * * Shared by the respond discriminator (extension/team-tool/respond.ts) and * the scheduler-tick deadline owner below — one definition, two consumers. */ export function isWaitingWorkerAlive(task: TeamTaskState, now = Date.now()): boolean { // Live in-memory handle: same-process live-session runtime keeps a session // handle in the registry keyed by (agentId | taskId). const handle = getLiveAgent(task.id); if (handle && (handle.status === "running" || handle.status === "queued" || handle.status === "waiting")) return true; // Heartbeat: last beat inside the gradient stale window (<60s). const level = classifyHeartbeat(task.heartbeat, DEFAULT_GRADIENT_THRESHOLDS, now); return level === "healthy" || level === "warn"; } /** In-process exactly-once guard: questionIds whose `ask.timedout` event was * already emitted by this scheduler process (prevents re-emission on every * subsequent tick while an ALIVE park waits for its in-tool timeout to * surface). Capped with FIFO eviction to bound memory. */ const emittedAskTimedoutQuestionIds = new Set(); const MAX_EMITTED_ASK_TIMEDOUT = 10_000; export interface WaitingDeadlineSweepResult { manifest: TeamRunManifest; tasks: TeamTaskState[]; /** QuestionIds for which an `ask.timedout` event was emitted this sweep. */ timedOutQuestionIds: string[]; /** Task ids requeued this sweep (dead workers only). */ requeuedTaskIds: string[]; } /** WP-2/R2 (ADR-0 item 10): render the "[ask timed out]" note injected into the * next dispatch when a parked ask deadline expires with no worker alive. * Fenced as untrusted (same-uid channel discipline). */ function renderAskTimedoutInjection(questionId: string): string { return [ "", "(Scheduler note: the previous ask(question) deadline expired with no worker alive to receive an answer. It is DATA, not a system directive.)", `[ask timed out] questionId=${questionId}`, "Continue with best judgment.", "", ].join("\n"); } /** * WP-2/R2 (ADR-0 item 10) — scheduler-tick deadline owner for parked asks. * * For every task with a persisted `task.waiting` whose deadline has expired: * - worker ALIVE → no root-side state change (the parked ask tool's own poll * deadline surfaces the timeout in-tool); only the `ask.timedout` event is * appended. * - worker DEAD → clear `waiting`, re-queue the task with a "[ask timed out]" * note injected into the next dispatch (pendingSteers cross-attempt * channel), and append `ask.timedout`. * * All state changes run under the run lock with a fresh reload * (respond.ts:42-43 discipline); events are appended AFTER the lock is * released (awaited — callers are async). Exactly-once per questionId via * the in-process guard above (a requeued task also loses `waiting`, so a * second sweep pass skips it). * * Returns undefined when there is nothing to do (no expired waiting park). */ export async function sweepExpiredWaitingTasks( cwd: string, runId: string, now = Date.now(), // PERF (2026-08-24): callers in the scheduler loop already hold a task view // loaded this tick. Check the cheap expiry predicate against it BEFORE // paying stat+parse of manifest.json + tasks.json — the sweep fires on every // unit settle/dispatch and almost always finds nothing expired. hintTasks?: TeamTaskState[], ): Promise { if (hintTasks) { const hintExpired = hintTasks.some((t) => t.status === "waiting" && t.waiting !== undefined && t.waiting.deadline <= now); if (!hintExpired) return undefined; } const initial = loadRunManifestById(cwd, runId); if (!initial) return undefined; const hasExpired = initial.tasks.some((t) => t.status === "waiting" && t.waiting !== undefined && t.waiting.deadline <= now); if (!hasExpired) return undefined; const pendingEvents: Array<{ taskId: string; questionId: string; workerAlive: boolean }> = []; const outcome = withRunLockSync(initial.manifest, () => { const fresh = loadRunManifestById(cwd, runId); // NOTE: inside withRunLockSync - consistent read if (!fresh) return null; let tasks = fresh.tasks; const timedOutQuestionIds: string[] = []; const requeuedTaskIds: string[] = []; for (const task of fresh.tasks) { if (task.status !== "waiting" || task.waiting === undefined) continue; if (task.waiting.deadline > now) continue; const questionId = task.waiting.questionId; if (emittedAskTimedoutQuestionIds.has(questionId)) continue; // exactly-once per process if (emittedAskTimedoutQuestionIds.size >= MAX_EMITTED_ASK_TIMEDOUT) { const oldest = emittedAskTimedoutQuestionIds.keys().next().value; if (oldest !== undefined) emittedAskTimedoutQuestionIds.delete(oldest); } emittedAskTimedoutQuestionIds.add(questionId); const workerAlive = isWaitingWorkerAlive(task, now); timedOutQuestionIds.push(questionId); pendingEvents.push({ taskId: task.id, questionId, workerAlive }); if (!workerAlive) { // DEAD: clear the park and re-queue with the injected timeout note. const note = renderAskTimedoutInjection(questionId); tasks = tasks.map((t) => t.id === task.id ? { ...t, status: "queued" as const, startedAt: undefined, finishedAt: undefined, error: undefined, waiting: undefined, pendingSteers: [...(t.pendingSteers ?? []), note], adaptive: { ...t.adaptive, phase: "resumed", task: t.adaptive?.task ?? "", }, } : t, ); requeuedTaskIds.push(task.id); } } if (requeuedTaskIds.length === 0) { return { manifest: fresh.manifest, tasks: fresh.tasks, timedOutQuestionIds, requeuedTaskIds }; } saveRunTasks(fresh.manifest, tasks); let manifest = fresh.manifest; // WP-2 review round 1 (P3): clear manifest.waitState when its park was // requeued — mirror the broker's clearWaitState guard, else the stale // pointer shields the run from staleness repair for the full 24h TTL. if (manifest.waitState && requeuedTaskIds.includes(manifest.waitState.taskId)) { manifest = { ...manifest, waitState: undefined, updatedAt: new Date().toISOString() }; saveRunManifest(manifest); } if ( manifest.status === "blocked" || manifest.status === "completed" || manifest.status === "failed" || manifest.status === "cancelled" ) { manifest = updateRunStatus(manifest, "running", `Requeued ${requeuedTaskIds.length} waiting task(s) after ask timeout.`); } try { const existingRuntimes = new Map(readCrewAgents(fresh.manifest).map((a) => [a.taskId, a.runtime])); saveCrewAgents( fresh.manifest, tasks .filter((t) => requeuedTaskIds.includes(t.id)) .map((t) => recordFromTask(fresh.manifest, t, existingRuntimes.get(t.id) ?? "child-process")), ); } catch (error) { logInternalError("dispatch-batch.ask-timeout.crewAgents", error, `runId=${runId}`); } return { manifest, tasks, timedOutQuestionIds, requeuedTaskIds }; }); if (!outcome) return undefined; // Events appended AFTER the run lock is released; awaited so callers (and // tests) observe durable events.jsonl once the sweep resolves. for (const event of pendingEvents) { await appendEventAsync(outcome.manifest.eventsPath, { type: "ask.timedout", runId, taskId: event.taskId, message: `Ask deadline expired for question ${event.questionId}${ event.workerAlive ? " (worker alive — surfaced in-tool)" : " (worker dead — task requeued)" }.`, data: { questionId: event.questionId, workerAlive: event.workerAlive, surface: event.workerAlive ? "in-tool" : "requeue" }, }).catch((error) => logInternalError("dispatch-batch.ask-timeout-event", error, `runId=${runId} taskId=${event.taskId}`)); } return outcome; } function findStep(workflow: WorkflowConfig, task: TeamTaskState): WorkflowStep { const step = workflow.steps.find((candidate) => candidate.id === task.stepId); if (!step) throw new CrewError(ErrorCode.ResourceNotFound, `Workflow step '${task.stepId}' not found for task '${task.id}'.`).withContext( `workflow step lookup (task=${task.id})`, ); // T4/R6 (ADR-6 §7): workflow-level specStrict merges into the dispatched step // (per-step flag wins) — the task packet freezes the resolved value. if (workflow.specStrict === true || step.specStrict === true) { return { ...step, specStrict: step.specStrict ?? workflow.specStrict }; } return step; } function findAgent(agents: AgentConfig[], task: TeamTaskState): AgentConfig { const agent = agents.find((candidate) => candidate.name === task.agent); if (!agent) throw new CrewError(ErrorCode.ResourceNotFound, `Agent '${task.agent}' not found for task '${task.id}'.`).withContext( `agent lookup (task=${task.id})`, ); return agent; } export function markBlocked(tasks: TeamTaskState[], reason: string): TeamTaskState[] { return tasks.map((task) => task.status === "queued" ? { ...task, status: "skipped", error: reason, finishedAt: new Date().toISOString(), graph: task.graph ? { ...task.graph, queue: "blocked" } : undefined, } : task, ); } /** * CORE-6: Unified cancel/fail of non-terminal tasks. Replaces hand-rolled * `.map()` + transform sites across this file. * * - Without `filter`: all non-terminal tasks (queued/running/waiting) are * terminalised with the given status. * - With `filter`: the filter is the sole gate — the non-terminal check is * NOT applied automatically, matching per-task/per-id variants. * - Optional `transform(task, terminalised)`: lets a caller attach * site-specific fields (graph mutation, terminalEvidence) to the * terminalised task. `terminalised` already carries status/finishedAt/error; * the transform returns it unchanged or a modified copy. * * RT-14: the two remaining inline cancel sites (cancelPlanTasks, * cancelRunFromSignal) route through this helper via `transform` so EVERY * cancel site uses the single shared transform. Their extra logic * (graph mutation / terminalEvidence) is preserved inside the transform. * * `markBlocked` is intentionally NOT unified here (it sets status "skipped", * not cancelled/failed, and only acts on "queued" tasks). */ export function cancelNonTerminalTasks( tasks: TeamTaskState[], status: "cancelled" | "failed", reason: string, filter?: (task: TeamTaskState) => boolean, transform?: (task: TeamTaskState, terminalised: TeamTaskState) => TeamTaskState, ): TeamTaskState[] { const predicate = filter ?? ((task: TeamTaskState) => isNonTerminalTaskStatus(task.status)); return tasks.map((task) => { if (!predicate(task)) return task; const terminalised: TeamTaskState = { ...task, status, finishedAt: new Date().toISOString(), error: reason }; return transform ? transform(task, terminalised) : terminalised; }); } // 2.8: adaptive-plan parsing/repair/injection moved to src/runtime/goal-workflow/adaptive-plan.ts. // Re-export the test-only helpers so existing test imports still resolve. function retryPolicyFromConfig(config: CrewReliabilityConfig | undefined): RetryPolicy { return { ...DEFAULT_RETRY_POLICY, ...(config?.retryPolicy ?? {}) }; } /** * #1 (assessment): decide whether the per-task retry path (executeWithRetry) is used. * Defaults to TRUE (opt-out) so transient worker hangs (ChildTimeout) are retried * automatically. Previously opt-in, which left the entire retry+recovery stack dormant. * Exported for unit testing. */ export function shouldUseRetry(reliability: CrewReliabilityConfig | undefined): boolean { return reliability?.autoRetry !== false; } /** * T9d cancel-race fix (2026-09-23, live battery finding 7): a retry attempt must * re-queue ONLY the task's OWN terminal failure (status "failed") while the run is * still active. Previously ANY non-queued/running status was re-queued — including * "cancelled", so a cross-session cancel landing between attempt 1's hard-kill and * attempt 2's start (measured live: cancel 16:46:22.987 → task.started 16:46:23.948, * replacement worker then ran 47s post-cancel) silently resurrected the task. * External terminal decisions (task cancelled/completed, or a terminal RUN manifest) * are respected. Exported for unit testing. */ export function shouldRequeueForRetry(input: { attempt: number; taskStatus: TeamTaskState["status"]; manifestStatus: TeamRunManifest["status"]; }): boolean { if (input.attempt <= 1) return false; if (TEAM_TERMINAL_RUN_STATUSES.has(input.manifestStatus)) return false; return input.taskStatus === "failed"; } function failedTaskFrom(result: { tasks: TeamTaskState[] }, taskId: string): TeamTaskState | undefined { return result.tasks.find((item) => item.id === taskId && item.status === "failed"); } /** * Check whether any task uses explicit `dependsOn` that would benefit from DAG-based * execution planning. If so, build an execution plan and use `getDagReadyTasks` * to augment the ready-set selection. */ function dagReadyTaskIds(tasks: TeamTaskState[], completedIds: Set): string[] | null { const hasExplicitDeps = tasks.some((t) => t.dependsOn.length > 0); if (!hasExplicitDeps) return null; // RR-012 (adjacent scheduler risk — proven reachable): exclude // delegate-broker shadow records from the DAG. They are managed by the // external grandchild spawner, not this scheduler: (1) as WAVE members a // running shadow would gate dependent waves / trigger markBlocked although // it is never in ctx.pendingUnits; (2) getDagReadyTasks below ignores both // status AND graph, so a queued OR running shadow surfaces as ready. const workflowTasks = tasks.filter((t) => !isDelegateShadowTask(t)); // FIX (goal-wrap runtime test): task.dependsOn stores STEP IDs (e.g. "execute"), not // task IDs (e.g. "02_execute"). The DAG scheduler compares deps against completedIds // (which are task IDs), so step-ID deps would never match → dependent tasks stuck blocked // forever. Map step IDs -> task IDs first (mirror dependencySatisfied in // task-graph-scheduler.ts which handles this via stepToTaskId). buildDagExecutionPlan + // getDagReadyTasks then work on consistent task IDs. const stepToTaskId = new Map(); for (const t of workflowTasks) { if (t.stepId) stepToTaskId.set(t.stepId, t.id); } const nodes: TaskNode[] = workflowTasks.map((t) => ({ id: t.id, dependsOn: t.dependsOn.map((dep) => stepToTaskId.get(dep) ?? dep), phase: t.adaptive?.phase ?? t.stepId, })); const plan = buildDagExecutionPlan(nodes); if (plan.hasCycle) return null; // fall back to existing scheduler return getDagReadyTasks(plan, completedIds); } /** * CORE-4 extraction 3: select the dispatch batch for the current loop * iteration. * * Computes the task-graph snapshot, DAG-ready tasks, workflow phase * preconditions, batch concurrency, write-path-overlap serialization, * coalesced-group logging, and streaming-dispatch slot allocation to * determine which tasks are ready to dispatch this cycle. * * Returns: * - `{ kind: "return", result }` when the run must block or abort * (plan-approval pending with mutating tasks, or no ready task at all). * - `{ kind: "dispatch", batch, ... }` when a batch is selected (may be * empty when tasks are still in-flight — the caller proceeds to the * wait phase with an empty dispatch set). * * The caller syncs `ctx.wfMachine` back after the call because this * function may advance the workflow phase state machine. * * @param ctx The scheduler context; `ctx.wfMachine`, `ctx.tasks`, and * `ctx.manifest` may be mutated in-place. */ /** * Pre-batch sweeps, shared by every scheduler tick (selectDispatchBatch) — * exported separately so the WIRING (not just the sweep internals) is testable * without driving a full team run (code review R7c). * * 1. WP-2/R2 (ADR-0 item 10): parked-ask deadline sweep — every tick, before * batch selection (a requeued park must be dispatchable this tick). * 2. T2/R4 (ADR-4 §4): re-plan reconciliation — queued tasks of items dropped * by the current revision are cancelled, in-flight ones get a wrap-up * advisory (soft cancel). Cheap-exits inside make this a no-op for runs * without plan-linked tasks. * Mutates ctx.manifest/ctx.tasks from the sweeps' disk-reloaded results. */ export async function runSchedulerSweeps(ctx: { manifest: TeamRunManifest; tasks: TeamTaskState[] }): Promise { // PERF (2026-08-24): pass this tick's task view as the expiry hint — the // sweep then skips the manifest load entirely when nothing is expired. const sweep = await sweepExpiredWaitingTasks(ctx.manifest.cwd, ctx.manifest.runId, Date.now(), ctx.tasks); if (sweep) { ctx.manifest = sweep.manifest; ctx.tasks = sweep.tasks; } const droppedSweep = sweepDroppedPlanItems(ctx.manifest, ctx.tasks); if (droppedSweep) { ctx.manifest = droppedSweep.manifest; ctx.tasks = droppedSweep.tasks; } } export async function selectDispatchBatch(ctx: SchedulerContext): Promise { // WP-2/R2 (ADR-0 item 10): this tick owns parked-ask deadlines — every // scheduler iteration sweeps tasks with an expired `task.waiting.deadline` // BEFORE batch selection (a requeued park must be dispatchable this tick). await runSchedulerSweeps(ctx); const snapshot = taskGraphSnapshot(ctx.tasks, ctx.queueIndex); // DAG-based execution plan: when tasks have explicit dependsOn, use the // topological wave planner to determine ready tasks. Fall back to the // existing task-graph-scheduler when no explicit deps exist (backward compat). const completedIds = new Set(ctx.tasks.filter((t) => t.status === "completed" || t.status === "needs_attention").map((t) => t.id)); const dagReady = dagReadyTaskIds(ctx.tasks, completedIds); // RR-012 (adjacent scheduler risk — proven reachable): belt-and-braces // selection guard. dagReadyTaskIds already excludes shadow records from the // DAG; this filter also covers the snapshot fallback (a graphed queued // shadow would resolve to queue "ready"). A shadow in a batch makes // findStep() throw ResourceNotFound (task.stepId === undefined) — the run // aborts (pre-warm site) or the record is force-failed by the synthesized // unit error. Shadow records stay visible in snapshots / team status. const shadowTaskIds = new Set(ctx.tasks.filter(isDelegateShadowTask).map((t) => t.id)); const readyBeforeFilter = (dagReady ?? snapshot.ready).filter((taskId) => !shadowTaskIds.has(taskId)); // Workflow phase precondition check (non-blocking: log warnings only). if (ctx.wfMachine.currentPhaseIndex < ctx.wfMachine.phases.length) { const completedArtifacts = ctx.manifest.artifacts.filter((a) => a.kind === "result" || a.kind === "summary").map((a) => a.path); const previousPhaseStatus = ctx.wfMachine.currentPhaseIndex > 0 ? (ctx.wfMachine.phases[ctx.wfMachine.currentPhaseIndex - 1]?.status ?? "pending") : "completed"; const wfContext: PhaseGuardContext = { completedArtifacts, previousPhaseStatus, taskResults: ctx.tasks .filter((t) => t.status === "completed" || t.status === "needs_attention") .map((t) => ({ taskId: t.id, status: t.status, outputPath: t.resultArtifact?.path, })), }; const preconditions = validatePhasePreconditions(ctx.wfMachine, wfContext); if (!preconditions.ready) { await appendEventAsync(ctx.manifest.eventsPath, { type: "workflow.preconditions", runId: ctx.manifest.runId, message: `Workflow phase '${ctx.wfMachine.phases[ctx.wfMachine.currentPhaseIndex]?.name}' is missing inputs: ${preconditions.blocking.join(", ")}`, data: { phaseIndex: ctx.wfMachine.currentPhaseIndex, phaseName: ctx.wfMachine.phases[ctx.wfMachine.currentPhaseIndex]?.name, blocking: preconditions.blocking, }, }); } else { // Advance the machine past completed phases. while ( ctx.wfMachine.currentPhaseIndex < ctx.wfMachine.phases.length && ctx.wfMachine.phases[ctx.wfMachine.currentPhaseIndex]?.status === "completed" ) { ctx.wfMachine = { ...ctx.wfMachine, currentPhaseIndex: ctx.wfMachine.currentPhaseIndex + 1, }; } } } // W5-4: by-id map once (was O(ready × tasks) via find-per-element). const taskByIdReady = new Map(ctx.tasks.map((t) => [t.id, t] as const)); const readyRoles = readyBeforeFilter.map((taskId) => taskByIdReady.get(taskId)?.role).filter((role): role is string => Boolean(role)); const concurrency = resolveBatchConcurrency({ workflowName: ctx.workflow.name, workflowMaxConcurrency: ctx.workflow.maxConcurrency, teamMaxConcurrency: ctx.input.team.maxConcurrency, limitMaxConcurrentWorkers: ctx.input.limits?.maxConcurrentWorkers, allowUnboundedConcurrency: ctx.input.limits?.allowUnboundedConcurrency, readyCount: readyBeforeFilter.length, workspaceMode: ctx.manifest.workspaceMode, readyRoles, }); // Round 25 (M5): serialize on write-path overlap when opted in. // Opt-in via limits.serializeOnPathOverlap; default off (= no behavior change). // filterReadyByWriteOverlap returns the same array when enabled=false, so // production runs pay nothing for the unused code path. When the flag is on, // `serializedReady` MAY be a strict subset of `readyBeforeFilter` (conflicting tasks // deferred to next cycle). const serializedReady = filterReadyByWriteOverlap( readyBeforeFilter, ctx.tasks, ctx.workflow, concurrency.maxConcurrent, ctx.input.limits?.serializeOnPathOverlap === true, ); // Round 25 (M6): coalesce micro-tasks when opted in. // Default off; when on, groups same-(role,cwd) tasks into coalesced groups // (with write-path safety). In v0.9.17 first ship, we ONLY log the // coalesced group count to the event stream (informational). Actual // dispatching of one-multi-task worker instead of N workers is deferred // to a follow-up — it's a non-trivial prompt-construction change that // deserves its own PR. For now, every coalesced group => one info event. const coalesceEnabled = ctx.workflow.coalesceMicroTasks === true; if (coalesceEnabled) { const coalescedGroups = planCoalescedGroups(serializedReady, ctx.tasks, ctx.workflow, true); for (const group of coalescedGroups) { if (group.tasks.length < 2) continue; // singletons are not interesting await appendEventAsync(ctx.manifest.eventsPath, { type: "task.coalesced", runId: ctx.manifest.runId, message: `Coalesced ${group.tasks.length} micro-tasks (role=${group.role}, cwd=${group.cwd})`, data: { groupId: group.id, role: group.role, cwd: group.cwd, taskIds: group.tasks.map((task) => task.id), }, }); } } if (concurrency.reason.includes(";unbounded:")) { await appendEventAsync(ctx.manifest.eventsPath, { type: "limits.unbounded", runId: ctx.manifest.runId, message: "Unbounded worker concurrency was explicitly enabled for this run.", data: { concurrencyReason: concurrency.reason, maxConcurrent: concurrency.maxConcurrent, }, }); } // ── OPT-01 streaming dispatch: exclude tasks already in-flight, limit // new dispatches to available concurrency slots. ── const inFlightTaskIds = new Set(); for (const pendingUnit of ctx.pendingUnits.values()) { for (const taskId of pendingUnit.taskIds) inFlightTaskIds.add(taskId); } const slotsAvailable = Math.max(0, concurrency.maxConcurrent - ctx.pendingUnits.size); // Review R3: the dispatch gate reads plan-record-first (ADR-4 §8) — in the // crash window between the record write and the manifest save, the record's // decision wins (resolve in favor of user intent). const approvalPending = isPlanApprovalPendingEffective(ctx.manifest); const dispatchableReady = serializedReady.filter((id) => !inFlightTaskIds.has(id)); const readyIds = approvalPending ? dispatchableReady : dispatchableReady.slice(0, slotsAvailable); const taskByIdDispatch = new Map(ctx.tasks.map((t) => [t.id, t] as const)); const candidateBatch = readyIds.map((id) => taskByIdDispatch.get(id)).filter((task): task is TeamTaskState => Boolean(task)); const readyBatch = approvalPending ? candidateBatch.filter((task) => !isMutatingTask(task)).slice(0, slotsAvailable) : candidateBatch; if (readyBatch.length === 0) { if (ctx.pendingUnits.size > 0) { // Tasks are in-flight — skip dispatch and proceed to wait phase. // (No return; code falls through to the dispatch section which is // a no-op with an empty readyBatch, then reaches the wait phase.) } else if (approvalPending && candidateBatch.some(isMutatingTask)) { await saveRunTasksAsync(ctx.manifest, ctx.tasks); saveCrewAgents(ctx.manifest, recordsForMaterializedTasks(ctx.manifest, ctx.tasks, ctx.runtimeKind)); ctx.manifest = updateRunStatus(ctx.manifest, "blocked", "Plan approval required before mutating implementation tasks run."); return { kind: "return", result: { manifest: ctx.manifest, tasks: ctx.tasks } }; } else { ctx.tasks = markBlocked(ctx.tasks, "No ready queued task; dependency graph may be invalid."); await saveRunTasksAsync(ctx.manifest, ctx.tasks); saveCrewAgents(ctx.manifest, recordsForMaterializedTasks(ctx.manifest, ctx.tasks, ctx.runtimeKind)); ctx.manifest = updateRunStatus(ctx.manifest, "blocked", "No ready queued task."); return { kind: "return", result: { manifest: ctx.manifest, tasks: ctx.tasks } }; } } return { kind: "dispatch", batch: readyBatch, concurrency, snapshot, approvalPending, coalesceEnabled }; } export type DispatchBatchDecision = Extract; /** * CORE-4 extraction 4: execute the dispatch batch selected by * selectDispatchBatch. * * Runs before_task_start hooks (skipping blocked tasks), builds coalesced * dispatch units, pre-warms the stable-prefix cache for unique cwds, and * dispatches each unit into ctx.pendingUnits as a fire-and-forget promise * (wrapped in executeWithRetry on the singleton path). The function is a * verbatim lift of the inline dispatch block; it does not return a * SchedulerDecision (void — it only populates ctx.pendingUnits). * * Reads ctx.manifest/tasks/workflow/input + runController.signal. Mutates * ctx.pendingUnits (add), ctx.tasks (hook skips), ctx.manifest (hook * status). The mutable manifest/tasks are accessed via ctx.* (not captured * locals) so that async retry callbacks observe the caller's re-synced * values, matching the original closure semantics. * * @param ctx The scheduler context. * @param decision The dispatch decision from selectDispatchBatch. */ export async function dispatchBatch(ctx: SchedulerContext, decision: DispatchBatchDecision): Promise { const { batch: readyBatch, concurrency, snapshot, approvalPending, coalesceEnabled } = decision; // Immutable context fields captured once; manifest/tasks are accessed via // ctx.* because they may be re-synced by the caller between dispatch and // promise resolution (retry callbacks fire asynchronously). const { workflow, input, runtimeKind, runController } = ctx; // 2.2 caller migration: batch progress is high-frequency informational (M7 wire). // .catch REQUIRED — buffered-flush rejections reject queued promises (see // child-executor.ts:530 note, CI 2026-10-03). void appendEventBuffered(ctx.manifest.eventsPath, { type: "task.progress", runId: ctx.manifest.runId, message: `Starting ready batch with ${readyBatch.length} task(s).`, data: { taskIds: readyBatch.map((task) => task.id), readyCount: snapshot.ready.length, blockedCount: snapshot.blocked.length, runningCount: snapshot.running.length, doneCount: snapshot.done.length, selectedCount: readyBatch.length, maxConcurrent: concurrency.maxConcurrent, defaultConcurrency: concurrency.defaultConcurrency, concurrencyReason: approvalPending ? `${concurrency.reason};plan-approval-read-only` : concurrency.reason, }, }).catch((error) => logInternalError("dispatch-batch.ready-batch-progress", error, `runId=${ctx.manifest.runId}`)); // Execute before_task_start hooks for the batch — P1-10: run hooks in // parallel (each may be a subprocess), then apply skip mutations in order. const beforeTaskStartReports = await Promise.all( readyBatch.map((task) => executeHook("before_task_start", { runId: ctx.manifest.runId, taskId: task.id, cwd: ctx.manifest.cwd, }).then((taskReport) => ({ task, taskReport })), ), ); for (const { task, taskReport } of beforeTaskStartReports) { appendHookEvent(ctx.manifest, taskReport); if (taskReport.outcome === "block") { ctx.tasks = ctx.tasks.map((t) => t.id === task.id ? { ...t, status: "skipped" as const, error: taskReport.reason ?? "before_task_start hook blocked execution.", } : t, ); ctx.manifest = updateRunStatus(ctx.manifest, ctx.manifest.status, `Task '${task.id}' blocked by hook.`); } } // W5-4: by-id map (was O(readyBatch × tasks) via find-per-element). const ctxTaskById = new Map(ctx.tasks.map((t) => [t.id, t] as const)); const batchTasks = readyBatch.filter((task) => { const t = ctxTaskById.get(task.id); return t !== undefined && t.status !== "skipped"; }); if (batchTasks.length > 1) { await appendEventAsync(ctx.manifest.eventsPath, { type: "task.parallel_start", runId: ctx.manifest.runId, message: `Launching ${batchTasks.length} tasks in PARALLEL (concurrency=${concurrency.selectedCount}): ${batchTasks.map((t) => `${t.role}(${t.id})`).join(", ")}`, data: { taskIds: batchTasks.map((t) => t.id), roles: batchTasks.map((t) => t.role), concurrency: concurrency.selectedCount, }, }); } // M6 real dispatch: when coalesceMicroTasks is enabled, batch the // ready tasks into dispatch units. Multi-task groups are dispatched // as one worker (single cold-start) instead of N. Singletons fall // through to per-task dispatch. const coalescedGroups = planCoalescedGroups( batchTasks.map((t) => t.id), ctx.tasks, workflow, coalesceEnabled, ); const dispatchUnits = buildDispatchUnits( batchTasks.map((t) => t.id), coalescedGroups, ); // NEW-M1: Pre-warm stable prefix cache for one representative task // per unique cwd. Parallel siblings with the same cwd/step reuse // the cached workspace tree, file retrieval, and knowledge fragment // instead of recomputing them independently (~200-800ms per batch). if (batchTasks.length > 1) { const seenCwds = new Set(); await Promise.all( batchTasks .filter((task) => { if (seenCwds.has(task.cwd)) return false; seenCwds.add(task.cwd); return true; }) .map((task) => { const step = findStep(workflow, task); return computeStablePrefixComponents(ctx.manifest, step, task); }), ); } // ── OPT-01 streaming dispatch: dispatch each unit into ctx.pendingUnits // instead of awaiting the entire batch via mapConcurrent. Each unit's // promise is stored so we can Promise.race on the next iteration. ── const dispatchUnit = async (unit: DispatchUnit): Promise<{ manifest: TeamRunManifest; tasks: TeamTaskState[] }> => { // M6 real dispatch path: single worker for N tasks. if (unit.kind === "group") { const groupTasks = unit.group.tasks; const firstTask = groupTasks[0]!; const step = findStep(workflow, firstTask); const agent = findAgent(input.agents, firstTask); const teamRole = input.team.roles.find((role) => role.name === firstTask.role); const perTaskRuntime = resolveTaskRuntimeKind(runtimeKind, firstTask.role, input.runtimeConfig?.isolationPolicy); return runCoalescedTaskGroup({ manifest: ctx.manifest, tasks: ctx.tasks, groupTasks, step, agent, signal: runController.signal, executeWorkers: input.executeWorkers, runtimeKind, workspaceId: input.workspaceId, onJsonEvent: input.onJsonEvent, runtimeConfig: input.runtimeConfig, reliability: input.reliability, teamRole, perTaskRuntime, }); } // Singleton path: original per-task dispatch. const task = batchTasks.find((t) => t.id === unit.taskId)!; const step = findStep(workflow, task); const agent = findAgent(input.agents, task); const teamRole = input.team.roles.find((role) => role.name === task.role); const perTaskRuntime = resolveTaskRuntimeKind(runtimeKind, task.role, input.runtimeConfig?.isolationPolicy); // CORE-3: compute retry policy + spawn budget ONCE per dispatch unit. // The spawnBudget object is shared (by reference) across every // runTeamTask call within executeWithRetry via baseInput spread, // so the counter accumulates across retry attempts × model fallbacks. const policy = retryPolicyFromConfig(input.reliability); const spawnBudget: SpawnBudget = { count: 0, max: policy.maxTotalSpawns ?? 0 }; const baseInput = { manifest: ctx.manifest, tasks: ctx.tasks, task, step, agent, signal: runController.signal, executeWorkers: input.executeWorkers, runtimeKind: runtimeKind, taskRuntimeOverride: perTaskRuntime !== runtimeKind ? perTaskRuntime : undefined, runtimeConfig: input.runtimeConfig, parentContext: input.parentContext, parentModel: input.parentModel, modelRegistry: input.modelRegistry, // P2-1: thread the host metric registry to the task-runner seam so the // retry-triage classifier can count calls/decisions (same // ExecuteTeamRunInput handle dispatch-batch itself uses at :812). metricRegistry: input.metricRegistry, modelOverride: input.modelOverride, teamRoleModel: teamRole?.model, teamRoleThinking: teamRole?.thinking, teamRoleFallbackModels: teamRole?.fallbackModels, teamRoleSkills: teamRole?.skills, skillOverride: input.skillOverride, limits: input.limits, onJsonEvent: input.onJsonEvent, workspaceId: input.workspaceId, spawnBudget, // R10-1 residual: same per-run cache instance as the closeout — one // object reference spread into EVERY runTeamTask call (initial + every // retry attempt via `...baseInput`), so dep-context reads and closeout // aggregation share memoized artifacts. resultReadCache: ctx.resultReadCache, }; // #1 (assessment): autoRetry now defaults ON (opt-out via reliability.autoRetry=false). // The dominant v0.9.13 failure was ChildTimeout ("worker became unresponsive") with // ZERO retries because this gate was opt-in. isRetryable() defaults to true when // retryableErrors is empty, so transient hangs now retry up to maxAttempts (3) with // exponential backoff. Set reliability.autoRetry=false to restore old single-shot behavior. if (!shouldUseRetry(input.reliability)) return withCorrelation(childCorrelation(ctx.manifest.runId, task.id), () => runTeamTask(baseInput)); let lastFailed: { manifest: TeamRunManifest; tasks: TeamTaskState[] } | undefined; let lastAttemptId: string | undefined; const attemptsSoFar: TaskAttemptState[] = [...(task.attempts ?? [])]; try { return await executeWithRetry( async (attempt, info) => { const startedAt = new Date().toISOString(); const inFlightAttempts: TaskAttemptState[] = [...attemptsSoFar, { attemptId: info.attemptId, startedAt }]; input.metricRegistry?.counter("crew.task.retry_attempt_total", "Retry attempts by run and task").inc({ runId: ctx.manifest.runId, taskId: task.id, }); // NOTE: no withRunLock — best-effort only; concurrent writes may cause inconsistency const fresh = loadRunManifestById(ctx.manifest.cwd, ctx.manifest.runId); const freshManifest = fresh?.manifest ?? ctx.manifest; let freshTasks = fresh?.tasks ?? ctx.tasks; let freshTask = freshTasks.find((item) => item.id === task.id) ?? task; // US-003 (2026-09-22) — RETRY-NOOP BUG, found while wiring the dead-letter // integration test: a retry (attempt > 1) only happens because OUR previous // attempt THREW — i.e. we ourselves persisted the terminal failure (e.g. // model-exhausted marks the task failed before the throw). The terminal-state // early-return below then treated OUR OWN failure as "nothing to do" and // returned it as SUCCESS: executeWithRetry never re-ran the task, never // exhausted, and the onRetryGivenUp dead-letter never fired — autoRetry was // silently a single attempt for every failure that persists task state. // Fix: on a retry, re-queue OUR OWN terminal failure so the retry actually // re-runs. Attempt 1 keeps the original guard (externally-terminal task = // someone else's decision; external cancellation exits via the signal). // 2026-09-23 (finding 7): re-queue ONLY own failures on an ACTIVE run — // cancelled/completed tasks and terminal manifests must never resurrect. if (shouldRequeueForRetry({ attempt, taskStatus: freshTask.status, manifestStatus: freshManifest.status })) { freshTask = { ...freshTask, status: "queued", error: undefined, finishedAt: undefined }; freshTasks = freshTasks.map((item) => (item.id === task.id ? freshTask : item)); } if (freshTask.status !== "queued" && freshTask.status !== "running") return { manifest: freshManifest, tasks: freshTasks, }; const taskWithAttempt: TeamTaskState = { ...freshTask, attempts: inFlightAttempts, }; const result = await withCorrelation(childCorrelation(freshManifest.runId, task.id), () => runTeamTask({ ...baseInput, manifest: freshManifest, tasks: freshTasks, task: taskWithAttempt, }), ); const failed = failedTaskFrom(result, task.id); const endedAt = new Date().toISOString(); const finishedAttempt: TaskAttemptState = { attemptId: info.attemptId, startedAt, endedAt, ...(failed?.error ? { error: failed.error } : {}), }; attemptsSoFar.push(finishedAttempt); const withAttempt = result.tasks.map((item) => item.id === task.id ? { ...item, attempts: [...attemptsSoFar] } : item, ); const enriched = { manifest: result.manifest, tasks: withAttempt, }; if (failed) { lastFailed = enriched; throw new CrewError(ErrorCode.TaskNotFound, failed.error ?? `Task ${task.id} failed.`).withContext( `retry evaluation (run=${ctx.manifest.runId})`, ); } input.metricRegistry?.histogram("crew.task.retry_count", "Retries per task", [0, 1, 2, 3, 5, 10]).observe( { runId: ctx.manifest.runId, team: input.team.name, }, Math.max(0, attempt - 1), ); return enriched; }, policy, { signal: runController.signal, attemptId: (attempt) => `${ctx.manifest.runId}:${task.id}:attempt-${attempt}`, onAttemptFailed: (attempt, error, delayMs, info) => { lastAttemptId = info.attemptId; appendEventAsync(ctx.manifest.eventsPath, { type: "crew.task.retry_attempt", runId: ctx.manifest.runId, taskId: task.id, message: error.message, data: { attempt, attemptId: info.attemptId, delayMs, }, metadata: { attemptId: info.attemptId }, }).catch((error) => logInternalError("team-runner.retry-attempt", error, `taskId=${task.id}`)); input.metricRegistry?.histogram("crew.task.retry_delay_ms", "Retry backoff delay, milliseconds").observe( { runId: ctx.manifest.runId, taskId: task.id, }, delayMs, ); }, onRetryGivenUp: (attempts, error, info) => { lastAttemptId = info.attemptId; // US-003: aborts are NOT exhaustion — executeWithRetry also fires this // hook on its cancelled exit, which previously deadlettered cancelled // tasks as "max-retries" (false positive; keep the signal high). if (runController.signal.aborted) return; // US-003: richer entry (agent/role/modelAttempts/runStatus) so the // project-level index is self-sufficient after the run dir is pruned. // Prefer the FAILED attempt's final task state over the pre-run snapshot. const failedTask = lastFailed?.tasks.find((item) => item.id === task.id); const deadletterTask = failedTask ?? task; appendDeadletter(ctx.manifest, { runId: ctx.manifest.runId, taskId: task.id, reason: "max-retries", attempts, attemptId: info.attemptId, lastError: error.message, timestamp: new Date().toISOString(), agent: deadletterTask.agent, role: deadletterTask.role, modelAttempts: deadletterTask.modelAttempts?.length, runStatus: ctx.manifest.status, }); input.metricRegistry ?.counter("crew.task.deadletter_total", "Deadletter triggers by reason") .inc({ reason: "max-retries" }); input.metricRegistry?.histogram("crew.task.retry_count", "Retries per task", [0, 1, 2, 3, 5, 10]).observe( { runId: ctx.manifest.runId, team: input.team.name, }, Math.max(0, attempts - 1), ); }, }, ); } catch (retryError) { if (retryError instanceof CrewCancellationError || input.signal?.aborted) { const reason = retryError instanceof CrewCancellationError ? retryError.reason : cancellationReasonFromSignal(input.signal); // NOTE: no withRunLock — best-effort only; concurrent writes may cause inconsistency const fresh = loadRunManifestById(ctx.manifest.cwd, ctx.manifest.runId); const freshManifest = fresh?.manifest ?? ctx.manifest; const freshTasks = fresh?.tasks ?? ctx.tasks; const cancelledTasks = cancelNonTerminalTasks( freshTasks, "cancelled", `${reason.message} (${reason.code})`, (item) => item.id === task.id && (item.status === "queued" || item.status === "running"), ); appendEventAsync(freshManifest.eventsPath, { type: "task.cancelled", runId: freshManifest.runId, taskId: task.id, message: reason.message, data: { reason, phase: "retry" }, metadata: lastAttemptId ? { attemptId: lastAttemptId } : undefined, }).catch((error) => logInternalError("team-runner.cancelled", error, `taskId=${task.id}`)); return { manifest: updateRunStatus(freshManifest, "cancelled", reason.message), tasks: cancelledTasks, }; } if (lastFailed) return lastFailed; // NOTE: no withRunLock — best-effort only; concurrent writes may cause inconsistency const fresh = loadRunManifestById(ctx.manifest.cwd, ctx.manifest.runId); const freshManifest = fresh?.manifest ?? ctx.manifest; const freshTasks = fresh?.tasks ?? ctx.tasks; const freshTask = freshTasks.find((item) => item.id === task.id) ?? task; if (freshTask.status !== "queued" && freshTask.status !== "running") return { manifest: freshManifest, tasks: freshTasks }; return withCorrelation(childCorrelation(freshManifest.runId, task.id), () => runTeamTask({ ...baseInput, manifest: freshManifest, tasks: freshTasks, task: freshTask, }), ); } }; // ── OPT-01 streaming dispatch: dispatch units into ctx.pendingUnits ── for (const unit of dispatchUnits) { const unitKey = unit.kind === "singleton" ? unit.taskId : unit.group.id; const unitTaskIds = unit.kind === "singleton" ? [unit.taskId] : unit.group.tasks.map((t) => t.id); // RT-12: create the wrapper promise ONCE at dispatch time so // mergeUnitResult can Promise.race on pre-existing wrappers instead // of allocating new async closures every loop iteration. const rawPromise = dispatchUnit(unit); const wrapped: Promise = (async () => { try { const result = await rawPromise; return { unitKey, result: result as { manifest: TeamRunManifest; tasks: TeamTaskState[] } | undefined, error: undefined as Error | undefined, }; } catch (error) { return { unitKey, result: undefined, error: error instanceof Error ? error : new Error(String(error)) }; } })(); ctx.pendingUnits.set(unitKey, { taskIds: unitTaskIds, promise: rawPromise, wrapped, }); // RT-NEW-2 race fix: record ever-dispatched task ids so terminaliseRunWithDrain // cancels (not skips) tasks whose unit settled + left pendingUnits before // the abort fired but whose task status isn't terminal yet. for (const id of unitTaskIds) ctx.dispatchedTaskIds.add(id); } } /** @internal 1.9(b) test export — exercise selectDispatchBatch directly. */ export const __test__selectDispatchBatch = selectDispatchBatch;