FN-8785: deduplicate queued dependency and scope logs

Persist queue episodes atomically so repeated scheduler and self-healing passes do not duplicate diagnostics.

- Add a queued-episode signature with PostgreSQL migration and task serialization support.
- Route dependency and file-scope queue transitions through the atomic deduplication API.
- Cover repeated and concurrent queue transitions, and update scheduler mocks for the new store API.

Files changed:
 .changeset/fn-8785-queued-log-deduplication.md     |   7 +
 docs/architecture.md                               |   1 +
 .../postgres/queued-episode-transition.pg.test.ts  | 158 +++++++++++++++++++++
 .../src/__tests__/postgres/schema-applier.test.ts  |   9 +-
 .../core/src/postgres/migrations/0000_initial.sql  |   1 +
 .../0044_fn_8785_queued_episode_signature.sql      |   3 +
 packages/core/src/postgres/schema-applier.ts       |  12 +-
 packages/core/src/postgres/schema/project.ts       |   1 +
 packages/core/src/store.ts                         |   5 +-
 packages/core/src/task-store/audit-ops.ts          |  80 +++++++++++
 packages/core/src/task-store/persistence.ts        |   2 +
 packages/core/src/task-store/serialization.ts      |   1 +
 packages/core/src/types/task/task-core.ts          |   5 +
 ...executor-outer-dispatch-dependency-gate.test.ts |  39 +++--
 .../engine/src/__tests__/executor-test-helpers.ts  |   7 +
 .../__tests__/scheduler-overlap-starvation.test.ts |  43 +++++-
 .../__tests__/scheduler-workflow-cutover.test.ts   |  35 +++--
 .../self-healing-completion-fanout.test.ts         |  37 +++++
 packages/engine/src/__tests__/self-healing.test.ts |  25 +++-
 packages/engine/src/executor.ts                    |  16 ++-
 packages/engine/src/scheduler.ts                   |  26 ++--
 packages/engine/src/self-healing.ts                | 117 +++++----------
 22 files changed, 491 insertions(+), 139 deletions(-)

Fusion-Task-Id: FN-8785
Fusion-Task-Lineage: 8d68c243-f24e-4db4-bd1d-b7eda01309f6
Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-04 11:52:15 -07:00
parent 750b4bedbc
commit 1dc636b8b4
22 changed files with 491 additions and 139 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Prevent repeated dependency and file-scope queue activity entries for unchanged blockers.
category: fix
dev: Queue episode signatures are durable and project/task transaction-serialized across scheduler, executor, and recovery producers.

View File

@@ -635,6 +635,7 @@ See [Memory Plugin Contract](./memory-plugin-contract.md) for the full plan.
- `blockedBy` invariant (FN-3924/FN-4091): the field is only durable when it references a current unresolved explicit dependency (or, for dependency-free tasks, an active overlap blocker). Completion gating now validates `blockedBy` through live task resolution: missing blockers and blockers already in `done`/`archived` are treated as stale, while only still-active blockers continue to prevent `fn_task_done`. If no current blocker remains, scheduler/event reconciliation clears `blockedBy` to `null` and re-evaluates from live task state. - `blockedBy` invariant (FN-3924/FN-4091): the field is only durable when it references a current unresolved explicit dependency (or, for dependency-free tasks, an active overlap blocker). Completion gating now validates `blockedBy` through live task resolution: missing blockers and blockers already in `done`/`archived` are treated as stale, while only still-active blockers continue to prevent `fn_task_done`. If no current blocker remains, scheduler/event reconciliation clears `blockedBy` to `null` and re-evaluates from live task state.
- Dependency-cycle invariant (FN-5256): task dependency graphs are acyclic at write time (`DependencyCycleError` in `TaskStore` for `createTask`, `createTaskWithReservedId`, `updateTask`, and `applyReplicatedTaskCreate`) with `task:dependency-cycle-rejected` audit evidence. Self-healing batch 2 adds `reconcileDependencyCycles`, which emits `task:dependency-cycle-detected`, auto-repairs only bounded umbrella-back-edge loops via `task:auto-reconciled-dependency-cycle`, and leaves ambiguous cycles untouched with `task:dependency-cycle-unrepaired` for operator inspection. - Dependency-cycle invariant (FN-5256): task dependency graphs are acyclic at write time (`DependencyCycleError` in `TaskStore` for `createTask`, `createTaskWithReservedId`, `updateTask`, and `applyReplicatedTaskCreate`) with `task:dependency-cycle-rejected` audit evidence. Self-healing batch 2 adds `reconcileDependencyCycles`, which emits `task:dependency-cycle-detected`, auto-repairs only bounded umbrella-back-edge loops via `task:auto-reconciled-dependency-cycle`, and leaves ambiguous cycles untouched with `task:dependency-cycle-unrepaired` for operator inspection.
- Dependency-blocking lease invariant (FN-6292): an `in-progress` task with unmet scheduling dependencies must not contribute an active file-scope lease in scheduler lease maps. This prevents a holder from queueing its own dependency behind its lease and creating a circular wait. - Dependency-blocking lease invariant (FN-6292): an `in-progress` task with unmet scheduling dependencies must not contribute an active file-scope lease in scheduler lease maps. This prevents a holder from queueing its own dependency behind its lease and creating a circular wait.
- **Queued blocker activity (FN-8785):** scheduler, executor, and self-healing use one project/task advisory-locked PostgreSQL transaction to persist queue fields, a durable full episode signature, and its task-log entry. The signature includes either the sorted unique full unmet dependency set or the overlap holder, so unchanged queued state logs once across producers and restarts; a changed kind/set, recovery/non-queued state, or a later re-block re-arms one observable entry.
- **Mission symbol admission (FN-8306):** autonomous mission implementation evaluates `evaluateMissionLineageApproval` (active Mission/Milestone/Slice, triaged or in-progress Feature, and required plan fingerprint). Approved work with durable declared symbols atomically acquires project-scoped locks after ordinary capacity, checkout, and dependency gates; same symbols queue with holder diagnostics while disjoint symbols can run in parallel. Mission-linked work lacking approved lineage is `lineage-blocked` and consumes neither symbol nor work lease. Non-mission work and approved mission work without resolvable symbols retain coarse file-scope serialization. Active scheduler heartbeats and long-running workflow processors renew their short crash-recoverable leases before expiry. Central `moveTask` exits from implementation release declared-symbol locks for review, cancellation/requeue, and terminal transitions; self-healing only expires terminal/missing/expired owners as a backstop. - **Mission symbol admission (FN-8306):** autonomous mission implementation evaluates `evaluateMissionLineageApproval` (active Mission/Milestone/Slice, triaged or in-progress Feature, and required plan fingerprint). Approved work with durable declared symbols atomically acquires project-scoped locks after ordinary capacity, checkout, and dependency gates; same symbols queue with holder diagnostics while disjoint symbols can run in parallel. Mission-linked work lacking approved lineage is `lineage-blocked` and consumes neither symbol nor work lease. Non-mission work and approved mission work without resolvable symbols retain coarse file-scope serialization. Active scheduler heartbeats and long-running workflow processors renew their short crash-recoverable leases before expiry. Central `moveTask` exits from implementation release declared-symbol locks for review, cancellation/requeue, and terminal transitions; self-healing only expires terminal/missing/expired owners as a backstop.
#### BlockedBy stamping invariants #### BlockedBy stamping invariants

View File

