U11: delete the unreachable legacy todo dispatcher from scheduler.schedule() (-929 lines, pure deletion) (#2505)

Based on `main`. **Pure deletion — no behavior change**, because the
deleted code cannot execute.

## Found while trying to convert it

This started as a U11 slice to make the scheduler's dispatch path
resolve its column by trait. Per the lesson from the dependency-blocked
feature I checked reachability *before* converting:

```ts
function shouldRunWorkflowColumnScheduler(_settings: Settings): boolean {
  return true;                       // parameter UNUSED, body a literal
}
...
if (shouldRunWorkflowColumnScheduler(settings)) {
  await this.runHoldReleaseSweepPass(tasks, settings);
  ...
  return;                            // UNCONDITIONAL, at the block's own depth
}
<929 lines of legacy pull-from-todo dispatcher>   // unreachable
```

The guard takes an **unused** parameter and returns a **literal**, so
the branch is statically always taken, and it ends in an **unconditional
`return`**. Everything after it in `schedule()` is unreachable.

`tsc` doesn't flag it because the condition is a function call rather
than a literal — which is exactly why 929 lines survived the U6 cutover.
The replacement was added *in front of* the old dispatcher rather than
*instead of* it, and the in-file comment says so outright:

> the hold/release sweep owns todo→in-progress pickup, so do not fall
through into the legacy pull-from-todo dispatcher after the sweep runs

## Why this matters beyond line count

**4 of the 15 `"todo"` literals in `scheduler.ts` live in this dead
region.** Converting them would have been pure waste — and worse, it
would have reported progress against the U11 critical path while
changing nothing. 11 live sites remain and are the real work.

## Corroborating evidence

Six imports became unused and are removed with it:
`resolveDependencyOrder`, `sortTasksByPriorityFanoutThenAgeAndId`,
`buildUnblockWeightMap`, `TransitionRejectionError`,
`isUnplannedSeedPrompt`, `DEFAULT_WORKFLOW_POOL_ID`.

That the dead region was their **only** consumer in this file is itself
evidence: a live dispatcher would still need dependency ordering and
priority sorting.

## Why no new test

The proof here is **static, not behavioral** — an unconditional `return`
before the code. A test cannot demonstrate absence of execution more
strongly than the control flow already does, and one that passed both
before and after would be theatre.

The evidence that nothing depended on it: **all 100 scheduler tests and
the full merge gate pass unchanged.**

## Measured

`scheduler.ts` **3,726 → 2,797 = −929 lines.**

Unlike every consolidation in this program, this is a **genuine net
reduction** — nothing was moved elsewhere.

## Verification

100 scheduler tests green across all 7 scheduler suites; merge gate
green (307 + 10 + 71); tsc clean; lint clean.

No changeset: `@fusion/engine` is private.

🤖 Generated with [Claude Code](https://claude.com/claude-code)
This commit is contained in:
gsxdsm
2026-07-28 17:15:23 -07:00
committed by GitHub
parent 6ee20d9817
commit a2b4ca76ac

View File

@@ -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<string>();
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<string, boolean>();
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<string, string[]>();
const activeScopeColumns = new Map<string, Task["column"]>();
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<string, Task>();
const overlapIgnorePaths = settings.overlapIgnorePaths ?? [];
const filteredScopeByTaskId = new Map<string, string[]>();
const getFilteredFileScope = async (taskId: string): Promise<string[]> => {
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<string, boolean>();
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 {