FN-5742: add dual-observe merge-request parity seam

Add a dual-observe merge-request seam that preserves legacy authority while emitting parity signals.

- Add scheduler parity helpers for dependency satisfaction and shadow lease state, keeping legacy scheduling decisions authoritative.
- Extend ProjectEngine/run-audit merge dequeue flow with shadow-candidate selection and parity audit emission.
- Add reliability interaction coverage for dual-observe dependency parity, shadow dequeue parity, lease parity, and runnable guard behavior.
- Document FN-5742 reliability backstop coverage in architecture docs.

Files changed:
 docs/architecture.md                               |   2 +
 packages/engine/src/__tests__/reliability-interactions/dual-observe-merge-seam.test.ts                | 193 +++++++++++++++++++++
 packages/engine/src/__tests__/scheduler.test.ts    |  64 +++++++
 packages/engine/src/project-engine.ts              |  42 +++++
 packages/engine/src/run-audit.ts                   |   3 +
 packages/engine/src/scheduler.ts                   | 158 +++++++++++++++--
 6 files changed, 448 insertions(+), 14 deletions(-)

Fusion-Task-Id: FN-5742

Fusion-Task-Lineage: dc4c0212-fe75-4860-bf06-2765a0b18d38
This commit is contained in:
gsxdsm
2026-05-30 18:52:19 -07:00
parent c0b749cee9
commit 8609669a56
6 changed files with 448 additions and 14 deletions

View File