@@ -0,0 +1,158 @@
import { afterAll, afterEach, beforeAll, expect, it } from "vitest";
import { and, eq, sql } from "drizzle-orm";
import { mkdtemp, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
type SharedPgTaskStoreHarness,
} from "../../__test-utils__/pg-test-harness.js";
import * as schema from "../../postgres/schema/index.js";
import type { ResolvedBackend } from "../../postgres/backend-resolver.js";
import { createConnectionSetFromUrl } from "../../postgres/connection.js";
import { createAsyncDataLayer } from "../../postgres/data-layer.js";
import { TaskStore } from "../../store.js";
import { insertTaskRow } from "../../task-store/async/async-persistence.js";
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_queued_episode" });
const dependency = (ids: string[]) => ({
signature: `dependency:${[...new Set(ids)].sort().join(",")}`,
blockedBy: ids[0] ?? null,
overlapBlockedBy: null,
action: `queued — unmet dependencies: ${ids.join(", ")}`,
});
const overlap = (id: string) => ({
signature: `file-scope:${id}`,
blockedBy: null,
overlapBlockedBy: id,
action: `queued — blocked by active file-scope lease ${id}`,
});
pgDescribe("queued episode transition", () => {
beforeAll(h.beforeAll);
afterAll(h.afterAll);
afterEach(h.afterEach);
it("logs once for an unchanged full dependency episode and again when its full set changes", async () => {
const task = await h.store().createTask({ description: "queued task" });
expect((await h.store().transitionQueuedEpisode(task.id, dependency(["FN-A", "FN-B"]))).appended).toBe(true);
expect((await h.store().transitionQueuedEpisode(task.id, dependency(["FN-A", "FN-B"]))).appended).toBe(false);
expect((await h.store().transitionQueuedEpisode(task.id, dependency(["FN-A", "FN-C"]))).appended).toBe(true);
const updated = await h.store().getTask(task.id);
expect(updated.blockedBy).toBe("FN-A");
expect(updated.log?.filter((entry) => entry.action.startsWith("queued — unmet dependencies"))).toHaveLength(2);
});
it("switches queue kind and re-arms after a cleared state", async () => {
const task = await h.store().createTask({ description: "queue boundaries" });
await h.store().transitionQueuedEpisode(task.id, dependency(["FN-A"]));
await h.store().transitionQueuedEpisode(task.id, overlap("FN-LOCK"));
await h.store().updateTask(task.id, { status: null, blockedBy: null, overlapBlockedBy: null });
expect((await h.store().transitionQueuedEpisode(task.id, overlap("FN-LOCK"))).appended).toBe(true);
const updated = await h.store().getTask(task.id);
expect(updated.log?.filter((entry) => entry.action.startsWith("queued —"))).toHaveLength(3);
});
it("serializes separate-store producers and keeps the committed episode suppressed after reconstruction", async () => {
const task = await h.store().createTask({ description: "concurrent queue" });
const backend: ResolvedBackend = {
mode: "external",
runtimeUrl: h.testUrl(),
migrationUrl: h.testUrl(),
migrationUrlOverridden: true,
directSessionUrl: h.testUrl(),
directSessionProvenance: "migration-override",
};
const [connections, root] = await Promise.all([
createConnectionSetFromUrl(backend, { projectId: "", useRuntimeRole: true }),
mkdtemp(join(tmpdir(), "fusion-queued-episode-reconstructed-")),
]);
try {
const reconstructed = new TaskStore(root, undefined, {
asyncLayer: createAsyncDataLayer(connections, { projectId: "" }),
});
const results = await Promise.all(Array.from(
{ length: 8 },
(_, index) => (index % 2 === 0 ? h.store() : reconstructed).transitionQueuedEpisode(task.id, overlap("FN-LOCK")),
));
expect(results.filter((result) => result.appended)).toHaveLength(1);
expect((await reconstructed.transitionQueuedEpisode(task.id, overlap("FN-LOCK"))).appended).toBe(false);
const row = (await h.adminDb().select().from(schema.project.tasks).where(and(
eq(schema.project.tasks.projectId, "__legacy_unscoped__"),
eq(schema.project.tasks.id, task.id),
)))[0];
expect(row?.queuedLogEpisodeSignature).toBe("file-scope:FN-LOCK");
expect((await reconstructed.getTask(task.id))?.log?.filter((entry) => entry.action.startsWith("queued —"))).toHaveLength(1);
} finally {
await Promise.allSettled([connections.close(), rm(root, { recursive: true, force: true })]);
}
});
it("keeps same task IDs isolated between projects", async () => {
const backend: ResolvedBackend = {
mode: "external",
runtimeUrl: h.testUrl(),
migrationUrl: h.testUrl(),
migrationUrlOverridden: true,
directSessionUrl: h.testUrl(),
directSessionProvenance: "migration-override",
};
const [connectionsA, connectionsB, rootA, rootB] = await Promise.all([
createConnectionSetFromUrl(backend, { projectId: "project-a", useRuntimeRole: true }),
createConnectionSetFromUrl(backend, { projectId: "project-b", useRuntimeRole: true }),
mkdtemp(join(tmpdir(), "fusion-queued-episode-a-")),
mkdtemp(join(tmpdir(), "fusion-queued-episode-b-")),
]);
try {
const layerA = createAsyncDataLayer(connectionsA, { projectId: "project-a" });
const layerB = createAsyncDataLayer(connectionsB, { projectId: "project-b" });
const storeA = new TaskStore(rootA, undefined, { asyncLayer: layerA });
const storeB = new TaskStore(rootB, undefined, { asyncLayer: layerB });
const now = new Date().toISOString();
const row = { id: "FN-SAME", description: "same id", column: "todo", currentStep: 0, createdAt: now, updatedAt: now };
await Promise.all([insertTaskRow(layerA, row, { lineageId: null }), insertTaskRow(layerB, row, { lineageId: null })]);
await expect(Promise.all([
storeA.transitionQueuedEpisode(row.id, overlap("FN-LOCK")),
storeB.transitionQueuedEpisode(row.id, overlap("FN-LOCK")),
])).resolves.toEqual(expect.arrayContaining([
expect.objectContaining({ appended: true }),
expect.objectContaining({ appended: true }),
]));
expect((await storeA.getTask(row.id))?.log?.filter((entry) => entry.action.startsWith("queued —"))).toHaveLength(1);
expect((await storeB.getTask(row.id))?.log?.filter((entry) => entry.action.startsWith("queued —"))).toHaveLength(1);
} finally {
await Promise.allSettled([
connectionsA.close(), connectionsB.close(), rm(rootA, { recursive: true, force: true }), rm(rootB, { recursive: true, force: true }),
]);
}
});
it("rolls back queue fields, marker, and log together when persistence fails", async () => {
const task = await h.store().createTask({ description: "rollback queue transition" });
await h.adminDb().execute(sql.raw(`
CREATE FUNCTION public.fail_queued_episode_transition_for_test() RETURNS trigger
LANGUAGE plpgsql AS $$ BEGIN
IF NEW.id = '${task.id}' THEN RAISE EXCEPTION 'forced queued episode failure'; END IF;
RETURN NEW;
END $$;
CREATE TRIGGER fail_queued_episode_transition_for_test
BEFORE UPDATE ON project.tasks FOR EACH ROW EXECUTE FUNCTION public.fail_queued_episode_transition_for_test();
`));
try {
await expect(h.store().transitionQueuedEpisode(task.id, overlap("FN-LOCK"))).rejects.toThrow("Failed query");
const rolledBack = await h.store().getTask(task.id);
expect(rolledBack?.status).not.toBe("queued");
expect(rolledBack?.queuedLogEpisodeSignature).toBeUndefined();
expect(rolledBack?.log?.filter((entry) => entry.action.startsWith("queued —"))).toHaveLength(0);
} finally {
await h.adminDb().execute(sql.raw(`
DROP TRIGGER fail_queued_episode_transition_for_test ON project.tasks;
DROP FUNCTION public.fail_queued_episode_transition_for_test();
`));
}
expect((await h.store().transitionQueuedEpisode(task.id, overlap("FN-LOCK"))).appended).toBe(true);
});
});

View File

@@ -86,6 +86,7 @@ import {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
} from "../../postgres/schema-applier.js"; } from "../../postgres/schema-applier.js";
import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js"; import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js";
import type { PluginSchemaInitHook } from "../../postgres/plugin-schema-hook.js"; import type { PluginSchemaInitHook } from "../../postgres/plugin-schema-hook.js";
@@ -105,7 +106,8 @@ describe("schema-applier: immutable migration identities", () => {
expect(TASK_LIFECYCLE_CONSUMERS_VERSION).toBe("0041"); expect(TASK_LIFECYCLE_CONSUMERS_VERSION).toBe("0041");
expect(VALIDATOR_INPUT_FINGERPRINT_VERSION).toBe("0042"); expect(VALIDATOR_INPUT_FINGERPRINT_VERSION).toBe("0042");
expect(UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION).toBe("0043"); expect(UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION).toBe("0043");
expect(SCHEMA_BASELINE_VERSION).toBe("0043"); expect(QUEUED_EPISODE_SIGNATURE_VERSION).toBe("0044");
expect(SCHEMA_BASELINE_VERSION).toBe("0044");
}); });
it("keeps monitor and approval isolation assigned to version 0003", () => { it("keeps monitor and approval isolation assigned to version 0003", () => {
@@ -1761,6 +1763,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
]); ]);
expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false); expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false);
}); });
@@ -1830,6 +1833,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
]); ]);
}); });
@@ -2032,6 +2036,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
]); ]);
}); });
@@ -2115,6 +2120,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
]); ]);
}); });
@@ -2198,6 +2204,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
TASK_LIFECYCLE_CONSUMERS_VERSION, TASK_LIFECYCLE_CONSUMERS_VERSION,
VALIDATOR_INPUT_FINGERPRINT_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION,
UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION,
QUEUED_EPISODE_SIGNATURE_VERSION,
]); ]);
}); });
}); });

View File

