diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index 956cb2bec4..fef1342985 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -1,8 +1,5 @@ import { getCurrentRepo, - resolveDependencyOrder, - sortTasksByPriorityFanoutThenAgeAndId, - buildUnblockWeightMap, computeBlockerFanoutMap, compareTasksByPriorityThenAgeAndId, HIGH_FANOUT_BLOCKER_TODO_THRESHOLD, @@ -14,8 +11,6 @@ import { type PrInfo, type AgentStore, type Settings, - TransitionRejectionError, - isUnplannedSeedPrompt, } from "@fusion/core"; import { existsSync } from "node:fs"; import { readFile } from "node:fs/promises"; @@ -44,7 +39,7 @@ import { StaleTaskReporter } from "./stale-task-reporter.js"; import { BacklogPressureReporter } from "./backlog-pressure-reporter.js"; import { UnlinkedMissionsAdvisoryReporter } from "./unlinked-missions-advisory-reporter.js"; import { createRunAuditor, generateSyntheticRunId } from "./run-audit.js"; -import { DEFAULT_WORKFLOW_POOL_ID, resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags } from "@fusion/core"; +import { resolveWorkflowIrForTask, resolveWorkflowIrById, resolveColumnFlags } from "@fusion/core"; import type { WorkflowIr, WorkflowIrV2 } from "@fusion/core"; import { runHoldReleaseSweep, isUnplannedForExecution, type SlotReservation } from "./hold-release.js"; import { moveTaskToReplanColumn } from "./replan-target.js"; @@ -1497,930 +1492,6 @@ export class Scheduler { } return; } - - const maxConcurrent = settings.maxConcurrent ?? this.options.maxConcurrent ?? 2; - const maxWorktrees = settings.maxWorktrees ?? this.options.maxWorktrees ?? 4; - - /* - FNXC:WorkflowReviewGates 2026-07-26-13:10: - DEAD CODE — deliberately NOT carrying the review-gate WIP-occupancy fix here. - `shouldRunWorkflowColumnScheduler()` returns an unconditional `true` and the branch above it - always returns, so this legacy dispatcher is unreachable. The live capacity accounting is in - `runHoldReleaseSweepPass` (`reservedWorktreeSlots`/`reservedConcurrentSlots`). Mirroring the - fix into this block would only imply coverage that never executes. - */ - // Count only in-progress tasks toward the worktree limit. - // In-review tasks with worktrees are idle (waiting to merge) and - // should not block new tasks from starting. - const activeWorktrees = tasks.filter( - (t) => t.column === "in-progress", - ).length; - - if (activeWorktrees >= maxWorktrees) { - if (!this.wasWorktreeLimited) { - schedulerLog.log(`Worktree limit reached (${activeWorktrees}/${maxWorktrees})`); - this.wasWorktreeLimited = true; - } - return; - } - - this.wasWorktreeLimited = false; - - const inProgress = tasks.filter((t) => t.column === "in-progress"); - - // Execution tasks occupy concurrency slots governed by maxConcurrent. - // Triage/specification tasks have their own limit (maxTriageConcurrent) - // and do not count against this slot. - const agentSlots = inProgress.length; - - // When a semaphore is provided, factor in its available slots so we - // don't schedule more tasks than the global limit allows. - const inProgressTaskIds = inProgress.map((task) => task.id); - /* - FNXC:ConcurrencyAdmission 2026-08-03-13:00: - FN-8453 capacity decisions must enrich task rows from their workflow IR. - A custom complete column can retain stale session metadata, so raw task - counting here would falsely exhaust executable capacity. - */ - const topLevelClaimedSlots = await computeTopLevelConcurrencyClaimedFromStore({ - store: this.store, - tasks, - }); - const computeDispatchCapacityDiagnostic = (startedThisTick: number): ConcurrencyGateDiagnostic => { - const started = Math.max(0, Math.floor(startedThisTick)); - // U6 (KTD-10): report the default workflow's in-progress capacity as a - // per-column gate — the generalization of the legacy maxConcurrent gate - // (which reads through to the same value). - // FNXC:WorkflowColumns 2026-07-27-09:42 (U2 / R9): the - // `isWorkflowColumnsEnabled` conditional is deleted — it returned a - // literal `true`, so the `undefined` arm never produced a diagnostic. - const perColumnGates = [{ - workflowId: DEFAULT_WORKFLOW_POOL_ID, - columnId: "in-progress", - used: agentSlots + started, - limit: maxConcurrent, - slack: maxConcurrent - (agentSlots + started), - }]; - return computeConcurrencyGateDiagnostic({ - agentSlots, - maxConcurrent, - activeWorktrees, - maxWorktrees, - semaphore: this.options.semaphore, - inProgressTaskIds, - topLevelClaimedSlots, - startedThisTick: started, - perColumnGates, - }); - }; - if (computeDispatchCapacityDiagnostic(0).available <= 0) return; - - const now = Date.now(); - let todo = tasks.filter((t) => { - if (t.column !== "todo" || t.paused || t.userPaused) return false; - // Skip tasks with a recovery backoff that hasn't elapsed yet - if (t.nextRecoveryAt && new Date(t.nextRecoveryAt).getTime() > now) return false; - // FNXC:CodingIdeasWorkflow 2026-07-04-10:45: a todo task with status "planning" is being specified in place by the triage service (merged planner/capacity column in Coding (Ideas)); it must not be dispatched until planning finishes and the status clears. - if (t.status === "planning") return false; - return true; - }); - /* - FNXC:CodingIdeasWorkflow 2026-07-04-10:46: - Exclude unplanned todo tasks whose PROMPT.md is still the bootstrap stub. In a merged planner/capacity column a freshly promoted card has no real spec yet; dispatching it would execute the stub. Normal-workflow todo tasks always carry a real spec (triage writes it before moving them to todo), so this filter is a no-op for them. This closes the gap between the operator promoting a card and the triage service picking it up. - - FNXC:CodingIdeasWorkflow 2026-07-25-11:20: - Use the shared isUnplannedSeedPrompt predicate instead of an open-coded strict stub compare. - The two disagreed on the refineTask seed shape: triage's todo-discovery treated a refinement - seed as unplanned (planning it) while this filter treated it as a real spec (keeping it as a - dispatch candidate), so the two lanes could race for the same card and only hold-release stood - between an executor and a prompt containing nothing but the operator's feedback text. One - predicate, one answer — and it also absorbs CRLF/trailing-newline drift. - */ - todo = ( - await Promise.all( - todo.map(async (t) => { - try { - const content = await readFile(getPromptPath(this.store.getTasksDir(), t.id), "utf-8"); - if (isUnplannedSeedPrompt(content, t.id, t.title, t.description)) return null; - } catch { - // Missing prompt is handled by filesystem validation below; keep the candidate. - } - return t; - }), - ) - ).filter((t): t is Task => t !== null); - - - // Filter out tasks belonging to blocked missions - if (todo.length > 0 && this.options.missionStore) { - const blockedSliceIds = new Set(); - for (const t of todo) { - if (t.sliceId && !blockedSliceIds.has(t.sliceId)) { - try { - const slice = await this.options.missionStore.getSlice(t.sliceId); - if (slice) { - const milestone = await this.options.missionStore.getMilestone(slice.milestoneId); - if (milestone) { - const mission = await this.options.missionStore.getMission(milestone.missionId); - if (mission && mission.status === "blocked") { - blockedSliceIds.add(t.sliceId); - } - } - } - } catch (err: unknown) { - const errorMessage = err instanceof Error ? err.message : String(err); - schedulerLog.warn( - `Mission/slice lookup failed during scheduling (task ${t.id}): ${errorMessage} — proceeding without blocked-slice check`, - ); - // If lookup fails, don't block the task - } - } - } - if (blockedSliceIds.size > 0) { - todo = todo.filter((t) => !t.sliceId || !blockedSliceIds.has(t.sliceId)); - } - } - - if (todo.length === 0) return; - - const maxAutoMergeRetries = - typeof settings.maxAutoMergeRetries === "number" ? settings.maxAutoMergeRetries : undefined; - const unblockWeights = buildUnblockWeightMap(tasks, { - maxAutoMergeRetries, - }); - todo = sortTasksByPriorityFanoutThenAgeAndId(todo, unblockWeights); - /* - FNXC:ConcurrencyAdmission 2026-08-03-14:00: - FN-8453 makes the coordinator, rather than this lane's priority/fanout - ordering, the authority for a free top-level slot. The scheduler exposes - its ready tasks to the shared project registry and dispatches only the - atomically admitted winner; other lanes refresh in the same admission pass. - */ - const projectId = this.store.getRootDir(); - this.coordinatorReadyTasks.clear(); - for (const task of todo) this.coordinatorReadyTasks.set(task.id, task); - await projectAdmissionCoordinator.admitOldest({ - projectId, - maxConcurrent, - claimed: async () => computeTopLevelConcurrencyClaimedFromStore({ store: this.store, tasks: await this.store.listTasks({ slim: true, includeArchived: false }) }), - semaphore: this.options.semaphore, - }); - todo = todo.filter((task) => this.coordinatorAdmittedTaskIds.has(task.id)); - // FNXC:ConcurrencyAdmission 2026-08-04-10:00: coordinator IDs select one - // handoff only. The pre-held semaphore slot remains the durable reservation; - // retaining the ID after a retry would bypass oldest-first re-admission. - for (const task of todo) this.coordinatorAdmittedTaskIds.delete(task.id); - if (todo.length === 0) return; - const topWeightedTask = todo.find((candidate) => (unblockWeights.get(candidate.id) ?? 0) >= 1); - if (topWeightedTask) { - schedulerLog.log( - `Dispatch ordering: priority+fanout (top: ${topWeightedTask.id}=${unblockWeights.get(topWeightedTask.id) ?? 0})`, - ); - } - - const mergeShadowEnabled = settings.mergeRequestContractShadowEnabled === true; - const markerAcceptedByTaskId = new Map(); - if (mergeShadowEnabled) { - const dependencyIds = new Set(tasks.flatMap((candidate) => candidate.dependencies)); - for (const depId of dependencyIds) { - markerAcceptedByTaskId.set(depId, (await this.store.getCompletionHandoffAcceptedMarker(depId)) !== null); - } - } - const schedulingDependencyOptions = mergeShadowEnabled - ? { - markerAcceptedByTaskId, - onParityDiff: (diff: SchedulingDependencyParityDiff) => { - this.emitDependencyParityDiff(diff); - }, - } - : undefined; - - /** - * Pre-compute file scopes for all currently active tasks (in-progress - * AND in-review with unmerged worktrees) so that todo tasks are never - * started when their files overlap with work already underway or - * awaiting merge. - * - * Including in-review tasks prevents a blocked task from starting on - * main HEAD when the blocker's changes haven't been merged yet. - * - * The re-entrance guard on this method ensures that this snapshot - * stays consistent throughout the pass — without it, a concurrent - * pass could read stale state and start conflicting tasks. - * - * Newly started tasks are appended to this map further below so that - * subsequent todo tasks in the same pass also see them. - */ - const activeScopes = new Map(); - const activeScopeColumns = new Map(); - const setActiveScopeLease = (taskId: string, scope: string[], column: Task["column"]): void => { - activeScopes.set(taskId, scope); - activeScopeColumns.set(taskId, column); - }; - const queuedHigherPriorityScopes: QueuedOverlapCandidate[] = []; - const queuedHigherPriorityTaskById = new Map(); - const overlapIgnorePaths = settings.overlapIgnorePaths ?? []; - const filteredScopeByTaskId = new Map(); - const getFilteredFileScope = async (taskId: string): Promise => { - const cached = filteredScopeByTaskId.get(taskId); - if (cached !== undefined) return cached; - const scope = await this.store.parseFileScopeFromPrompt(taskId); - const filteredScope = filterPathsByIgnoreList(scope, overlapIgnorePaths, { ignoreHiddenOverlapPaths: settings.ignoreHiddenOverlapPaths }); - filteredScopeByTaskId.set(taskId, filteredScope); - return filteredScope; - }; - if (settings.groupOverlappingFiles) { - // In-progress tasks - for (const t of inProgress) { - if (!shouldHoldActiveFileScopeLease(t, tasks, { schedulingDependencyOptions })) continue; - const filteredScope = await getFilteredFileScope(t.id); - if (isCoordinationOnlyTask(t, filteredScope)) continue; - if (filteredScope.length === 0) continue; - setActiveScopeLease(t.id, filteredScope, "in-progress"); - } - // Only live in-review tasks with a worktree belong in activeScopes. - // Paused in-review tasks (e.g., failed-merge tasks awaiting human triage) cannot - // make progress, so they must not contribute to overlap blockers; including them - // caused a deadlock pattern where a paused task indefinitely re-stamped - // `blockedBy` on overlapping todo tasks every scheduler tick. (FN-3867 / FN-3857) - // Permanently-failed in-review tasks from SelfHealingManager.checkStuckBudget() - // also keep their worktree, but after the stuck-kill budget is exhausted they - // will never merge, so superseding re-implementation tasks (for example FN-4177 - // replaced by FN-4198) must not stay queued behind them. (FN-4200) - // FNXC:PostgresCutover 2026-06-27-09:30: - // Pre-compute handoff markers before the .filter() because - // getCompletionHandoffAcceptedMarker is async and cannot be awaited - // inside a synchronous filter callback. Without this, the Promise - // object is always !== null, making handoffAccepted incorrectly true. - const handoffMarkerMap = new Map(); - if (settings.mergeRequestContractShadowEnabled === true) { - for (const t of tasks) { - if (t.column === "in-review") { - handoffMarkerMap.set(t.id, (await this.store.getCompletionHandoffAcceptedMarker(t.id)) !== null); - } - } - } - const inReviewWithWorktree = tasks.filter( - (t) => t.column === "in-review" && shouldHoldActiveFileScopeLease(t, tasks, { - mergeRequestContractShadowEnabled: settings.mergeRequestContractShadowEnabled, - handoffAccepted: handoffMarkerMap.get(t.id) ?? false, - schedulingDependencyOptions, - }), - ); - for (const t of inReviewWithWorktree) { - const filteredScope = await getFilteredFileScope(t.id); - if (isCoordinationOnlyTask(t, filteredScope)) continue; - if (filteredScope.length === 0) continue; - - const handoffAccepted = settings.mergeRequestContractShadowEnabled === true - ? (await this.store.getCompletionHandoffAcceptedMarker(t.id)) !== null - : false; - if (!handoffAccepted) { - setActiveScopeLease(t.id, filteredScope, "in-review"); - } - - if (settings.mergeRequestContractShadowEnabled === true) { - const mergeRequestRecord = await this.store.getMergeRequestRecordAsync(t.id); - const { shadowExecutorLeaseApplied, shadowMergeLockApplied, shadowLeaseApplied } = - computeShadowLeaseParityState(mergeRequestRecord?.state ?? null); - if (shadowLeaseApplied !== !handoffAccepted) { - void this.store.recordRunAuditEvent?.({ - taskId: t.id, - agentId: "scheduler", - runId: generateSyntheticRunId("scheduler", t.id), - domain: "database", - mutationType: "merge:lease-parity-diff", - target: t.id, - metadata: { - taskId: t.id, - legacyLeaseColumn: "in-review", - legacyLeaseApplied: !handoffAccepted, - shadowLeaseApplied, - shadowExecutorLeaseApplied, - shadowMergeLockApplied, - mergeRequestState: mergeRequestRecord?.state ?? null, - }, - }); - } - } - } - - for (const t of todo) { - const filteredScope = await getFilteredFileScope(t.id); - if (isCoordinationOnlyTask(t, filteredScope)) continue; - if (filteredScope.length === 0) continue; - if (!isRunnableQueuedOverlapCandidate(t, tasks, now, activeScopes, filteredScope)) continue; - queuedHigherPriorityScopes.push({ - id: t.id, - priority: t.priority, - createdAt: t.createdAt, - scope: filteredScope, - }); - queuedHigherPriorityTaskById.set(t.id, t); - } - } - - // Resolve dependency order among todo tasks - const ordered = resolveDependencyOrder(todo); - let started = 0; - let loggedMissingAgentStoreThisPass = false; - - for (const taskId of ordered) { - const task = tasks.find((t) => t.id === taskId)!; - - if (task.checkedOutBy && this.options.leaseManager) { - const recovered = await this.options.leaseManager.recoverAbandonedLease( - task.id, - "scheduler detected stale todo lease", - { preserveProgress: true }, - ); - if (!recovered) { - await this.options.leaseManager.reconcileLeaseRow(task.id); - await this.store.updateTask(task.id, { status: "queued" }); - await this.logDispatchQueuedReason(task.id, "queued — checkout lease recovery blocked dispatch"); - continue; - } - } - - // Check all deps are satisfied (done, in-review, or archived) - const unmetDeps = getUnmetSchedulingDependencies(task, tasks, schedulingDependencyOptions); - - if (unmetDeps.length > 0) { - await this.store.updateTask(task.id, { - status: "queued", - blockedBy: unmetDeps[0], - }); - await this.logDispatchQueuedReason(task.id, `queued — unmet dependencies: ${unmetDeps.join(", ")}`); - this.options.onBlocked?.(task, unmetDeps); - continue; - } - - if (task.userPaused === true) { - if (task.status !== "queued") { - await this.store.updateTask(task.id, { status: "queued" }); - } - await this.logDispatchQueuedReason(task.id, "queued — user paused (manual move to todo)"); - continue; - } - - // Validate filesystem state before starting (only for tasks with satisfied deps) - const validation = await this.validateTaskFilesystem(task.id); - if (!validation.valid) { - schedulerLog.warn(`Task ${task.id} filesystem validation failed: ${validation.reason}`); - /* - FNXC:WorkflowScheduling 2026-07-13-11:25: - The filesystem-validation rebound must set `needs-replan`, not just move. For a - plan-in-place workflow the replan column IS "todo", so the move is a no-op — without - the status write triage cannot rediscover the card (its PROMPT.md is missing or - unreadable, so the seed check throws and skips it) and this branch re-fires every - scheduler tick forever, appending a misleading log line each time. - */ - const replanColumn = await moveTaskToReplanColumn(this.store, task); - await this.store.updateTask(task.id, { status: "needs-replan" }); - await this.store.logEntry(task.id, `Task rebounded to ${replanColumn} for re-specification — filesystem validation failed`, validation.reason); - continue; - } - - // Stale spec enforcement: check if PROMPT.md has aged beyond the configured threshold. - // When enabled, stale tasks are rebounded to the workflow-aware replan column with - // status "needs-replan" so they receive fresh specification before execution - // (workflows without a "triage" column replan in place in todo). This guard runs - // after filesystem validation so missing/unreadable files skip staleness checks entirely. - const promptPath = getPromptPath(this.store.getTasksDir(), task.id); - const staleness = await evaluateSpecStaleness({ settings, promptPath, task }); - if (staleness.isStale) { - schedulerLog.warn(`Task ${task.id} specification is stale — ${staleness.reason}`); - await moveTaskToReplanColumn(this.store, task); - await this.store.updateTask(task.id, { status: "needs-replan" }); - await this.store.logEntry(task.id, staleness.reason); - continue; - } - // If staleness evaluation was skipped (missing/unreadable file), continue to - // existing scheduler logic which handles filesystem validation separately. - - // FNXC:MissionSymbolAdmission 2026-08-01-00:00: Apply the same three-way - // admission before the early overlap/priority pass and the final release - // reservation. Otherwise this pass would serialize approved disjoint symbols - // (or misdiagnose unapproved mission work as an overlap) before lock admission. - const earlyMissionAdmission = await decideMissionSymbolAdmission( - task, - this.options.missionStore, - { planApprovalRequired: settings.planApprovalMode === "require-all" }, - ); - if (earlyMissionAdmission.kind === "lineage-blocked") { - await this.store.updateTask(task.id, { status: "queued", blockedBy: null, overlapBlockedBy: null }); - await this.logDispatchQueuedReason(task.id, `queued — mission lineage blocked: ${earlyMissionAdmission.reason}`); - continue; - } - - // Check file scope overlap when enabled. Approved, resolvable mission work - // bypasses this coarse pass and is mutually excluded by its durable symbols. - if (settings.groupOverlappingFiles && earlyMissionAdmission.kind === "coarse-fallback") { - const taskScope = await getFilteredFileScope(task.id); - const coordinationOnlyTask = isCoordinationOnlyTask(task, taskScope); - if (taskScope.length > 0 && !coordinationOnlyTask) { - const activeScopeEntries = Array.from(activeScopes.entries()).sort(([aId], [bId]) => aId.localeCompare(bId)); - const overlapBlockerId = task.overlapBlockedBy || task.blockedBy; - const currentBlockerScope = overlapBlockerId ? activeScopes.get(overlapBlockerId) : undefined; - const hasValidCurrentBlocker = - Boolean(overlapBlockerId) - && Boolean(currentBlockerScope) - && this.pathsOverlap(taskScope, currentBlockerScope!); - - /** - * blockedBy stamping invariants: - * - sticky when still valid: preserve an existing active overlapping blocker - * - deterministic when changing: pick the first overlapping active task by sorted taskId - * - idempotent writes only: update DB only when blockedBy/status must change - */ - const overlappingTaskId = hasValidCurrentBlocker - ? overlapBlockerId - : activeScopeEntries.find(([, ipScope]) => this.pathsOverlap(taskScope, ipScope))?.[0] ?? null; - - const runnableQueuedHigherPriorityScopes = queuedHigherPriorityScopes.filter((queuedCandidate) => { - const queuedTask = queuedHigherPriorityTaskById.get(queuedCandidate.id); - if (!queuedTask) return false; - return isRunnableQueuedOverlapCandidate(queuedTask, tasks, now, activeScopes, queuedCandidate.scope); - }); - - const higherPriorityQueuedOverlap = findHigherPriorityQueuedOverlap( - { - id: task.id, - priority: task.priority, - createdAt: task.createdAt, - scope: taskScope, - }, - runnableQueuedHigherPriorityScopes, - this.pathsOverlap.bind(this), - ); - - if (higherPriorityQueuedOverlap) { - const dependencyBlocker = unmetDeps[0] ?? null; - if ( - task.status !== "queued" - || task.blockedBy !== dependencyBlocker - || task.overlapBlockedBy !== higherPriorityQueuedOverlap.id - ) { - await this.store.updateTask(task.id, { - status: "queued", - blockedBy: dependencyBlocker, - overlapBlockedBy: higherPriorityQueuedOverlap.id, - }); - } - await this.rollbackRunningAgentsForQueuedTodoTask(task.id); - await this.logDispatchQueuedReason( - task.id, - `queued — deferred for higher-priority runnable queued task ${higherPriorityQueuedOverlap.id} (overlap)`, - ); - continue; - } - - if (overlappingTaskId) { - const dependencyBlocker = unmetDeps[0] ?? null; - if ( - task.status !== "queued" - || task.blockedBy !== dependencyBlocker - || task.overlapBlockedBy !== overlappingTaskId - ) { - await this.store.updateTask(task.id, { - status: "queued", - blockedBy: dependencyBlocker, - overlapBlockedBy: overlappingTaskId, - }); - } - - const overlapBlockerTask = tasks.find((candidate) => candidate.id === overlappingTaskId); - await this.rollbackRunningAgentsForQueuedTodoTask(task.id); - const activeLeaseColumn = activeScopeColumns.get(overlappingTaskId) ?? overlapBlockerTask?.column ?? "unknown"; - await this.logDispatchQueuedReason( - task.id, - `queued — blocked by active file-scope lease ${overlappingTaskId} (column=${activeLeaseColumn})`, - ); - continue; - } - - if (task.overlapBlockedBy) { - await this.store.updateTask(task.id, { overlapBlockedBy: null }); - } - } else if (coordinationOnlyTask && task.overlapBlockedBy) { - await this.store.updateTask(task.id, { overlapBlockedBy: null }); - await this.store.logEntry( - task.id, - "coordination/no-commit task bypassed non-implementation overlap lease", - ); - } - } - - /** - * FNXC:Scheduler-Concurrency 2026-06-13-20:08: - * FN-6423 fixes the FN-6420 evidence where queue logs reported `gate=maxWorktrees` with `maxWorktrees used=1/3` and `semaphore used=-9/3`. Recompute capacity at the queue decision point, including tasks already started this tick, so the gate label, memo key, and `started` decision share one authoritative snapshot. - */ - const queuePointCapacity = computeDispatchCapacityDiagnostic(started); - if (queuePointCapacity.available <= 0) { - const reason = formatConcurrencyLimitReason(queuePointCapacity); - const concurrencySignature = formatConcurrencyLimitMemoKey(queuePointCapacity); - await this.logDispatchQueuedReason( - task.id, - reason, - concurrencySignature, - ); - continue; - } - - // Dependencies met — resolve base branch from in-review deps. - // Worktree allocation is deferred to moveTask below, where it - // runs under TaskStore's cross-task allocation lock so it can't - // race against a concurrent manual-move. - const baseBranch = this.resolveBaseBranch(task, tasks); - - // Compare-and-swap: re-read the task to verify it's still in "todo" before dispatching. - // This prevents dispatching a task twice if another schedule() call or user action - // moved it away from "todo" between our initial snapshot and this dispatch attempt. - // The re-entrance guard prevents overlapping schedule() passes, but external events - // (user moves, API calls) can still trigger concurrent state changes. - const freshTask = await this.store.getTask(task.id); - if (!freshTask || freshTask.column !== "todo") { - schedulerLog.log(`Task ${task.id} no longer in "todo" (column=${freshTask?.column ?? "N/A"}) — skipping dispatch`); - continue; - } - /* - FNXC:UserPausedDispatch 2026-07-21-21:30: - The final fresh-read dispatch gate must honor both pause representations because an operator can set userPaused after the scheduler's initial queue snapshot but before worktree allocation. - */ - if (freshTask.paused || freshTask.userPaused) { - schedulerLog.log(`Task ${task.id} is paused — skipping dispatch`); - continue; - } - - const latestSettings = await this.store.getSettings(); - const oscillationSettings = latestSettings as Settings & { - dispatchOscillationSettleMs?: number; - dispatchOscillationThreshold?: number; - dispatchOscillationWindowMs?: number; - }; - const dispatchSettleMs = oscillationSettings.dispatchOscillationSettleMs - ?? DEFAULT_DISPATCH_OSCILLATION_SETTLE_MS; - const dispatchOscillationThreshold = oscillationSettings.dispatchOscillationThreshold - ?? DEFAULT_DISPATCH_OSCILLATION_THRESHOLD; - const dispatchOscillationWindowMs = oscillationSettings.dispatchOscillationWindowMs - ?? DEFAULT_DISPATCH_OSCILLATION_WINDOW_MS; - const recentEngineTodoMovedAt = this.recentEngineTodoRequeues.get(task.id); - if (recentEngineTodoMovedAt) { - if (freshTask.columnMovedAt !== recentEngineTodoMovedAt) { - this.recentEngineTodoRequeues.delete(task.id); - } else { - const movedAtMs = Date.parse(recentEngineTodoMovedAt); - const settleAgeMs = Number.isFinite(movedAtMs) ? Math.max(0, Date.now() - movedAtMs) : dispatchSettleMs; - if (settleAgeMs < dispatchSettleMs) { - schedulerLog.log(`Task ${task.id} was engine-requeued ${settleAgeMs}ms ago — waiting ${dispatchSettleMs}ms settle window before redispatch`); - continue; - } - this.recentEngineTodoRequeues.delete(task.id); - } - } - if (latestSettings.globalPause) { - schedulerLog.log(`Task ${task.id} dispatch aborted — globalPause became active mid-pass`); - continue; - } - if (latestSettings.enginePaused) { - schedulerLog.log(`Task ${task.id} dispatch aborted — enginePaused became active mid-pass`); - continue; - } - - // Resolve effective node for routing - let effectiveNode = resolveEffectiveNode(freshTask, settings); - logTaskRouting(task.id, effectiveNode); - - // Enforce dispatch configuration validation before node-health fallback logic. - if (effectiveNode.nodeId !== undefined && this.options.validateNodeDispatch) { - const validation = await this.options.validateNodeDispatch(effectiveNode.nodeId); - if (!validation.allowed) { - if (!this.wasNodeDispatchValidationBlocked.has(task.id)) { - this.wasNodeDispatchValidationBlocked.add(task.id); - schedulerLog.log(`Task ${task.id} dispatch blocked — ${validation.reason}`); - await this.store.logEntry(task.id, validation.reason); - } - continue; - } - - this.wasNodeDispatchValidationBlocked.delete(task.id); - } - - // Enforce unavailable-node policy + owning-node handoff policy - // FN-4832: this guard currently applies only when node routing is explicit; local routing still relies on checkout 409 claim backstops. - if (effectiveNode.nodeId !== undefined && this.options.nodeHealthMonitor) { - const localNodeId = this.options.localNodeId ?? "local"; - if (freshTask.checkoutNodeId && freshTask.checkedOutBy && freshTask.checkoutNodeId !== localNodeId) { - const ownerNodeHealth = this.options.nodeHealthMonitor.getNodeHealth(freshTask.checkoutNodeId); - // FN-4832 + AGENTS.md Checkout Leasing: never dispatch tasks with active foreign ownership; policy decides park/reassign. - const handoffDecision = decideOwningNodeHandoff({ - task: freshTask, - ownerNodeId: freshTask.checkoutNodeId, - ownerNodeHealth, - localNodeId, - handoffPolicy: settings.owningNodeHandoffPolicy, - }); - - if (handoffDecision.action === "park") { - if (!this.wasNodeBlocked.has(task.id)) { - this.wasNodeBlocked.add(task.id); - if (ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online") { - await this.emitNodeUnreachableRecoveryAudit(freshTask, { - ownerNodeId: freshTask.checkoutNodeId, - ownerNodeHealth, - handoffAction: handoffDecision.action, - handoffReason: handoffDecision.reason, - decisionPath: "scheduler-handoff-park", - newColumn: freshTask.column, - dispatchNodeBefore: effectiveNode.nodeId, - dispatchNodeAfter: effectiveNode.nodeId, - }); - } - const reason = `Owning-node handoff parked dispatch: ${handoffDecision.reason}`; - schedulerLog.log(`Task ${task.id} dispatch blocked — ${reason}`); - await this.store.logEntry(task.id, reason); - try { - await this.store.recordRunAuditEvent?.({ - taskId: freshTask.id, - agentId: "scheduler", - runId: generateSyntheticRunId("scheduler", freshTask.id), - domain: "database", - mutationType: "node:handoff:parked", - target: freshTask.id, - metadata: { - taskId: freshTask.id, - ownerNodeId: freshTask.checkoutNodeId, - ownerNodeHealth: - ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online" - ? ownerNodeHealth - : "unknown", - localNodeId, - handoffPolicy: settings.owningNodeHandoffPolicy, - decisionReason: handoffDecision.reason, - source: "scheduler.dispatch", - }, - }); - } catch (error) { - schedulerLog.warn(`Task ${task.id} failed to emit node:handoff:parked audit: ${error instanceof Error ? error.message : String(error)}`); - } - } - continue; - } - - await this.store.logEntry(task.id, `Owning-node handoff applied: ${handoffDecision.reason}`); - try { - await this.store.recordRunAuditEvent?.({ - taskId: freshTask.id, - agentId: "scheduler", - runId: generateSyntheticRunId("scheduler", freshTask.id), - domain: "database", - mutationType: handoffDecision.action === "reassign-local" ? "node:handoff:reassign-local" : "node:handoff:reassign-any", - target: freshTask.id, - metadata: { - taskId: freshTask.id, - ownerNodeId: freshTask.checkoutNodeId, - ownerNodeHealth: - ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online" - ? ownerNodeHealth - : "unknown", - localNodeId, - handoffPolicy: settings.owningNodeHandoffPolicy, - decisionReason: handoffDecision.reason, - source: "scheduler.dispatch", - }, - }); - } catch (error) { - schedulerLog.warn(`Task ${task.id} failed to emit node:handoff audit: ${error instanceof Error ? error.message : String(error)}`); - } - const dispatchNodeBefore = effectiveNode.nodeId; - if (handoffDecision.action === "reassign-local") { - effectiveNode = { nodeId: undefined, source: "local" }; - } - if (ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online") { - await this.emitNodeUnreachableRecoveryAudit(freshTask, { - ownerNodeId: freshTask.checkoutNodeId, - ownerNodeHealth, - handoffAction: handoffDecision.action, - handoffReason: handoffDecision.reason, - decisionPath: - handoffDecision.action === "reassign-local" - ? "scheduler-handoff-reassign-local" - : "scheduler-handoff-reassign-any", - newColumn: freshTask.column, - dispatchNodeBefore, - dispatchNodeAfter: effectiveNode.nodeId, - }); - } - } - - const nodeHealth = effectiveNode.nodeId - ? this.options.nodeHealthMonitor.getNodeHealth(effectiveNode.nodeId) - : undefined; - const decision = applyUnavailableNodePolicy({ - effectiveNode, - nodeHealth, - policy: settings.unavailableNodePolicy, - }); - - if (!decision.allowed) { - if (!this.wasNodeBlocked.has(task.id)) { - this.wasNodeBlocked.add(task.id); - schedulerLog.log(`Task ${task.id} dispatch blocked — ${decision.reason}`); - await this.store.logEntry(task.id, decision.reason); - } - continue; - } - - this.wasNodeBlocked.delete(task.id); - - if (decision.fallbackToLocal) { - schedulerLog.log(`Task ${task.id} falling back to local — ${decision.reason}`); - await this.store.logEntry(task.id, decision.reason); - effectiveNode = { nodeId: undefined, source: "local" }; - } - } - - if (latestSettings.ephemeralAgentsEnabled === false && !freshTask.assignedAgentId) { - if (!this.options.agentStore) { - if (!loggedMissingAgentStoreThisPass) { - loggedMissingAgentStoreThisPass = true; - schedulerLog.warn("ephemeralAgentsEnabled=false but scheduler has no agentStore; falling back to legacy dispatch behavior"); - } - } else { - const selectedAgent = await selectPermanentAgentForTask({ - task: freshTask, - agentStore: this.options.agentStore, - taskStore: this.store, - }); - - if (!selectedAgent) { - await this.store.updateTask(task.id, { status: "queued" }); - if (!this.wasPermanentAgentUnavailable.has(task.id)) { - await this.logDispatchQueuedReason( - task.id, - "queued — no permanent executor available (ephemeral agents disabled)", - ); - this.wasPermanentAgentUnavailable.add(task.id); - } - continue; - } - - await this.store.updateTask(task.id, { assignedAgentId: selectedAgent.id }); - await this.store.logEntry( - task.id, - `Auto-assigned to permanent agent ${selectedAgent.id} (ephemeral agents disabled)`, - ); - this.wasPermanentAgentUnavailable.delete(task.id); - } - } else { - this.wasPermanentAgentUnavailable.delete(task.id); - } - - // Clear status, reserve worktree path, and then move to in-progress. - // Reset mergeRetries so a fresh execution gets a fresh merge budget — - // otherwise a task whose previous run exhausted its 3 retries (e.g. - // verification failure that was later cleared) lands back in in-review - // with mergeRetries=MAX, the merger refuses it (canMergeTask false), - // and the ghost-review fallback bounces it back to todo every 10 min - // before the 30-min cooldown can elapse — infinite loop. See FN-3305. - const dispatchTimestamp = new Date().toISOString(); - const lastDispatchAtMs = freshTask.lastDispatchAt ? Date.parse(freshTask.lastDispatchAt) : Number.NaN; - const priorDispatchWithinWindow = Number.isFinite(lastDispatchAtMs) - && Date.now() - lastDispatchAtMs <= dispatchOscillationWindowMs; - const nextDispatchStormCount = priorDispatchWithinWindow - ? (freshTask.dispatchStormCount ?? 0) + 1 - : 1; - - if (nextDispatchStormCount > dispatchOscillationThreshold) { - const oscillationError = freshTask.error - ?? `DISPATCH_OSCILLATION: detected ${nextDispatchStormCount} todo↔in-progress cycles within ${dispatchOscillationWindowMs}ms. Task auto-paused for operator review.`; - await this.store.updateTask(task.id, { - dispatchStormCount: nextDispatchStormCount, - lastDispatchAt: dispatchTimestamp, - paused: true, - pausedReason: "dispatch-oscillation", - status: freshTask.status ?? "queued", - error: oscillationError, - }); - await this.store.logEntry( - task.id, - `Dispatch oscillation auto-paused after ${nextDispatchStormCount} cycles within ${dispatchOscillationWindowMs}ms`, - ); - await this.store.appendAgentLog?.( - task.id, - "Dispatch oscillation detected — task auto-paused for operator review", - "text", - `cycleCount=${nextDispatchStormCount} windowMs=${dispatchOscillationWindowMs}`, - ); - await this.store.recordRunAuditEvent?.({ - taskId: task.id, - agentId: "scheduler", - runId: generateSyntheticRunId("scheduler-dispatch-oscillation", task.id), - domain: "database", - mutationType: "task:dispatch-oscillation-terminalized", - target: task.id, - metadata: { - taskId: task.id, - cycleCount: nextDispatchStormCount, - windowMs: dispatchOscillationWindowMs, - lastMoveSource: recentEngineTodoMovedAt ? "engine" : "scheduler", - }, - }); - schedulerLog.warn(`Task ${task.id} auto-paused after dispatch oscillation threshold ${dispatchOscillationThreshold} was exceeded (${nextDispatchStormCount} cycles)`); - continue; - } - - try { - /* - FNXC:UserPausedDispatch 2026-07-21-21:45: - Scheduler dispatch predicates the todo-to-in-progress transition on both pause flags under the task lock. No awaited routing, metadata, or worktree preparation gap may let a concurrent operator pause lose to stale scheduler state. - */ - const move = await this.store.moveTaskIf( - task.id, - "in-progress", - (live) => live.column === "todo" && live.paused !== true && live.userPaused !== true, - { - moveSource: "scheduler", - allocateWorktree: (reservedNames) => - this.planWorktreePath(task, settings.worktreeNaming, reservedNames, settings), - }, - ); - if (!move.moved) { - schedulerLog.log(`Task ${task.id} became paused or left todo before dispatch — skipping`); - continue; - } - Object.assign(task, move.task); - } catch (error) { - if (error instanceof TransitionRejectionError && error.rejection.code === "capacity-exhausted") { - await this.store.updateTask(task.id, { status: "queued" }); - const reason = error.message || "queued — in-progress column at capacity"; - await this.logDispatchQueuedReason(task.id, reason, `capacity-exhausted:${reason}`); - continue; - } - throw error; - } - schedulerLog.log(`Starting ${task.id}: ${task.title || task.id} (deps satisfied)`); - await this.store.updateTask(task.id, { - status: null, - blockedBy: null, - executionStartBranch: baseBranch ?? undefined, - effectiveNodeId: effectiveNode.nodeId ?? null, - effectiveNodeSource: effectiveNode.source, - mergeRetries: 0, - }); - await this.store.updateTask(task.id, { - dispatchStormCount: nextDispatchStormCount, - lastDispatchAt: dispatchTimestamp, - }); - this.recentEngineTodoRequeues.delete(task.id); - this.wasNodeBlocked.delete(task.id); - this.wasNodeDispatchValidationBlocked.delete(task.id); - this.wasPermanentAgentUnavailable.delete(task.id); - this.clearDispatchQueuedReasonMemo(task.id); - await this.store.logEntry(task.id, `Node routing resolved: ${effectiveNode.nodeId ?? "local"} (source: ${effectiveNode.source})`); - this.options.onSchedule?.(task); - started++; - - // Track newly started task's file scope for overlap with remaining todo tasks - if (settings.groupOverlappingFiles) { - const scope = await getFilteredFileScope(task.id); - if (scope.length > 0 && !isCoordinationOnlyTask(task, scope)) setActiveScopeLease(task.id, scope, "in-progress"); - } - } - - await this.emitHighOverlapFanoutWarnings(tasks); - - const staleWarningWindows = [settings.staleInProgressWarningMs, settings.staleInReviewWarningMs] - .filter((value): value is number => typeof value === "number" && Number.isFinite(value) && value > 0); - const minWarningMs = staleWarningWindows.length > 0 ? Math.min(...staleWarningWindows) : 0; - if (minWarningMs > 0 && Date.now() - this.lastStaleTaskReportAt >= minWarningMs) { - try { - await this.staleTaskReporter.report(); - this.lastStaleTaskReportAt = Date.now(); - } catch (error) { - schedulerLog.warn("Stale task reporter failed", error); - } - } - - if (settings.backlogPressureAlertEnabled !== false && Date.now() - this.lastBacklogPressureReportAt >= 60_000) { - try { - await this.backlogPressureReporter.report(); - } catch (error) { - schedulerLog.warn("Backlog pressure reporter failed", error); - } finally { - this.lastBacklogPressureReportAt = Date.now(); - } - } - - if (Date.now() - this.lastUnlinkedMissionsAdvisoryReportAt >= 60_000) { - try { - await this.unlinkedMissionsAdvisoryReporter.report(); - } catch (error) { - schedulerLog.warn("Unlinked missions advisory reporter failed", error); - } finally { - this.lastUnlinkedMissionsAdvisoryReportAt = Date.now(); - } - } } catch (err) { schedulerLog.error("Scheduling error:", err); } finally {