@@ -1741,6 +1741,7 @@ This section preserves the detailed lifecycle/self-healing contracts that were f
- **Orphaned execution sweep is observation-only (FN-5337)**: `recoverOrphanedExecutions` only annotates stale in-progress candidates with `task:orphan-detected-no-action` and `[orphan-detected] ... no action (operator-decides)` logs. It must never move `in-progress`/`in-review` backward to `todo` or mutate lease/worktree metadata. Proof-based backward recovery remains exclusively in `recoverInProgressLimbo` (FN-5219), `RestartRecoveryCoordinator`, `recoverMissingWorktreeReviewFailures`, and explicit executor/merger failure paths. Reintroducing lifecycle mutation here requires hard git/session proof gating plus CEO+CTO+PM sign-off.
- **Self-owned reclaim resume-limbo escalation (FN-5704)**: `reclaimSelfOwnedBranchConflicts` tracks `resumeLimboCount`, `resumeLimboTipSha`, and `resumeLimboStepSignature` for in-progress reclaim/unpause loops. If reclaim finds no progress (same tip, same step-status signature, and no active-session signal) for `MAX_NO_PROGRESS_RESUME_ATTEMPTS` consecutive sweeps, self-healing escalates by moving the task to `todo` with `preserveWorktree: true`, `preserveProgress: true`, and `preserveResumeState: true` instead of endlessly re-arming resume. Escalation emits `task:resume-limbo-escalated` run-audit metadata (`frozenTipSha`, `idleMs`, `resumeAttemptCount`, `currentStep`) and resets the limbo counter.
- **Merge-request shadow contract (FN-5741 Phase 1)**: `mergeRequestContractShadowEnabled` defaults OFF. OFF means no writes to `merge_requests` or `completion_handoff_markers` and the legacy lifecycle remains authoritative. ON enables write-only shadow persistence: executor/self-healing append `task:completion-handoff-accepted` marker+record writes only after successful legacy `handoffToReview`, and merger mirrors `merge:request-enqueued` plus `queued → running → succeeded` (or `manual-required` for `autoMerge:false`) transitions. Phase 1 never reads these shadow records for column movement, dependency checks, lease arbitration, merge dequeue, or FN-5479/FN-5704 limbo recovery decisions.
- **Dual-observe parity seam (FN-5742 Phase 2)**: with the same flag ON, legacy remains authoritative while shadow reads compute/emit parity telemetry only. Scheduler emits `merge:dependency-parity-diff` when `in-review|done|archived` dependency satisfaction diverges from completion-handoff marker satisfaction, and `merge:lease-parity-diff` when legacy in-review overlap leasing diverges from shadow lease decomposition. Merger emits `merge:request-dequeued-shadow` (agree/disagree metadata) by comparing legacy dequeue selection to shadow merge-request selection while explicitly skipping `manual-required` rows. Phase 3 dequeue cutover is gated on sustained parity (low disagreement rate) from these additive events; no lifecycle authority changes in Phase 2.
- **No-progress churn terminalization (FN-5168)**: `StuckTaskDetector` now tracks ignored `fn_task_update` rebuffs via `recordIgnoredStepUpdate(taskId)` and, after one loop/compact-and-resume recovery has already fired in the same `execute()` lifecycle, escalates `ignoredStepUpdateCount >= 25` to the terminal reason `no-progress-churn`. `SelfHealingManager.checkStuckBudget()` maps that reason directly to `STUCK_NO_PROGRESS_CHURN`, emits `task:stuck-no-progress-churn-terminalized` with `{ taskId, ignoredStepUpdateCount, stuckKillStreak, lastReason }`, and parks the task in `in-review` without consuming the normal stuck-kill budget. Under FN-5147 `autoMerge: false`, that failed in-review task remains terminal-until-merged just like `STUCK_LOOP_EXHAUSTED`; the new class adds an earlier bounded exit, not a re-execution path.
- **Landed-files attribution (FN-5103)**: Rebase-strategy `mergeDetails.landedFiles` / `filesChanged` / `insertions` / `deletions` are captured from task-attributable commits only via `filterFilesToOwnTaskCommits` (subject-prefix + trailer + bracket-prefix evidence), tagged `landedFilesAttributionRestricted: true`. Zero own commits → `landedFiles: []` and `noOpVerifiedShortCircuit: true`. FN-5304 guard: when `<rebaseBaseSha>..HEAD` reports zero own commits, merger must also validate the source `fusion/<id>` tip; if that source tip still has attributable own commits relative to `rebaseBaseSha`, throw `SilentNoOpAttributionMismatchError`, refuse writing `mergeConfirmed: true`, park the task in `in-review` with `status: "failed"`, and emit `merge:no-op-attribution-mismatch`. If source ref is unavailable, skip with diagnostic + `merge:no-op-attribution-mismatch-skipped` (`reason: "source-ref-unavailable"`). Attribution-helper failures fall back to the unrestricted `rebaseBaseSha..sha` walk and set `landedFilesCaptureFallback: 'attribution-failed'`. Self-healing `recoverDoneTaskMergeMetadata` skips reconcile when `landedFilesAttributionRestricted` or `noOpVerifiedShortCircuit` is set so the narrower set is not overwritten with the full range. Squash-strategy capture is unchanged.
- **Soft-delete scheduler invalidation (FN-5137)**: `task:deleted` events must invalidate `AutoClaimSnapshotManager` and clear scheduler bookkeeping (`pausedTaskIds`, `failedTaskIds`, `wasNodeDispatchValidationBlocked`, `wasNodeBlocked`); `executor.execute()` / `resumeOrphaned()` / `resumeTaskForAgent()` refuse any task with `deletedAt` set.
@@ -1771,6 +1772,7 @@ Reliability-layer changes are in scope. Interaction regression backstops live in
- FN-5715 backstop: `packages/engine/src/__tests__/reliability-interactions/mission-validation-trigger-gap.test.ts` locks the mission-validation trigger invariant so done mission-linked tasks still start validation when the mission loop was stopped, startup recovery replays done implementing features with unpassed assertions, and recovery remains idempotent for already-passed features.
- FN-5738 backstop: `packages/engine/src/__tests__/reliability-interactions/mission-validation-trigger-gap.test.ts` extends coverage so zero-assertion auto-pass deterministically advances `loopState` to `passed`, sets `lastValidatorStatus="passed"`, emits `validation_auto_passed_no_assertions`, and does not re-fire on repeated recovery passes.
- FN-5741 backstop: `packages/engine/src/__tests__/reliability-interactions/merge-request-shadow-handoff.test.ts` guards Phase-1 merge-request contract shadow writes: flag OFF is a no-op, flag ON writes marker/record strictly after legacy handoff, and `autoMerge:false` remains `manual-required` without shadow running transitions.
- FN-5742 backstop: `packages/engine/src/__tests__/reliability-interactions/dual-observe-merge-seam.test.ts` guards Phase-2 dual-observe invariants: legacy dependency satisfaction remains authoritative while parity diffs emit, and shadow dequeue selection never advances `manual-required` rows.
- FN-5337 backstop: `packages/engine/src/__tests__/reliability-interactions/orphan-detected-no-requeue.test.ts` locks observation-only orphan detection across FN-5279 repro metadata desync, worktree-present and worktree-missing candidates, FN-5219 ordering, FN-5147 in-review isolation, FN-5083 branch-cleared composition, lease-manager non-invocation, and per-sweep idempotent audit emission.
- FN-5256 backstop: `packages/engine/src/__tests__/reliability-interactions/dependency-cycle-reconcile.test.ts` covers persisted dependency-cycle detection via `reconcileDependencyCycles`, bounded umbrella-back-edge auto-repair, ambiguous-cycle observe-only behavior, composition ordering with `reconcileSelfDefeatingDependencies`, and the post-sweep write-time guard invariant. Core write-boundary regressions (FN-5240/5241/5242 signature, indirect cycle, umbrella back-edge rejection) live in `packages/core/src/__tests__/store-dependency-cycle.test.ts`.
- FN-5325 backstop: `packages/engine/src/__tests__/reliability-interactions/scheduler-overlap-priority-inversion.test.ts` covers queued-overlap priority/age deferral, equal-priority age ordering, FN-4969 fanout composition, and one-shot per-pass `scheduler:overlap-priority-inversion` audit surfacing against running lower-priority blockers.