@@ -52,6 +52,7 @@ CREATE TABLE IF NOT EXISTS project.tasks (
worktree text, worktree text,
blocked_by text, blocked_by text,
overlap_blocked_by text, overlap_blocked_by text,
queued_log_episode_signature text,
paused integer DEFAULT 0, paused integer DEFAULT 0,
user_paused integer DEFAULT 0, user_paused integer DEFAULT 0,
paused_reason text, paused_reason text,

View File

@@ -0,0 +1,3 @@
-- FNXC:QueuedTaskLogging 2026-08-04-18:03: retain one durable full blocker signature per live task queue episode.
ALTER TABLE project.tasks
ADD COLUMN IF NOT EXISTS queued_log_episode_signature text;

View File

@@ -56,7 +56,7 @@ capacity-model table drop that landed while this PR was open.
*/ */
/* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: advance the schema ceiling so durable consumer state exists before observers begin polling FN-8684's outbox. */ /* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: advance the schema ceiling so durable consumer state exists before observers begin polling FN-8684's outbox. */
/* FNXC:MissionValidation 2026-08-01-16:21: advance the schema ceiling before validator admission reads durable content fingerprints. */ /* FNXC:MissionValidation 2026-08-01-16:21: advance the schema ceiling before validator admission reads durable content fingerprints. */
export const SCHEMA_BASELINE_VERSION = "0043"; export const SCHEMA_BASELINE_VERSION = "0044";
/** FNXC:SymbolLock 2026-07-20-10:00: upgrades need durable task declarations before admission resolves symbols. */ /** FNXC:SymbolLock 2026-07-20-10:00: upgrades need durable task declarations before admission resolves symbols. */
export const TASK_DECLARED_SYMBOLS_VERSION = "0028"; export const TASK_DECLARED_SYMBOLS_VERSION = "0028";
const INITIAL_SCHEMA_VERSION = "0000"; const INITIAL_SCHEMA_VERSION = "0000";
@@ -185,6 +185,8 @@ export const TASK_LIFECYCLE_CONSUMERS_VERSION = "0041";
export const VALIDATOR_INPUT_FINGERPRINT_VERSION = "0042"; export const VALIDATOR_INPUT_FINGERPRINT_VERSION = "0042";
/** FNXC:PlanningDependencyReseed 2026-08-04-02:14: durable per-episode unplanned-dispatch diagnostics. */ /** FNXC:PlanningDependencyReseed 2026-08-04-02:14: durable per-episode unplanned-dispatch diagnostics. */
export const UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION = "0043"; export const UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION = "0043";
/** FNXC:QueuedTaskLogging 2026-08-04-18:03: upgraded databases need the durable full queue-episode signature before concurrent producers can suppress repeats safely. */
export const QUEUED_EPISODE_SIGNATURE_VERSION = "0044";
/** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */ /** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */
export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained"; export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained";
@@ -402,6 +404,7 @@ const TASK_LIFECYCLE_OUTBOX_MIGRATION_PATH = join(MIGRATIONS_DIR, "0040_fn_8684_
const TASK_LIFECYCLE_CONSUMERS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0041_fn_8685_task_lifecycle_consumers.sql"); const TASK_LIFECYCLE_CONSUMERS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0041_fn_8685_task_lifecycle_consumers.sql");
const VALIDATOR_INPUT_FINGERPRINT_MIGRATION_PATH = join(MIGRATIONS_DIR, "0042_fn_8694_validator_input_fingerprint.sql"); const VALIDATOR_INPUT_FINGERPRINT_MIGRATION_PATH = join(MIGRATIONS_DIR, "0042_fn_8694_validator_input_fingerprint.sql");
const UNPLANNED_EXECUTION_BLOCK_DEDUPE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0043_fn8768_dispatch_dedupe.sql"); const UNPLANNED_EXECUTION_BLOCK_DEDUPE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0043_fn8768_dispatch_dedupe.sql");
const QUEUED_EPISODE_SIGNATURE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0044_fn_8785_queued_episode_signature.sql");
/** /**
* Ensure the migration bookkeeping table exists. Lives in the public schema so * Ensure the migration bookkeeping table exists. Lives in the public schema so
@@ -515,6 +518,7 @@ export async function applySchemaBaseline(
const taskLifecycleConsumersAlreadyApplied = applied.includes(TASK_LIFECYCLE_CONSUMERS_VERSION); const taskLifecycleConsumersAlreadyApplied = applied.includes(TASK_LIFECYCLE_CONSUMERS_VERSION);
const validatorInputFingerprintAlreadyApplied = applied.includes(VALIDATOR_INPUT_FINGERPRINT_VERSION); const validatorInputFingerprintAlreadyApplied = applied.includes(VALIDATOR_INPUT_FINGERPRINT_VERSION);
const unplannedExecutionBlockDedupeAlreadyApplied = applied.includes(UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION); const unplannedExecutionBlockDedupeAlreadyApplied = applied.includes(UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION);
const queuedEpisodeSignatureAlreadyApplied = applied.includes(QUEUED_EPISODE_SIGNATURE_VERSION);
assertBinaryNotOlderThanDatabase(applied); assertBinaryNotOlderThanDatabase(applied);
let schemaChanged = false; let schemaChanged = false;
@@ -1100,6 +1104,12 @@ export async function applySchemaBaseline(
await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION}) ON CONFLICT (version) DO NOTHING`); await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION}) ON CONFLICT (version) DO NOTHING`);
schemaChanged = true; schemaChanged = true;
} }
if (!queuedEpisodeSignatureAlreadyApplied) {
const migrationSql = await readFile(QUEUED_EPISODE_SIGNATURE_MIGRATION_PATH, "utf8");
await tx.execute(sql.raw(migrationSql));
await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${QUEUED_EPISODE_SIGNATURE_VERSION}) ON CONFLICT (version) DO NOTHING`);
schemaChanged = true;
}
return { applied: schemaChanged, pluginHooksRun: pluginHooks.length }; return { applied: schemaChanged, pluginHooksRun: pluginHooks.length };
}); });

View File

@@ -80,6 +80,7 @@ export const tasks = projectSchema.table("tasks", {
worktree: text("worktree"), worktree: text("worktree"),
blockedBy: text("blocked_by"), blockedBy: text("blocked_by"),
overlapBlockedBy: text("overlap_blocked_by"), overlapBlockedBy: text("overlap_blocked_by"),
queuedLogEpisodeSignature: text("queued_log_episode_signature"),
paused: integer("paused").default(0), paused: integer("paused").default(0),
userPaused: integer("user_paused").default(0), userPaused: integer("user_paused").default(0),
pausedReason: text("paused_reason"), pausedReason: text("paused_reason"),

View File

@@ -115,7 +115,7 @@ import { queryRunAuditEvents } from "./task-store/async/async-audit.js";
import { isValidMergeRequestTransitionImpl, releaseMergeQueueLeaseImpl, collectMergeDetailsImpl, applyPrMergedTransitionImpl } from "./task-store/merge-queue-ops-2.js"; import { isValidMergeRequestTransitionImpl, releaseMergeQueueLeaseImpl, collectMergeDetailsImpl, applyPrMergedTransitionImpl } from "./task-store/merge-queue-ops-2.js";
import { upsertWorkflowWorkItemImpl, replaceActiveTaskWorkflowContinuationImpl, seedStrandedPlanReviewContinuationImpl, transitionWorkflowWorkItemImpl, acquireWorkflowWorkItemLeaseImpl } from "./task-store/workflow-workitems-ops-2.js"; import { upsertWorkflowWorkItemImpl, replaceActiveTaskWorkflowContinuationImpl, seedStrandedPlanReviewContinuationImpl, transitionWorkflowWorkItemImpl, acquireWorkflowWorkItemLeaseImpl } from "./task-store/workflow-workitems-ops-2.js";
import { getSettingsImpl, getSettingsFastImpl, getSettingsByScopeImpl, getSettingsByScopeFastImpl } from "./task-store/settings-ops-2.js"; import { getSettingsImpl, getSettingsFastImpl, getSettingsByScopeImpl, getSettingsByScopeFastImpl } from "./task-store/settings-ops-2.js";
import { runPluginColumnTransitionHooksImpl, checkAndRecordUnplannedExecutionBlockImpl, logEntryImpl } from "./task-store/audit-ops.js"; import { runPluginColumnTransitionHooksImpl, checkAndRecordUnplannedExecutionBlockImpl, logEntryImpl, transitionQueuedEpisodeImpl, type QueuedEpisodeTransition } from "./task-store/audit-ops.js";
import { clearWorkflowRunBranchesImpl, projectMergeRequestToWorkflowWorkItemImpl, createCompletionHandoffWorkflowWorkImpl } from "./task-store/workflow-workitems-ops.js"; import { clearWorkflowRunBranchesImpl, projectMergeRequestToWorkflowWorkItemImpl, createCompletionHandoffWorkflowWorkImpl } from "./task-store/workflow-workitems-ops.js";
import { flushAgentLogBufferImpl, appendAgentLogBatchImpl } from "./task-store/agent-logs.js"; import { flushAgentLogBufferImpl, appendAgentLogBatchImpl } from "./task-store/agent-logs.js";
import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js"; import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js";
@@ -1712,6 +1712,9 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
async checkAndRecordUnplannedExecutionBlock(id: string, episode: string): Promise<boolean> { async checkAndRecordUnplannedExecutionBlock(id: string, episode: string): Promise<boolean> {
return checkAndRecordUnplannedExecutionBlockImpl(this, id, episode); return checkAndRecordUnplannedExecutionBlockImpl(this, id, episode);
} }
async transitionQueuedEpisode(id: string, transition: QueuedEpisodeTransition): Promise<import("./task-store/audit-ops.js").QueuedEpisodeTransitionResult> {
return transitionQueuedEpisodeImpl(this, id, transition);
}
async logEntry(id: string, action: string, outcome?: string, runContext?: RunMutationContext): Promise<Task> { async logEntry(id: string, action: string, outcome?: string, runContext?: RunMutationContext): Promise<Task> {
return logEntryImpl(this, id, action, outcome, runContext); return logEntryImpl(this, id, action, outcome, runContext);
} }

View File

@@ -18,6 +18,7 @@ import "../builtin-traits.js";
import {__setTaskActivityLogLimitsForTesting, truncateTaskLogOutcome, getTaskActivityLogEntryLimit} from "../task-store/comments.js"; import {__setTaskActivityLogLimitsForTesting, truncateTaskLogOutcome, getTaskActivityLogEntryLimit} from "../task-store/comments.js";
import {readTaskRow, updateTaskColumns} from "../task-store/async/async-persistence.js"; import {readTaskRow, updateTaskColumns} from "../task-store/async/async-persistence.js";
import { getLiveTaskColumn } from "./async/async-comments-attachments.js"; import { getLiveTaskColumn } from "./async/async-comments-attachments.js";
import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js";
import { resolveArchivedLanes } from "../project-lane-vocabulary.js"; import { resolveArchivedLanes } from "../project-lane-vocabulary.js";
import * as schema from "../postgres/schema/index.js"; import * as schema from "../postgres/schema/index.js";
@@ -122,6 +123,85 @@ Release gates can be evaluated by multiple schedulers. Claim the project/task
episode and append its diagnostic in one transaction so a crash cannot leave a episode and append its diagnostic in one transaction so a crash cannot leave a
suppression marker without the operator-visible task-log entry. suppression marker without the operator-visible task-log entry.
*/ */
export interface QueuedEpisodeTransition {
/** Canonical complete blocker identity, e.g. dependency:FN-1,FN-2. */
signature: string;
blockedBy: string | null;
overlapBlockedBy: string | null;
action: string;
outcome?: string;
runContext?: RunMutationContext;
}
export interface QueuedEpisodeTransitionResult {
appended: boolean;
task: Task;
}
/*
FNXC:QueuedTaskLogging 2026-08-04-18:03:
Dependency and file-scope producers share this full-signature transition so queue activity is
edge-triggered across schedulers, executors, self-healing, and process restarts. Acquire the
project/task advisory transaction lock before reading or updating the row; atomically persist the
marker, queue fields, and sole log entry. A matching signature suppresses only an already queued
row with matching blocker fields, so recovery/non-queued state and any blocker-kind/full-set change
re-arm reporting. Do not call public TaskStore mutation methods in this transaction.
*/
export async function transitionQueuedEpisodeImpl(
store: TaskStore,
id: string,
transition: QueuedEpisodeTransition,
): Promise<QueuedEpisodeTransitionResult> {
const layer = store.asyncLayer!;
const projectId = layer.projectId?.trim() || "__legacy_unscoped__";
const now = new Date().toISOString();
const result = await layer.transactionImmediate(async (tx) => {
await acquireTaskAdvisoryXactLock(tx, projectId, id);
const rows = await tx.select().from(schema.project.tasks).where(and(
eq(schema.project.tasks.projectId, projectId),
eq(schema.project.tasks.id, id),
isNull(schema.project.tasks.deletedAt),
));
const current = rows[0];
if (!current) throw new Error(`Task ${id} not found or archived while queuing`);
const appended = !(
current.status === "queued"
&& (current.blockedBy ?? null) === transition.blockedBy
&& (current.overlapBlockedBy ?? null) === transition.overlapBlockedBy
&& (current.queuedLogEpisodeSignature ?? null) === transition.signature
);
const log = Array.isArray(current.log) ? [...current.log as TaskLogEntry[]] : [];
if (appended) {
log.push({
timestamp: now,
action: transition.action,
outcome: truncateTaskLogOutcome(transition.outcome),
...(transition.runContext ? { runContext: transition.runContext } : {}),
});
const limit = getTaskActivityLogEntryLimit();
if (log.length > limit) log.splice(0, log.length - limit);
}
const updated = await tx.update(schema.project.tasks).set({
status: "queued",
blockedBy: transition.blockedBy,
overlapBlockedBy: transition.overlapBlockedBy,
queuedLogEpisodeSignature: transition.signature,
...(appended ? { log } : {}),
updatedAt: now,
}).where(and(
eq(schema.project.tasks.projectId, projectId),
eq(schema.project.tasks.id, id),
)).returning();
return { appended, task: updated[0]! };
});
const task = store.rowToTask(store.pgRowToTaskRow(result.task as unknown as Record<string, unknown>));
await store.writeTaskJsonFile(store.taskDir(id), task);
if (store.isWatching) store.taskCache.set(id, { ...task });
store.emitTaskLifecycleEventSafely("task:updated", [task]);
return { appended: result.appended, task };
}
export async function checkAndRecordUnplannedExecutionBlockImpl( export async function checkAndRecordUnplannedExecutionBlockImpl(
store: TaskStore, store: TaskStore,
id: string, id: string,

View File

@@ -26,6 +26,7 @@ export interface TaskRow {
worktree: string | null; worktree: string | null;
blockedBy: string | null; blockedBy: string | null;
overlapBlockedBy: string | null; overlapBlockedBy: string | null;
queuedLogEpisodeSignature: string | null;
paused: number | null; paused: number | null;
pausedReason: string | null; pausedReason: string | null;
wedgeNotification: string | null; wedgeNotification: string | null;
@@ -247,6 +248,7 @@ export const TASK_COLUMN_DESCRIPTORS: TaskColumnDescriptor[] = [
defineTaskColumn("worktree", (task) => task.worktree ?? null), defineTaskColumn("worktree", (task) => task.worktree ?? null),
defineTaskColumn("blockedBy", (task) => task.blockedBy ?? null), defineTaskColumn("blockedBy", (task) => task.blockedBy ?? null),
defineTaskColumn("overlapBlockedBy", (task) => task.overlapBlockedBy ?? null), defineTaskColumn("overlapBlockedBy", (task) => task.overlapBlockedBy ?? null),
defineTaskColumn("queuedLogEpisodeSignature", (task) => task.queuedLogEpisodeSignature ?? null),
defineTaskColumn("paused", (task) => task.paused ? 1 : 0), defineTaskColumn("paused", (task) => task.paused ? 1 : 0),
defineTaskColumn("pausedReason", (task) => task.pausedReason ?? null), defineTaskColumn("pausedReason", (task) => task.pausedReason ?? null),
defineTaskColumn("wedgeNotification", (task) => toJsonNullable(task.wedgeNotification)), defineTaskColumn("wedgeNotification", (task) => toJsonNullable(task.wedgeNotification)),

View File

@@ -76,6 +76,7 @@ export function rowToTask(row: TaskRow): Task {
worktree: row.worktree || undefined, worktree: row.worktree || undefined,
blockedBy: row.blockedBy || undefined, blockedBy: row.blockedBy || undefined,
overlapBlockedBy: row.overlapBlockedBy || undefined, overlapBlockedBy: row.overlapBlockedBy || undefined,
queuedLogEpisodeSignature: row.queuedLogEpisodeSignature || undefined,
paused: row.paused ? true : undefined, paused: row.paused ? true : undefined,
pausedReason: row.pausedReason || undefined, pausedReason: row.pausedReason || undefined,
wedgeNotification: fromJson<Task["wedgeNotification"]>(row.wedgeNotification) ?? undefined, wedgeNotification: fromJson<Task["wedgeNotification"]>(row.wedgeNotification) ?? undefined,

View File

@@ -660,6 +660,11 @@ export interface Task {
* Cleared when the overlap resolves (the blocker task moves to done or its * Cleared when the overlap resolves (the blocker task moves to done or its
* scope no longer overlaps). */ * scope no longer overlaps). */
overlapBlockedBy?: string; overlapBlockedBy?: string;
/**
* Durable identity of the currently reported dependency/file-scope queue episode.
* Internal producers use it to make persisted queue activity edge-triggered.
*/
queuedLogEpisodeSignature?: string;
/** When true, all automated agent and scheduler interaction is suspended. */ /** When true, all automated agent and scheduler interaction is suspended. */
paused?: boolean; paused?: boolean;
/** When true, this task was explicitly moved back to todo by a user and should not auto-dispatch. */ /** When true, this task was explicitly moved back to todo by a user and should not auto-dispatch. */

View File

@@ -119,17 +119,14 @@ describe("executor outer dispatch dependency gate", () => {
preserveResumeState: true, preserveResumeState: true,
recoveryRehome: true, recoveryRehome: true,
})); }));
expect(store.updateTask).toHaveBeenCalledWith( expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({
child.id, signature: "dependency:FN-PARENT",
expect.objectContaining({ status: "queued", blockedBy: parent.id }), blockedBy: parent.id,
undefined, }));
); expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({
expect(store.logEntry).toHaveBeenCalledWith( action: expect.stringContaining("queued — unmet dependencies: FN-PARENT"),
child.id, outcome: expect.stringContaining("dependency gate blocked"),
expect.stringContaining("queued — unmet dependencies: FN-PARENT"), }));
expect.stringContaining("dependency gate blocked"),
undefined,
);
expect(graph).not.toHaveBeenCalled(); expect(graph).not.toHaveBeenCalled();
// FNXC:DependencyGating 2026-07-16-00:00: A dependency-gated outer return // FNXC:DependencyGating 2026-07-16-00:00: A dependency-gated outer return
// must drop the scheduler's reservation because no downstream owner can take it. // must drop the scheduler's reservation because no downstream owner can take it.
@@ -148,11 +145,10 @@ describe("executor outer dispatch dependency gate", () => {
await executor.execute(child); await executor.execute(child);
expect(store.updateTask).toHaveBeenCalledWith( expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({
child.id, signature: "dependency:FN-PARENT",
expect.objectContaining({ status: "queued", blockedBy: parent.id }), blockedBy: parent.id,
undefined, }));
);
expect(graph).not.toHaveBeenCalled(); expect(graph).not.toHaveBeenCalled();
}); });
@@ -167,7 +163,7 @@ describe("executor outer dispatch dependency gate", () => {
await executor.execute(child); await executor.execute(child);
expect(store.moveTask).not.toHaveBeenCalled(); expect(store.moveTask).not.toHaveBeenCalled();
expect(store.updateTask).not.toHaveBeenCalledWith(child.id, expect.objectContaining({ status: "queued" }), undefined); expect(store.transitionQueuedEpisode).not.toHaveBeenCalled();
expect(graph).toHaveBeenCalledWith(child, { alreadyClaimed: true }); expect(graph).toHaveBeenCalledWith(child, { alreadyClaimed: true });
}); });
@@ -196,11 +192,10 @@ describe("executor outer dispatch dependency gate", () => {
await executor.execute(child); await executor.execute(child);
expect(store.getCompletionHandoffAcceptedMarker).toHaveBeenCalledWith(parent.id); expect(store.getCompletionHandoffAcceptedMarker).toHaveBeenCalledWith(parent.id);
expect(store.updateTask).toHaveBeenCalledWith( expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({
child.id, signature: "dependency:FN-PARENT",
expect.objectContaining({ status: "queued", blockedBy: parent.id }), blockedBy: parent.id,
undefined, }));
);
expect(graph).not.toHaveBeenCalled(); expect(graph).not.toHaveBeenCalled();
}); });