View File

@@ -0,0 +1,193 @@
import { describe, expect, it, vi } from "vitest";
import {
computeShadowLeaseParityState,
getUnmetSchedulingDependencies,
isRunnableQueuedOverlapCandidate,
} from "../../scheduler.js";
import { ProjectEngine } from "../../project-engine.js";
import { classifyTransientMergeError } from "../../transient-merge-error-classifier.js";
describe("FN-5742 dual-observe merge seam", () => {
it("emits no dependency parity diff when legacy and marker agree", () => {
const task = { id: "FN-T", dependencies: ["FN-DEP"] } as any;
const dep = { id: "FN-DEP", column: "done" } as any;
const onParityDiff = vi.fn();
const unmet = getUnmetSchedulingDependencies(task, [task, dep], {
markerAcceptedByTaskId: new Map([["FN-DEP", true]]),
onParityDiff,
});
expect(unmet).toEqual([]);
expect(onParityDiff).not.toHaveBeenCalled();
});
it("keeps legacy dependency satisfaction authoritative while emitting parity diff", () => {
const task = { id: "FN-T", dependencies: ["FN-DEP"] } as any;
const dep = { id: "FN-DEP", column: "in-review" } as any;
const onParityDiff = vi.fn();
const unmet = getUnmetSchedulingDependencies(task, [task, dep], {
markerAcceptedByTaskId: new Map([["FN-DEP", false]]),
onParityDiff,
});
expect(unmet).toEqual([]);
expect(onParityDiff).toHaveBeenCalledWith(
expect.objectContaining({
taskId: "FN-T",
dependencyId: "FN-DEP",
legacySatisfied: true,
markerSatisfied: false,
}),
);
});
it("shadow dequeue selector skips manual-required records", () => {
const fakeEngine = {
mergeQueue: ["FN-A", "FN-B", "FN-C"],
runtime: {
getTaskStore: () => ({
getMergeRequestRecord: (taskId: string) => {
if (taskId === "FN-A") return { state: "manual-required" };
if (taskId === "FN-B") return { state: "queued" };
return null;
},
}),
},
};
const candidate = (ProjectEngine.prototype as any).getShadowMergeRequestCandidateId.call(fakeEngine);
expect(candidate).toBe("FN-B");
});
it("emits shadow dequeue parity audit with agree metadata", () => {
const recordRunAuditEvent = vi.fn();
const fakeEngine = {
runtime: {
getTaskStore: () => ({ recordRunAuditEvent }),
},
};
(ProjectEngine.prototype as any).emitMergeRequestShadowDequeueParity.call(fakeEngine, "FN-LEGACY", "FN-SHADOW");
expect(recordRunAuditEvent).toHaveBeenCalledWith(
expect.objectContaining({
taskId: "FN-LEGACY",
mutationType: "merge:request-dequeued-shadow",
metadata: expect.objectContaining({
legacyTaskId: "FN-LEGACY",
shadowTaskId: "FN-SHADOW",
agree: false,
}),
}),
);
});
it("marks dequeue parity as agree when legacy and shadow match", () => {
const recordRunAuditEvent = vi.fn();
const fakeEngine = {
runtime: {
getTaskStore: () => ({ recordRunAuditEvent }),
},
};
(ProjectEngine.prototype as any).emitMergeRequestShadowDequeueParity.call(fakeEngine, "FN-LEGACY", "FN-LEGACY");
expect(recordRunAuditEvent).toHaveBeenCalledWith(
expect.objectContaining({
mutationType: "merge:request-dequeued-shadow",
metadata: expect.objectContaining({
agree: true,
}),
}),
);
});
it("is a no-op when shadow dequeue API is unavailable", () => {
const fakeEngine = {
mergeQueue: ["FN-A"],
runtime: {
getTaskStore: () => ({}),
},
};
const candidate = (ProjectEngine.prototype as any).getShadowMergeRequestCandidateId.call(fakeEngine);
expect(candidate).toBeNull();
});
it("computes lease parity so manual-required does not create a merge lock", () => {
const parity = computeShadowLeaseParityState("manual-required");
expect(parity.shadowExecutorLeaseApplied).toBe(false);
expect(parity.shadowMergeLockApplied).toBe(false);
expect(parity.shadowLeaseApplied).toBe(false);
});
it("computes lease parity so queued requests create a merge lock", () => {
const parity = computeShadowLeaseParityState("queued");
expect(parity.shadowExecutorLeaseApplied).toBe(false);
expect(parity.shadowMergeLockApplied).toBe(true);
expect(parity.shadowLeaseApplied).toBe(true);
});
it("keeps dependency checks unchanged when parity observer options are absent", () => {
const task = { id: "FN-T", dependencies: ["FN-DEP"] } as any;
const dep = { id: "FN-DEP", column: "in-review" } as any;
const unmet = getUnmetSchedulingDependencies(task, [task, dep]);
expect(unmet).toEqual([]);
});
it("does not block unrelated executor dispatch when merge lane is busy", () => {
const now = Date.now();
const activeScopes = new Map<string, string[]>([["FN-MERGE", ["packages/engine/src/merger.ts"]]]);
const todo = {
id: "FN-UNRELATED",
column: "todo",
status: "queued",
paused: false,
userPaused: false,
dependencies: [],
nextRecoveryAt: undefined,
} as any;
const runnable = isRunnableQueuedOverlapCandidate(todo, [todo], now, activeScopes, ["docs/architecture.md"]);
expect(runnable).toBe(true);
});
it("classifies transient merge errors without changing scheduler runnable decisions", () => {
const transient = classifyTransientMergeError(
"lease-handoff-failed: target-not-queued while attempting merge handoff",
);
expect(transient).toBe("lease-handoff-target-not-queued");
const todo = {
id: "FN-RUN",
column: "todo",
status: "queued",
paused: false,
userPaused: false,
dependencies: [],
nextRecoveryAt: undefined,
} as any;
expect(isRunnableQueuedOverlapCandidate(todo, [todo], Date.now())).toBe(true);
});
it("shadow dequeue helpers do not touch limbo recovery counters", () => {
const updateTask = vi.fn();
const fakeEngine = {
runtime: {
getTaskStore: () => ({
recordRunAuditEvent: vi.fn(),
updateTask,
}),
},
mergeQueue: ["FN-A"],
};
(ProjectEngine.prototype as any).emitMergeRequestShadowDequeueParity.call(fakeEngine, "FN-A", "FN-A");
(ProjectEngine.prototype as any).getShadowMergeRequestCandidateId.call(fakeEngine);
expect(updateTask).not.toHaveBeenCalled();
});
});

View File

@@ -8,6 +8,7 @@ import {
findHigherPriorityQueuedOverlap,
isCoordinationOnlyTask,
isRunnableQueuedOverlapCandidate,
getUnmetSchedulingDependencies,
} from "../scheduler.js";
import { AgentSemaphore } from "../concurrency.js";
import type { TaskStore, Task, TaskDetail } from "@fusion/core";
@@ -328,6 +329,69 @@ describe("isCoordinationOnlyTask", () => {
});
});
describe("getUnmetSchedulingDependencies", () => {
it("keeps legacy in-review satisfaction authoritative while emitting parity diff", () => {
const task = createMockTask({ id: "FN-T", dependencies: ["FN-DEP"] });
const dep = createMockTask({ id: "FN-DEP", column: "in-review" });
const diffs: Array<{ dependencyId: string; legacySatisfied: boolean; markerSatisfied: boolean }> = [];
const unmet = getUnmetSchedulingDependencies(task, [task, dep], {
markerAcceptedByTaskId: new Map([["FN-DEP", false]]),
onParityDiff: (diff) => {
diffs.push({
dependencyId: diff.dependencyId,
legacySatisfied: diff.legacySatisfied,
markerSatisfied: diff.markerSatisfied,
});
},
});
expect(unmet).toEqual([]);
expect(diffs).toEqual([
{
dependencyId: "FN-DEP",
legacySatisfied: true,
markerSatisfied: false,
},
]);
});
it("does not emit parity diff when legacy and marker paths agree dependency is satisfied", () => {
const task = createMockTask({ id: "FN-T", dependencies: ["FN-DEP"] });
const dep = createMockTask({ id: "FN-DEP", column: "done" });
const onParityDiff = vi.fn();
const unmet = getUnmetSchedulingDependencies(task, [task, dep], {
markerAcceptedByTaskId: new Map([["FN-DEP", false]]),
onParityDiff,
});
expect(unmet).toEqual([]);
expect(onParityDiff).not.toHaveBeenCalled();
});
it("does not emit parity diff when legacy and marker paths agree dependency is unmet", () => {
const task = createMockTask({ id: "FN-T", dependencies: ["FN-DEP"] });
const dep = createMockTask({ id: "FN-DEP", column: "in-progress" });
const onParityDiff = vi.fn();
const unmet = getUnmetSchedulingDependencies(task, [task, dep], {
markerAcceptedByTaskId: new Map([["FN-DEP", false]]),
onParityDiff,
});
expect(unmet).toEqual(["FN-DEP"]);
expect(onParityDiff).not.toHaveBeenCalled();
});
it("preserves legacy behavior when parity options are omitted", () => {
const task = createMockTask({ id: "FN-T", dependencies: ["FN-DEP"] });
const dep = createMockTask({ id: "FN-DEP", column: "in-review" });
expect(getUnmetSchedulingDependencies(task, [task, dep])).toEqual([]);
});
});
describe("isRunnableQueuedOverlapCandidate", () => {
const now = new Date("2026-01-01T00:00:00.000Z").getTime();

View File

@@ -1300,6 +1300,43 @@ export class ProjectEngine {
return undefined;
}
private getShadowMergeRequestCandidateId(): string | null {
const store = this.runtime.getTaskStore() as TaskStore & {
getMergeRequestRecord?: (taskId: string) => { state: string } | null;
};
if (typeof store.getMergeRequestRecord !== "function") {
return null;
}
for (const queuedTaskId of this.mergeQueue) {
const record = store.getMergeRequestRecord(queuedTaskId);
if (!record) continue;
if (record.state === "manual-required") continue;
if (record.state === "queued" || record.state === "retrying" || record.state === "running") {
return queuedTaskId;
}
}
return null;
}
private emitMergeRequestShadowDequeueParity(legacyTaskId: string, shadowTaskId: string | null): void {
const agree = shadowTaskId === legacyTaskId;
const store = this.runtime.getTaskStore();
void store.recordRunAuditEvent?.({
taskId: legacyTaskId,
agentId: "merger",
runId: generateSyntheticRunId("merger-shadow-dequeue", legacyTaskId),
domain: "database",
mutationType: "merge:request-dequeued-shadow",
target: legacyTaskId,
metadata: {
legacyTaskId,
shadowTaskId,
agree,
},
});
}
private internalEnqueueMerge(taskId: string): boolean {
if (this.shuttingDown) return false;
if (this.mergeActive.has(taskId)) {
@@ -1405,8 +1442,13 @@ export class ProjectEngine {
const cwd = this.config.workingDirectory;
while (this.mergeQueue.length > 0 && !this.shuttingDown) {
const shadowCandidateTaskId = this.getShadowMergeRequestCandidateId();
const taskId = await this.pickNextMergeTaskId(store);
if (!taskId) break;
const shadowSettings = await store.getSettings();
if (shadowSettings.mergeRequestContractShadowEnabled === true) {
this.emitMergeRequestShadowDequeueParity(taskId, shadowCandidateTaskId);
}
// pickNextMergeTaskId awaits store.getTask; re-check shutdown so we
// don't start a merge whose queue entry was cleared by stop().
if (this.shuttingDown) break;

View File

@@ -401,6 +401,9 @@ export type DatabaseMutationType =
| "task:unpause"
| "task:dependency:add"
| "merge:request-enqueued"
| "merge:dependency-parity-diff"
| "merge:lease-parity-diff"
| "merge:request-dequeued-shadow"
| "mergeQueue:lease-target-unavailable"
| "mergeQueue:enqueue-rejected"
| "mergeQueue:stale-lease-on-column-exit"

View File

@@ -167,10 +167,63 @@ export function isCoordinationOnlyTask(task: Task, scope: string[]): boolean {
return isCoordinationSafeScope(scope);
}
export function getUnmetSchedulingDependencies(task: Task, tasks: Task[]): string[] {
function isLegacyDependencySatisfied(dep: Task | undefined): boolean {
return !!dep && (dep.column === "done" || dep.column === "in-review" || dep.column === "archived");
}
function isMarkerDependencySatisfied(dep: Task | undefined, markerAccepted: boolean): boolean {
if (!dep) return false;
if (dep.column === "done" || dep.column === "archived") return true;
return markerAccepted;
}
export function computeShadowLeaseParityState(mergeRequestState: string | null): {
shadowExecutorLeaseApplied: boolean;
shadowMergeLockApplied: boolean;
shadowLeaseApplied: boolean;
} {
const shadowExecutorLeaseApplied = false;
const shadowMergeLockApplied = mergeRequestState !== null
&& mergeRequestState !== "succeeded"
&& mergeRequestState !== "cancelled"
&& mergeRequestState !== "exhausted"
&& mergeRequestState !== "manual-required";
return {
shadowExecutorLeaseApplied,
shadowMergeLockApplied,
shadowLeaseApplied: shadowExecutorLeaseApplied || shadowMergeLockApplied,
};
}
export interface SchedulingDependencyParityDiff {
taskId: string;
dependencyId: string;
legacySatisfied: boolean;
markerSatisfied: boolean;
}
export function getUnmetSchedulingDependencies(
task: Task,
tasks: Task[],
options?: {
markerAcceptedByTaskId?: Map<string, boolean>;
onParityDiff?: (diff: SchedulingDependencyParityDiff) => void;
},
): string[] {
return task.dependencies.filter((depId) => {
const dep = tasks.find((candidate) => candidate.id === depId);
return dep && dep.column !== "done" && dep.column !== "in-review" && dep.column !== "archived";
if (!dep) return false;
const legacySatisfied = isLegacyDependencySatisfied(dep);
const markerSatisfied = isMarkerDependencySatisfied(dep, options?.markerAcceptedByTaskId?.get(depId) === true);
if (options?.onParityDiff && legacySatisfied !== markerSatisfied) {
options.onParityDiff({
taskId: task.id,
dependencyId: depId,
legacySatisfied,
markerSatisfied,
});
}
return !legacySatisfied;
});
}
@@ -534,16 +587,26 @@ export class Scheduler {
if (!settings.globalPause && !settings.enginePaused) {
const todoTasks = await this.store.listTasks({ column: "todo", slim: true });
const allTasks = await this.store.listTasks({ slim: true, includeArchived: true });
const taskById = new Map(allTasks.map((candidate) => [candidate.id, candidate]));
for (const dependent of todoTasks) {
const mentionsCompletedTask = dependent.dependencies.includes(task.id);
const currentlyBlockedByCompletedTask = dependent.blockedBy === task.id;
if (!mentionsCompletedTask && !currentlyBlockedByCompletedTask) continue;
const unresolvedDeps = dependent.dependencies.filter((depId) => {
const dep = taskById.get(depId);
return dep && dep.column !== "done" && dep.column !== "in-review" && dep.column !== "archived";
});
const markerAcceptedByTaskId = settings.mergeRequestContractShadowEnabled === true
? new Map(dependent.dependencies.map((depId) => [depId, this.store.getCompletionHandoffAcceptedMarker(depId) !== null]))
: undefined;
const unresolvedDeps = getUnmetSchedulingDependencies(
dependent,
[dependent, ...allTasks],
markerAcceptedByTaskId
? {
markerAcceptedByTaskId,
onParityDiff: (diff) => {
this.emitDependencyParityDiff(diff);
},
}
: undefined,
);
try {
if (unresolvedDeps.length > 0) {
@@ -655,17 +718,27 @@ export class Scheduler {
const inProgressTasks = await this.store.listTasks({ column: "in-progress", slim: true });
const dependents = [...todoTasks, ...inProgressTasks];
const allTasks = await this.store.listTasks({ slim: true, includeArchived: true });
const taskById = new Map(allTasks.map((candidate) => [candidate.id, candidate]));
for (const dependent of dependents) {
const mentionsDeletedTask = dependent.dependencies.includes(task.id);
const currentlyBlockedByDeletedTask = dependent.blockedBy === task.id;
if (!mentionsDeletedTask && !currentlyBlockedByDeletedTask) continue;
const unresolvedDeps = dependent.dependencies.filter((depId) => {
const dep = taskById.get(depId);
return dep && dep.column !== "done" && dep.column !== "in-review" && dep.column !== "archived";
});
const markerAcceptedByTaskId = settings.mergeRequestContractShadowEnabled === true
? new Map(dependent.dependencies.map((depId) => [depId, this.store.getCompletionHandoffAcceptedMarker(depId) !== null]))
: undefined;
const unresolvedDeps = getUnmetSchedulingDependencies(
dependent,
[dependent, ...allTasks],
markerAcceptedByTaskId
? {
markerAcceptedByTaskId,
onParityDiff: (diff) => {
this.emitDependencyParityDiff(diff);
},
}
: undefined,
);
try {
if (unresolvedDeps.length > 0) {
@@ -800,6 +873,22 @@ export class Scheduler {
return true;
}
private emitDependencyParityDiff(diff: SchedulingDependencyParityDiff): void {
void this.store.recordRunAuditEvent?.({
taskId: diff.taskId,
agentId: "scheduler",
runId: generateSyntheticRunId("scheduler", diff.taskId),
domain: "database",
mutationType: "merge:dependency-parity-diff",
target: diff.dependencyId,
metadata: {
depId: diff.dependencyId,
legacyResult: diff.legacySatisfied,
markerResult: diff.markerSatisfied,
},
});
}
private async emitDispatchQueuedConcurrencyAudit(task: Task, diagnostic: ConcurrencyGateDiagnostic): Promise<void> {
try {
await this.store.recordRunAuditEvent?.({
@@ -1206,7 +1295,33 @@ export class Scheduler {
for (const t of inReviewWithWorktree) {
const filteredScope = await getFilteredFileScope(t.id);
if (isCoordinationOnlyTask(t, filteredScope)) continue;
if (filteredScope.length > 0) setActiveScopeLease(t.id, filteredScope, "in-review");
if (filteredScope.length === 0) continue;
setActiveScopeLease(t.id, filteredScope, "in-review");
if (settings.mergeRequestContractShadowEnabled === true) {
const mergeRequestRecord = this.store.getMergeRequestRecord(t.id);
const { shadowExecutorLeaseApplied, shadowMergeLockApplied, shadowLeaseApplied } =
computeShadowLeaseParityState(mergeRequestRecord?.state ?? null);
if (shadowLeaseApplied !== true) {
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: true,
shadowLeaseApplied,
shadowExecutorLeaseApplied,
shadowMergeLockApplied,
mergeRequestState: mergeRequestRecord?.state ?? null,
},
});
}
}
}
for (const t of todo) {
@@ -1226,6 +1341,14 @@ export class Scheduler {
// Resolve dependency order among todo tasks
const ordered = resolveDependencyOrder(todo);
const mergeShadowEnabled = settings.mergeRequestContractShadowEnabled === true;
const markerAcceptedByTaskId = new Map<string, boolean>();
if (mergeShadowEnabled) {
const dependencyIds = new Set(todo.flatMap((candidate) => candidate.dependencies));
for (const depId of dependencyIds) {
markerAcceptedByTaskId.set(depId, this.store.getCompletionHandoffAcceptedMarker(depId) !== null);
}
}
let started = 0;
let loggedMissingAgentStoreThisPass = false;
@@ -1247,7 +1370,14 @@ export class Scheduler {
}
// Check all deps are satisfied (done, in-review, or archived)
const unmetDeps = getUnmetSchedulingDependencies(task, tasks);
const unmetDeps = getUnmetSchedulingDependencies(task, tasks, mergeShadowEnabled
? {
markerAcceptedByTaskId,
onParityDiff: (diff) => {
this.emitDependencyParityDiff(diff);
},
}
: undefined);
if (unmetDeps.length > 0) {
await this.store.updateTask(task.id, {