View File

@@ -646,6 +646,13 @@ export function createMockStore() {
updatedAt: new Date().toISOString(), updatedAt: new Date().toISOString(),
})), })),
logEntry: vi.fn().mockResolvedValue(undefined), logEntry: vi.fn().mockResolvedValue(undefined),
transitionQueuedEpisode: vi.fn(async (id: string, transition: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string; outcome?: string }) => {
const prior = patches.get(id) ?? {};
const appended = !(prior.status === "queued" && prior.blockedBy === transition.blockedBy && prior.overlapBlockedBy === transition.overlapBlockedBy && prior.queuedLogEpisodeSignature === transition.signature);
await store.updateTask(id, { status: "queued", blockedBy: transition.blockedBy, overlapBlockedBy: transition.overlapBlockedBy, queuedLogEpisodeSignature: transition.signature });
if (appended) await store.logEntry(id, transition.action, transition.outcome);
return { appended, task: { id, ...patches.get(id) } };
}),
addTaskComment: vi.fn().mockResolvedValue(undefined), addTaskComment: vi.fn().mockResolvedValue(undefined),
parseStepsFromPrompt: vi.fn().mockResolvedValue([]), parseStepsFromPrompt: vi.fn().mockResolvedValue([]),
parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]), parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]),

View File

@@ -57,6 +57,7 @@ function createStore(tasks: Task[], scopes: Record<string, string[]>, settings:
if (task) Object.assign(task, patch); if (task) Object.assign(task, patch);
return task as Task; return task as Task;
}); });
const logEntry = vi.fn(async () => undefined);
const moveTask = vi.fn(async (id: string, column: Task["column"]) => { const moveTask = vi.fn(async (id: string, column: Task["column"]) => {
const task = tasks.find((candidate) => candidate.id === id); const task = tasks.find((candidate) => candidate.id === id);
if (task) task.column = column; if (task) task.column = column;
@@ -88,7 +89,22 @@ function createStore(tasks: Task[], scopes: Record<string, string[]>, settings:
moveTask, moveTask,
moveTaskIf, moveTaskIf,
getTask: vi.fn(async (id: string) => tasks.find((task) => task.id === id) ?? null), getTask: vi.fn(async (id: string) => tasks.find((task) => task.id === id) ?? null),
logEntry: vi.fn(async () => undefined), logEntry,
transitionQueuedEpisode: vi.fn(async (id: string, transition: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string }) => {
const task = tasks.find((candidate) => candidate.id === id)!;
const appended = !(
task.status === "queued"
&& (task.blockedBy ?? null) === transition.blockedBy
&& (task.overlapBlockedBy ?? null) === transition.overlapBlockedBy
&& task.queuedLogEpisodeSignature === transition.signature
);
await updateTask(id, transition.signature.startsWith("dependency:")
? { status: "queued", blockedBy: transition.blockedBy ?? undefined }
: { status: "queued", blockedBy: transition.blockedBy, overlapBlockedBy: transition.overlapBlockedBy });
task.queuedLogEpisodeSignature = transition.signature;
if (appended) await logEntry(id, transition.action);
return { appended, task };
}),
getRootDir: vi.fn(() => "/tmp/project"), getRootDir: vi.fn(() => "/tmp/project"),
getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"), getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"),
on: vi.fn(), on: vi.fn(),
@@ -567,6 +583,31 @@ describe("scheduler overlap starvation regression (FN-057)", () => {
expect(store.moveTask).not.toHaveBeenCalledWith("FN-900", "in-progress", expect.anything()); expect(store.moveTask).not.toHaveBeenCalledWith("FN-900", "in-progress", expect.anything());
}); });
it("persists one dependency queue log per full blocker signature across repeated scheduler passes", async () => {
const tasks = [
makeTask({ id: "FN-A", column: "todo" }),
makeTask({ id: "FN-B", column: "todo" }),
makeTask({ id: "FN-C", column: "todo" }),
makeTask({ id: "FN-QUEUED", column: "todo", dependencies: ["FN-A", "FN-B"] }),
];
const store = createStore(tasks, {});
const scheduler = new Scheduler(store);
(scheduler as any).running = true;
await scheduler.schedule();
await scheduler.schedule();
tasks.find((task) => task.id === "FN-QUEUED")!.dependencies = ["FN-A", "FN-C"];
await scheduler.schedule();
const queueLogs = (store.logEntry as ReturnType<typeof vi.fn>).mock.calls
.filter((call) => call[0] === "FN-QUEUED" && String(call[1]).startsWith("queued — unmet dependencies"));
expect(queueLogs).toHaveLength(2);
expect(queueLogs.map((call) => call[1])).toEqual([
"queued — unmet dependencies: FN-A, FN-B",
"queued — unmet dependencies: FN-A, FN-C",
]);
});
it("clears an absent overlap blocker only after confirming no current overlap remains", async () => { it("clears an absent overlap blocker only after confirming no current overlap remains", async () => {
const tasks = [ const tasks = [
makeTask({ id: "FN-901", column: "todo", status: "queued", priority: "normal", overlapBlockedBy: "FN-MISSING" }), makeTask({ id: "FN-901", column: "todo", status: "queued", priority: "normal", overlapBlockedBy: "FN-MISSING" }),

View File

@@ -53,6 +53,12 @@ function storeWith(
workflows: { selections?: Record<string, string>; definitions?: Record<string, WorkflowIr> } = {}, workflows: { selections?: Record<string, string>; definitions?: Record<string, WorkflowIr> } = {},
): TaskStore { ): TaskStore {
const byId = new Map(tasks.map((candidate) => [candidate.id, candidate])); const byId = new Map(tasks.map((candidate) => [candidate.id, candidate]));
const updateTask = vi.fn(async (id: string, patch: Partial<Task>) => {
const current = byId.get(id);
if (current) Object.assign(current, patch);
return current as Task;
});
const logEntry = vi.fn(async () => undefined);
return { return {
listTasks: vi.fn(async () => [...byId.values()]), listTasks: vi.fn(async () => [...byId.values()]),
getTask: vi.fn(async (id: string) => byId.get(id) ?? null), getTask: vi.fn(async (id: string) => byId.get(id) ?? null),
@@ -63,11 +69,7 @@ function storeWith(
...settings, ...settings,
})), })),
updateSettings: vi.fn(async (patch: Record<string, unknown>) => ({ ...settings, ...patch })), updateSettings: vi.fn(async (patch: Record<string, unknown>) => ({ ...settings, ...patch })),
updateTask: vi.fn(async (id: string, patch: Partial<Task>) => { updateTask,
const current = byId.get(id);
if (current) Object.assign(current, patch);
return current as Task;
}),
moveTask: vi.fn(async (id: string, column: Task["column"]) => { moveTask: vi.fn(async (id: string, column: Task["column"]) => {
const current = byId.get(id); const current = byId.get(id);
if (current) current.column = column; if (current) current.column = column;
@@ -80,7 +82,22 @@ function storeWith(
return { task: current, moved: true }; return { task: current, moved: true };
}), }),
parseFileScopeFromPrompt: vi.fn(async () => []), parseFileScopeFromPrompt: vi.fn(async () => []),
logEntry: vi.fn(async () => undefined), logEntry,
transitionQueuedEpisode: vi.fn(async (id: string, transition: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string }) => {
const current = byId.get(id)!;
const appended = !(current.status === "queued"
&& (current.blockedBy ?? null) === transition.blockedBy
&& (current.overlapBlockedBy ?? null) === transition.overlapBlockedBy
&& current.queuedLogEpisodeSignature === transition.signature);
await updateTask(id, {
status: "queued",
blockedBy: transition.blockedBy,
overlapBlockedBy: transition.overlapBlockedBy,
queuedLogEpisodeSignature: transition.signature,
});
if (appended) await logEntry(id, transition.action);
return { appended, task: current };
}),
getRootDir: vi.fn(() => "/tmp/project"), getRootDir: vi.fn(() => "/tmp/project"),
getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"), getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"),
on: vi.fn(), on: vi.fn(),
@@ -388,9 +405,11 @@ describe("Scheduler workflow cutover", () => {
await scheduler.schedule(); await scheduler.schedule();
expect(store.updateTask).toHaveBeenCalledWith("FN-002", { expect(store.transitionQueuedEpisode).toHaveBeenCalledWith("FN-002", {
status: "queued", signature: "dependency:FN-001",
blockedBy: "FN-001", blockedBy: "FN-001",
overlapBlockedBy: null,
action: "queued — unmet dependencies: FN-001",
}); });
expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-002", "in-progress", expect.anything(), expect.anything()); expect(store.moveTaskIf).not.toHaveBeenCalledWith("FN-002", "in-progress", expect.anything(), expect.anything());
expect(onBlocked).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002" }), ["FN-001"]); expect(onBlocked).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002" }), ["FN-001"]);

View File

@@ -70,6 +70,23 @@ function createStore(tasks: Task[], settings?: Partial<Settings>): TaskStore & E
map.set(id, { ...task, ...patch } as Task); map.set(id, { ...task, ...patch } as Task);
return map.get(id); return map.get(id);
}), }),
transitionQueuedEpisode: vi.fn(async (id: string, transition: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string }) => {
const task = map.get(id)!;
const appended = !(task.status === "queued"
&& (task.blockedBy ?? null) === transition.blockedBy
&& (task.overlapBlockedBy ?? null) === transition.overlapBlockedBy
&& task.queuedLogEpisodeSignature === transition.signature);
const updated = {
...task,
status: "queued",
blockedBy: transition.blockedBy,
overlapBlockedBy: transition.overlapBlockedBy,
queuedLogEpisodeSignature: transition.signature,
log: appended ? [...(task.log ?? []), { timestamp: new Date().toISOString(), action: transition.action }] : task.log,
} as Task;
map.set(id, updated);
return { appended, task: updated };
}),
moveTask: vi.fn(async (id: string, column: Task["column"]) => { moveTask: vi.fn(async (id: string, column: Task["column"]) => {
const task = map.get(id)!; const task = map.get(id)!;
const from = task.column; const from = task.column;
@@ -113,6 +130,26 @@ describe("self-healing completion fan-out", () => {
); );
}); });
it("deduplicates concurrent completion fanout that leaves a dependent behind the same queue episode", async () => {
const blocker = makeTask("FN-B", { column: "done" });
const other = makeTask("FN-OTHER", { column: "todo" });
const dependent = makeTask("FN-DEPENDENT", {
column: "todo",
status: "queued" as any,
blockedBy: "FN-B",
dependencies: ["FN-B", "FN-OTHER"],
});
const store = createStore([blocker, other, dependent]);
const mgr = new SelfHealingManager(store, { rootDir: "/repo" });
await Promise.all([mgr.reconcileCompletedTask("FN-B"), mgr.reconcileCompletedTask("FN-B")]);
const updated = await store.getTask("FN-DEPENDENT");
expect(updated?.blockedBy).toBe("FN-OTHER");
expect(updated?.queuedLogEpisodeSignature).toBe("dependency:FN-OTHER");
expect(updated?.log?.filter((entry) => entry.action.includes("FN-4523"))).toHaveLength(1);
});
it("prefers worktree hint and is idempotent when missing", async () => { it("prefers worktree hint and is idempotent when missing", async () => {
(existsSyncMock as any).mockImplementation((p: string) => p === "/wt/fn-b"); (existsSyncMock as any).mockImplementation((p: string) => p === "/wt/fn-b");
const blocker = makeTask("FN-B", { column: "done", branch: "fusion/fn-b" }); const blocker = makeTask("FN-B", { column: "done", branch: "fusion/fn-b" });

View File

@@ -204,6 +204,7 @@ function createMockStore(overrides: Record<string, unknown> = {}): TaskStore & E
} as unknown as Task), } as unknown as Task),
updateTask: vi.fn().mockResolvedValue({} as Task), updateTask: vi.fn().mockResolvedValue({} as Task),
logEntry: vi.fn().mockResolvedValue(undefined), logEntry: vi.fn().mockResolvedValue(undefined),
transitionQueuedEpisode: vi.fn().mockResolvedValue({ appended: true }),
moveTask: vi.fn().mockResolvedValue(undefined), moveTask: vi.fn().mockResolvedValue(undefined),
handoffToReview: vi.fn().mockResolvedValue(undefined), handoffToReview: vi.fn().mockResolvedValue(undefined),
enqueueMergeQueue: vi.fn().mockResolvedValue(undefined), enqueueMergeQueue: vi.fn().mockResolvedValue(undefined),
@@ -9477,6 +9478,22 @@ describe("FN-4538 overlapBlockedBy self-healing", () => {
if (options?.column === "in-review") return tasks.filter((task) => task.column === "in-review"); if (options?.column === "in-review") return tasks.filter((task) => task.column === "in-review");
return tasks; return tasks;
}); });
(store.updateTask as ReturnType<typeof vi.fn>).mockImplementation(async (id: string, patch: Record<string, unknown>) => {
const task = tasks.find((candidate) => candidate.id === id)!;
Object.assign(task, patch);
return task;
});
(store.transitionQueuedEpisode as ReturnType<typeof vi.fn>).mockImplementation(async (id: string, transition: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string }) => {
const task = tasks.find((candidate) => candidate.id === id)!;
const appended = !(task.status === "queued" && task.blockedBy === transition.blockedBy && task.overlapBlockedBy === transition.overlapBlockedBy && task.queuedLogEpisodeSignature === transition.signature);
await store.updateTask(id, { blockedBy: transition.blockedBy, status: "queued" });
task.blockedBy = transition.blockedBy;
task.overlapBlockedBy = transition.overlapBlockedBy;
task.status = "queued";
task.queuedLogEpisodeSignature = transition.signature;
if (appended) await store.logEntry(id, transition.action);
return { appended, task };
});
return store; return store;
} }
@@ -9555,7 +9572,7 @@ describe("FN-4538 overlapBlockedBy self-healing", () => {
manager.stop(); manager.stop();
}); });
it("FN-6276: clearStaleBlockedBy resets preserved queued memo after blocker resolves", async () => { it("FN-6276: recovery does not repeat a restored queued overlap episode", async () => {
const overlapBlocker = makeTask("FN-ACTIVE", { column: "in-progress" }); const overlapBlocker = makeTask("FN-ACTIVE", { column: "in-progress" });
const target = makeTask("FN-TARGET", { const target = makeTask("FN-TARGET", {
column: "todo", column: "todo",
@@ -9579,11 +9596,11 @@ describe("FN-4538 overlapBlockedBy self-healing", () => {
overlapBlocker.column = "in-progress"; overlapBlocker.column = "in-progress";
await manager.clearStaleBlockedBy(); await manager.clearStaleBlockedBy();
expect((store.logEntry as ReturnType<typeof vi.fn>).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(2); expect((store.logEntry as ReturnType<typeof vi.fn>).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(1);
manager.stop(); manager.stop();
}); });
it("FN-6276: stop clears preserved queued memo so next pass logs again", async () => { it("FN-8785: restart preserves an unchanged overlap episode without another log", async () => {
const overlapBlocker = makeTask("FN-ACTIVE", { column: "in-progress" }); const overlapBlocker = makeTask("FN-ACTIVE", { column: "in-progress" });
const target = makeTask("FN-TARGET", { const target = makeTask("FN-TARGET", {
column: "todo", column: "todo",
@@ -9600,7 +9617,7 @@ describe("FN-4538 overlapBlockedBy self-healing", () => {
manager.stop(); manager.stop();
await manager.clearStaleBlockedBy(); await manager.clearStaleBlockedBy();
expect((store.logEntry as ReturnType<typeof vi.fn>).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(2); expect((store.logEntry as ReturnType<typeof vi.fn>).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(1);
manager.stop(); manager.stop();
}); });

View File

@@ -12664,13 +12664,15 @@ export class TaskExecutor {
recoveryRehome: true, recoveryRehome: true,
}); });
} }
await this.store.updateTask(liveTask.id, { status: "queued", blockedBy: unmetDeps[0] }, this.getRunContextFor(liveTask.id)); const normalizedUnmetDeps = [...new Set(unmetDeps)].sort();
await this.store.logEntry( await this.store.transitionQueuedEpisode(liveTask.id, {
liveTask.id, signature: `dependency:${normalizedUnmetDeps.join(",")}`,
`queued — unmet dependencies: ${unmetDeps.join(", ")}`, blockedBy: unmetDeps[0] ?? null,
"Executor pre-dispatch dependency gate blocked workflow/authoritative execution.", overlapBlockedBy: liveTask.overlapBlockedBy ?? null,
this.getRunContextFor(liveTask.id), action: `queued — unmet dependencies: ${unmetDeps.join(", ")}`,
); outcome: "Executor pre-dispatch dependency gate blocked workflow/authoritative execution.",
runContext: this.getRunContextFor(liveTask.id),
});
executorLog.log(`${liveTask.id}: executor dispatch blocked by unmet dependencies: ${unmetDeps.join(", ")}`); executorLog.log(`${liveTask.id}: executor dispatch blocked by unmet dependencies: ${unmetDeps.join(", ")}`);
return true; return true;
} }

View File

@@ -1617,6 +1617,13 @@ export class Scheduler {
} }
} }
private async transitionQueuedEpisode(
task: Task,
input: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string },
): Promise<boolean> {
return (await this.store.transitionQueuedEpisode(task.id, input)).appended;
}
private async logDispatchQueuedReason(taskId: string, reason: string, memoKey?: string): Promise<boolean> { private async logDispatchQueuedReason(taskId: string, reason: string, memoKey?: string): Promise<boolean> {
const key = `${taskId}:${memoKey ?? reason}`; const key = `${taskId}:${memoKey ?? reason}`;
if (this.wasDispatchQueuedReasonLogged.has(key)) { if (this.wasDispatchQueuedReasonLogged.has(key)) {
@@ -2398,11 +2405,13 @@ export class Scheduler {
const unmetDeps = getUnmetSchedulingDependencies(task, tasks, schedulingDependencyOptions); const unmetDeps = getUnmetSchedulingDependencies(task, tasks, schedulingDependencyOptions);
if (unmetDeps.length > 0) { if (unmetDeps.length > 0) {
await this.store.updateTask(task.id, { const normalizedUnmetDeps = [...new Set(unmetDeps)].sort();
status: "queued", await this.transitionQueuedEpisode(task, {
blockedBy: unmetDeps[0], signature: `dependency:${normalizedUnmetDeps.join(",")}`,
blockedBy: unmetDeps[0] ?? null,
overlapBlockedBy: task.overlapBlockedBy ?? null,
action: `queued — unmet dependencies: ${unmetDeps.join(", ")}`,
}); });
await this.logDispatchQueuedReason(task.id, `queued — unmet dependencies: ${unmetDeps.join(", ")}`);
this.options.onBlocked?.(task, unmetDeps); this.options.onBlocked?.(task, unmetDeps);
return null; return null;
} }
@@ -2876,16 +2885,13 @@ export class Scheduler {
if (overlappingTaskId) { if (overlappingTaskId) {
const activeLeaseColumn = activeScopeColumns.get(overlappingTaskId) ?? "in-progress"; const activeLeaseColumn = activeScopeColumns.get(overlappingTaskId) ?? "in-progress";
await this.store.updateTask(task.id, { await this.transitionQueuedEpisode(task, {
status: "queued", signature: `file-scope:${overlappingTaskId}`,
blockedBy: null, blockedBy: null,
overlapBlockedBy: overlappingTaskId, overlapBlockedBy: overlappingTaskId,
action: `queued — blocked by active file-scope lease ${overlappingTaskId} (column=${activeLeaseColumn})`,
}); });
await this.rollbackRunningAgentsForQueuedTodoTask(task.id); await this.rollbackRunningAgentsForQueuedTodoTask(task.id);
await this.logDispatchQueuedReason(
task.id,
`queued — blocked by active file-scope lease ${overlappingTaskId} (column=${activeLeaseColumn})`,
);
return null; return null;
} }

View File

@@ -896,7 +896,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
private strandedHoldContinuationNoActionAudited = new Set<string>(); private strandedHoldContinuationNoActionAudited = new Set<string>();
/* FNXC:SymbolLock 2026-07-30-14:20: idle symbol-lock sweeps emit one no-action audit until a stale lock re-arms the diagnostic. */ /* FNXC:SymbolLock 2026-07-30-14:20: idle symbol-lock sweeps emit one no-action audit until a stale lock re-arms the diagnostic. */
private symbolLockNoActionAudited = false; private symbolLockNoActionAudited = false;
private preservedQueuedOverlapLogged = new Map<string, string>();
private maintenanceTickCounter = 0; private maintenanceTickCounter = 0;
private readonly taskLifecycleRetentionLastPrunedAt = new Map<string, number>(); private readonly taskLifecycleRetentionLastPrunedAt = new Map<string, number>();
private readonly processBootStartedAt = Date.now(); private readonly processBootStartedAt = Date.now();
@@ -1809,7 +1808,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
this.finalizeUnprovenWarned.clear(); this.finalizeUnprovenWarned.clear();
this.strandedCompletedFailureProvenanceWarned.clear(); this.strandedCompletedFailureProvenanceWarned.clear();
this.preservedQueuedOverlapLogged.clear();
log.debug("Stopped"); log.debug("Stopped");
} }
@@ -5050,19 +5048,28 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
const hasActiveOverlapBlocker = await hasActiveFileScopeOverlapBlocker(dependent, overlapBlockedBy); const hasActiveOverlapBlocker = await hasActiveFileScopeOverlapBlocker(dependent, overlapBlockedBy);
if (todoTaskIds.has(dependent.id)) { if (todoTaskIds.has(dependent.id)) {
/*
FNXC:QueuedTaskLogging 2026-08-04-18:32:
Completion fanout can race duplicate completion callbacks. Route retained dependency
and overlap queue episodes through the durable transition so its state repair and
one visible event commit together rather than reintroducing poll-driven task logs.
*/
if (unresolvedDeps.length > 0) { if (unresolvedDeps.length > 0) {
const nextBlocker = unresolvedDeps[0]!; const nextBlocker = unresolvedDeps[0]!;
await this.store.updateTask(dependent.id, { blockedBy: nextBlocker, overlapBlockedBy, status: "queued" }); const normalizedUnresolvedDeps = [...new Set(unresolvedDeps)].sort();
await this.store.logEntry( await this.store.transitionQueuedEpisode(dependent.id, {
dependent.id, signature: `dependency:${normalizedUnresolvedDeps.join(",")}`,
`Auto-recovered (FN-4523): cleared stale blockedBy — blocker ${taskId} is done; now blocked by ${nextBlocker}`, blockedBy: nextBlocker,
); overlapBlockedBy,
action: `Auto-recovered (FN-4523): cleared stale blockedBy — blocker ${taskId} is done; now blocked by ${nextBlocker}`,
});
} else if (hasActiveOverlapBlocker) { } else if (hasActiveOverlapBlocker) {
await this.store.updateTask(dependent.id, { blockedBy: null, overlapBlockedBy, status: "queued" }); await this.store.transitionQueuedEpisode(dependent.id, {
await this.store.logEntry( signature: `file-scope:${overlapBlockedBy}`,
dependent.id, blockedBy: null,
`Auto-recovered (FN-4523): preserved queued status — still blocked by file scope overlap with ${overlapBlockedBy}`, overlapBlockedBy,
); action: `Auto-recovered (FN-4523): preserved queued status — still blocked by file scope overlap with ${overlapBlockedBy}`,
});
} else { } else {
await this.store.updateTask(dependent.id, { blockedBy: null, overlapBlockedBy: null, status: null }); await this.store.updateTask(dependent.id, { blockedBy: null, overlapBlockedBy: null, status: null });
await this.store.logEntry( await this.store.logEntry(
@@ -5681,18 +5688,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
return tasks.filter((task) => typeof task.blockedBy === "string" && task.blockedBy.trim().length > 0).length; return tasks.filter((task) => typeof task.blockedBy === "string" && task.blockedBy.trim().length > 0).length;
} }
private shouldLogPreservedQueuedOverlap(taskId: string, overlapBlockedBy: string | null | undefined): overlapBlockedBy is string {
if (!overlapBlockedBy) return false;
const previous = this.preservedQueuedOverlapLogged.get(taskId);
if (previous === overlapBlockedBy) return false;
this.preservedQueuedOverlapLogged.set(taskId, overlapBlockedBy);
return true;
}
private clearPreservedQueuedOverlapMemo(taskId: string): void {
this.preservedQueuedOverlapLogged.delete(taskId);
}
/** /**
* #1401: periodic transitionPending recovery sweep. Delegates to the store's * #1401: periodic transitionPending recovery sweep. Delegates to the store's
* idempotent recovery method (a no-op when no stale markers exist), keeping * idempotent recovery method (a no-op when no stale markers exist), keeping
@@ -5996,7 +5991,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
); );
if (blockedTasks.length === 0 && queuedDependencyTasks.length === 0) { if (blockedTasks.length === 0 && queuedDependencyTasks.length === 0) {
this.preservedQueuedOverlapLogged.clear();
return 0; return 0;
} }
@@ -6105,50 +6099,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
hold: new Set(["todo"]), review: new Set(["in-review"]), hold: new Set(["todo"]), review: new Set(["in-review"]),
}; };
for (const [taskId, lastLoggedBlockerId] of this.preservedQueuedOverlapLogged) {
const memoTask = taskById.get(taskId);
const memoHasActiveOverlapBlocker = memoTask
? await hasActiveFileScopeOverlapBlocker(memoTask, memoTask.overlapBlockedBy)
: false;
if (
!candidates.has(taskId)
/*
FNXC:WorkflowResolvedColumns 2026-07-30-21:40 (FLAGGED AND LEFT COUNTED):
This sits in a log-dedup closure defined BEFORE the per-referenced-task lane prefetch below, so
the resolved sets are not in scope here and tsc says so. Hoisting the prefetch above the closure
is not available either — it is keyed on `candidates`, which this closure helps build.
Left as the literal rather than restructured: the closure only decides whether to re-log an
already-logged blocker, so the degraded answer costs a duplicate log line on a renamed board, not
a wrong lifecycle decision. Restructuring a sweep's control flow to convert a logging guard is
the wrong trade.
*/
/*
FNXC:WorkflowResolvedColumns 2026-07-31-12:10 (u12 — CONVERTED; the stated blocker was not real):
The prior note said the resolved sets could not be in scope because the lane prefetch is keyed on
`candidates`, "which this closure helps build". Measured: this loop does not build `candidates` —
it is fully populated two statements above (`blockedTasks` then `queuedDependencyTasks`) and this
loop only CLEARS memo entries. So the prefetch was hoistable, and it now sits above.
Reaching this clause already proves the id is a candidate: `!candidates.has(taskId)` is the first
arm of the same `||` chain, so short-circuit means we only ask the lane question for ids the
prefetch covered (`referencedIds.add(task.id)` for every candidate). `lanesOf` still falls back to
the legacy set, so an unresolvable workflow answers exactly as the literal did.
`|| !memoTask` is explicit rather than implied: `memoTask?.column !== "todo"` was ALSO the undefined
check, and tsc narrowed the clauses after it on that basis. Dropping it silently broke the
narrowing (TS18048 on `memoTask.status`), so the undefined arm is now stated on its own line.
*/
|| !memoTask
|| !lanesOf(taskId).hold.has(memoTask.column)
|| memoTask.status !== "queued"
|| memoTask.overlapBlockedBy !== lastLoggedBlockerId
|| !memoHasActiveOverlapBlocker
) {
this.clearPreservedQueuedOverlapMemo(taskId);
}
}
for (const task of candidates.values()) { for (const task of candidates.values()) {
const blockerId = task.blockedBy; const blockerId = task.blockedBy;
@@ -6246,7 +6196,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
let didRecover = false; let didRecover = false;
if (todoTaskIds.has(task.id)) { if (todoTaskIds.has(task.id)) {
if (unresolvedDeps.length > 0) { if (unresolvedDeps.length > 0) {
this.clearPreservedQueuedOverlapMemo(task.id);
const nextBlocker = unresolvedDeps[0]!; const nextBlocker = unresolvedDeps[0]!;
if (nextBlocker === blockerId) { if (nextBlocker === blockerId) {
continue; continue;
@@ -6255,19 +6204,19 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
await this.store.logEntry(task.id, `Auto-recovered (FN-5488): refreshed stale blockedBy — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; ${reason}; now blocked by ${nextBlocker}`); await this.store.logEntry(task.id, `Auto-recovered (FN-5488): refreshed stale blockedBy — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; ${reason}; now blocked by ${nextBlocker}`);
didRecover = true; didRecover = true;
} else if (hasActiveOverlapBlocker) { } else if (hasActiveOverlapBlocker) {
await this.store.updateTask(task.id, { blockedBy: null, status: "queued" }); const transition = await this.store.transitionQueuedEpisode(task.id, {
if (this.shouldLogPreservedQueuedOverlap(task.id, task.overlapBlockedBy)) { signature: `file-scope:${task.overlapBlockedBy}`,
await this.store.logEntry(task.id, `Auto-recovered (FN-5488): preserved queued status — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; still blocked by file scope overlap with ${task.overlapBlockedBy}`); blockedBy: null,
didRecover = true; overlapBlockedBy: task.overlapBlockedBy ?? null,
} action: `Auto-recovered (FN-5488): preserved queued status — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; still blocked by file scope overlap with ${task.overlapBlockedBy}`,
});
didRecover = transition.appended;
} else { } else {
this.clearPreservedQueuedOverlapMemo(task.id);
await this.store.updateTask(task.id, { blockedBy: null, overlapBlockedBy: null, status: null }); await this.store.updateTask(task.id, { blockedBy: null, overlapBlockedBy: null, status: null });
await this.store.logEntry(task.id, `Auto-recovered (FN-5488): cleared stale blockedBy — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; ${reason}`); await this.store.logEntry(task.id, `Auto-recovered (FN-5488): cleared stale blockedBy — blocker=${blockerId} blockerStatus=${blocker?.status ?? "none"} reason=${reasonCode ?? "unspecified"}; ${reason}`);
didRecover = true; didRecover = true;
} }
} else { } else {
this.clearPreservedQueuedOverlapMemo(task.id);
await this.store.updateTask(task.id, { blockedBy: null }); await this.store.updateTask(task.id, { blockedBy: null });
await this.store.logEntry(task.id, `Auto-recovered (FN-4091): cleared stale blockedBy — ${reason}`); await this.store.logEntry(task.id, `Auto-recovered (FN-4091): cleared stale blockedBy — ${reason}`);
didRecover = true; didRecover = true;
@@ -6289,13 +6238,14 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
if (queuedDependencyTaskIds.has(task.id)) { if (queuedDependencyTaskIds.has(task.id)) {
try { try {
if (hasActiveOverlapBlocker) { if (hasActiveOverlapBlocker) {
await this.store.updateTask(task.id, { blockedBy: null, status: "queued" }); const transition = await this.store.transitionQueuedEpisode(task.id, {
if (this.shouldLogPreservedQueuedOverlap(task.id, task.overlapBlockedBy)) { signature: `file-scope:${task.overlapBlockedBy}`,
await this.store.logEntry(task.id, `Auto-recovered: preserved queued status — still blocked by file scope overlap with ${task.overlapBlockedBy}`); blockedBy: null,
recovered++; overlapBlockedBy: task.overlapBlockedBy ?? null,
} action: `Auto-recovered: preserved queued status — still blocked by file scope overlap with ${task.overlapBlockedBy}`,
});
if (transition.appended) recovered++;
} else { } else {
this.clearPreservedQueuedOverlapMemo(task.id);
// FN-5434: routine scheduler↔self-healing queued-status churn should stay silent; keep state cleanup only. // FN-5434: routine scheduler↔self-healing queued-status churn should stay silent; keep state cleanup only.
await this.store.updateTask(task.id, { blockedBy: null, overlapBlockedBy: null, status: null }); await this.store.updateTask(task.id, { blockedBy: null, overlapBlockedBy: null, status: null });
} }
@@ -6307,7 +6257,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
continue; continue;
} }
this.clearPreservedQueuedOverlapMemo(task.id);
const nextBlocker = unresolvedDeps[0] ?? null; const nextBlocker = unresolvedDeps[0] ?? null;
if (nextBlocker && task.blockedBy !== nextBlocker) { if (nextBlocker && task.blockedBy !== nextBlocker) {
try { try {