From 1dc636b8b455d869a5531323fcc7f1937a2d8307 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Tue, 4 Aug 2026 11:52:15 -0700 Subject: [PATCH] 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) --- .../fn-8785-queued-log-deduplication.md | 7 + docs/architecture.md | 1 + .../queued-episode-transition.pg.test.ts | 158 ++++++++++++++++++ .../__tests__/postgres/schema-applier.test.ts | 9 +- .../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 + ...tor-outer-dispatch-dependency-gate.test.ts | 39 ++--- .../src/__tests__/executor-test-helpers.ts | 7 + .../scheduler-overlap-starvation.test.ts | 43 ++++- .../scheduler-workflow-cutover.test.ts | 35 +++- .../self-healing-completion-fanout.test.ts | 37 ++++ .../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(-) create mode 100644 .changeset/fn-8785-queued-log-deduplication.md create mode 100644 packages/core/src/__tests__/postgres/queued-episode-transition.pg.test.ts create mode 100644 packages/core/src/postgres/migrations/0044_fn_8785_queued_episode_signature.sql diff --git a/.changeset/fn-8785-queued-log-deduplication.md b/.changeset/fn-8785-queued-log-deduplication.md new file mode 100644 index 0000000000..9c22ea1f6f --- /dev/null +++ b/.changeset/fn-8785-queued-log-deduplication.md @@ -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. diff --git a/docs/architecture.md b/docs/architecture.md index 043f24196d..25a43b8931 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -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. - 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. + - **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. #### BlockedBy stamping invariants diff --git a/packages/core/src/__tests__/postgres/queued-episode-transition.pg.test.ts b/packages/core/src/__tests__/postgres/queued-episode-transition.pg.test.ts new file mode 100644 index 0000000000..294dfe41e3 --- /dev/null +++ b/packages/core/src/__tests__/postgres/queued-episode-transition.pg.test.ts @@ -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); + }); +}); diff --git a/packages/core/src/__tests__/postgres/schema-applier.test.ts b/packages/core/src/__tests__/postgres/schema-applier.test.ts index a15ea1e883..e9a2696c09 100644 --- a/packages/core/src/__tests__/postgres/schema-applier.test.ts +++ b/packages/core/src/__tests__/postgres/schema-applier.test.ts @@ -86,6 +86,7 @@ import { TASK_LIFECYCLE_CONSUMERS_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, + QUEUED_EPISODE_SIGNATURE_VERSION, } from "../../postgres/schema-applier.js"; import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.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(VALIDATOR_INPUT_FINGERPRINT_VERSION).toBe("0042"); 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", () => { @@ -1761,6 +1763,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { TASK_LIFECYCLE_CONSUMERS_VERSION, VALIDATOR_INPUT_FINGERPRINT_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, + QUEUED_EPISODE_SIGNATURE_VERSION, ]); expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false); }); @@ -1830,6 +1833,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { TASK_LIFECYCLE_CONSUMERS_VERSION, VALIDATOR_INPUT_FINGERPRINT_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, VALIDATOR_INPUT_FINGERPRINT_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, VALIDATOR_INPUT_FINGERPRINT_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, VALIDATOR_INPUT_FINGERPRINT_VERSION, UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION, + QUEUED_EPISODE_SIGNATURE_VERSION, ]); }); }); diff --git a/packages/core/src/postgres/migrations/0000_initial.sql b/packages/core/src/postgres/migrations/0000_initial.sql index b534efbe14..92b5b7cb67 100644 --- a/packages/core/src/postgres/migrations/0000_initial.sql +++ b/packages/core/src/postgres/migrations/0000_initial.sql @@ -52,6 +52,7 @@ CREATE TABLE IF NOT EXISTS project.tasks ( worktree text, blocked_by text, overlap_blocked_by text, + queued_log_episode_signature text, paused integer DEFAULT 0, user_paused integer DEFAULT 0, paused_reason text, diff --git a/packages/core/src/postgres/migrations/0044_fn_8785_queued_episode_signature.sql b/packages/core/src/postgres/migrations/0044_fn_8785_queued_episode_signature.sql new file mode 100644 index 0000000000..432921305a --- /dev/null +++ b/packages/core/src/postgres/migrations/0044_fn_8785_queued_episode_signature.sql @@ -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; diff --git a/packages/core/src/postgres/schema-applier.ts b/packages/core/src/postgres/schema-applier.ts index 698b995292..4c6e2a1651 100644 --- a/packages/core/src/postgres/schema-applier.ts +++ b/packages/core/src/postgres/schema-applier.ts @@ -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: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. */ export const TASK_DECLARED_SYMBOLS_VERSION = "0028"; const INITIAL_SCHEMA_VERSION = "0000"; @@ -185,6 +185,8 @@ export const TASK_LIFECYCLE_CONSUMERS_VERSION = "0041"; export const VALIDATOR_INPUT_FINGERPRINT_VERSION = "0042"; /** FNXC:PlanningDependencyReseed 2026-08-04-02:14: durable per-episode unplanned-dispatch diagnostics. */ 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. */ 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 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 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 @@ -515,6 +518,7 @@ export async function applySchemaBaseline( const taskLifecycleConsumersAlreadyApplied = applied.includes(TASK_LIFECYCLE_CONSUMERS_VERSION); const validatorInputFingerprintAlreadyApplied = applied.includes(VALIDATOR_INPUT_FINGERPRINT_VERSION); const unplannedExecutionBlockDedupeAlreadyApplied = applied.includes(UNPLANNED_EXECUTION_BLOCK_DEDUPE_VERSION); + const queuedEpisodeSignatureAlreadyApplied = applied.includes(QUEUED_EPISODE_SIGNATURE_VERSION); assertBinaryNotOlderThanDatabase(applied); 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`); 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 }; }); diff --git a/packages/core/src/postgres/schema/project.ts b/packages/core/src/postgres/schema/project.ts index 86d8f189bb..0b2fcf977e 100644 --- a/packages/core/src/postgres/schema/project.ts +++ b/packages/core/src/postgres/schema/project.ts @@ -80,6 +80,7 @@ export const tasks = projectSchema.table("tasks", { worktree: text("worktree"), blockedBy: text("blocked_by"), overlapBlockedBy: text("overlap_blocked_by"), + queuedLogEpisodeSignature: text("queued_log_episode_signature"), paused: integer("paused").default(0), userPaused: integer("user_paused").default(0), pausedReason: text("paused_reason"), diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 7013662f06..4baf4d94dc 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -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 { 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 { 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 { flushAgentLogBufferImpl, appendAgentLogBatchImpl } from "./task-store/agent-logs.js"; import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js"; @@ -1712,6 +1712,9 @@ export class TaskStore extends EventEmitter { async checkAndRecordUnplannedExecutionBlock(id: string, episode: string): Promise { return checkAndRecordUnplannedExecutionBlockImpl(this, id, episode); } + async transitionQueuedEpisode(id: string, transition: QueuedEpisodeTransition): Promise { + return transitionQueuedEpisodeImpl(this, id, transition); + } async logEntry(id: string, action: string, outcome?: string, runContext?: RunMutationContext): Promise { return logEntryImpl(this, id, action, outcome, runContext); } diff --git a/packages/core/src/task-store/audit-ops.ts b/packages/core/src/task-store/audit-ops.ts index 89965abeb1..164b8bbb2a 100644 --- a/packages/core/src/task-store/audit-ops.ts +++ b/packages/core/src/task-store/audit-ops.ts @@ -18,6 +18,7 @@ import "../builtin-traits.js"; import {__setTaskActivityLogLimitsForTesting, truncateTaskLogOutcome, getTaskActivityLogEntryLimit} from "../task-store/comments.js"; import {readTaskRow, updateTaskColumns} from "../task-store/async/async-persistence.js"; import { getLiveTaskColumn } from "./async/async-comments-attachments.js"; +import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js"; import { resolveArchivedLanes } from "../project-lane-vocabulary.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 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 { + 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)); + 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( store: TaskStore, id: string, diff --git a/packages/core/src/task-store/persistence.ts b/packages/core/src/task-store/persistence.ts index b62b51b043..336b168eb5 100644 --- a/packages/core/src/task-store/persistence.ts +++ b/packages/core/src/task-store/persistence.ts @@ -26,6 +26,7 @@ export interface TaskRow { worktree: string | null; blockedBy: string | null; overlapBlockedBy: string | null; + queuedLogEpisodeSignature: string | null; paused: number | null; pausedReason: string | null; wedgeNotification: string | null; @@ -247,6 +248,7 @@ export const TASK_COLUMN_DESCRIPTORS: TaskColumnDescriptor[] = [ defineTaskColumn("worktree", (task) => task.worktree ?? null), defineTaskColumn("blockedBy", (task) => task.blockedBy ?? null), defineTaskColumn("overlapBlockedBy", (task) => task.overlapBlockedBy ?? null), + defineTaskColumn("queuedLogEpisodeSignature", (task) => task.queuedLogEpisodeSignature ?? null), defineTaskColumn("paused", (task) => task.paused ? 1 : 0), defineTaskColumn("pausedReason", (task) => task.pausedReason ?? null), defineTaskColumn("wedgeNotification", (task) => toJsonNullable(task.wedgeNotification)), diff --git a/packages/core/src/task-store/serialization.ts b/packages/core/src/task-store/serialization.ts index 2445587a71..e584e84e04 100644 --- a/packages/core/src/task-store/serialization.ts +++ b/packages/core/src/task-store/serialization.ts @@ -76,6 +76,7 @@ export function rowToTask(row: TaskRow): Task { worktree: row.worktree || undefined, blockedBy: row.blockedBy || undefined, overlapBlockedBy: row.overlapBlockedBy || undefined, + queuedLogEpisodeSignature: row.queuedLogEpisodeSignature || undefined, paused: row.paused ? true : undefined, pausedReason: row.pausedReason || undefined, wedgeNotification: fromJson(row.wedgeNotification) ?? undefined, diff --git a/packages/core/src/types/task/task-core.ts b/packages/core/src/types/task/task-core.ts index fb6b8a6262..57909b70de 100644 --- a/packages/core/src/types/task/task-core.ts +++ b/packages/core/src/types/task/task-core.ts @@ -660,6 +660,11 @@ export interface Task { * Cleared when the overlap resolves (the blocker task moves to done or its * scope no longer overlaps). */ 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. */ paused?: boolean; /** When true, this task was explicitly moved back to todo by a user and should not auto-dispatch. */ diff --git a/packages/engine/src/__tests__/executor-outer-dispatch-dependency-gate.test.ts b/packages/engine/src/__tests__/executor-outer-dispatch-dependency-gate.test.ts index db74a32bfb..a5e6657b3f 100644 --- a/packages/engine/src/__tests__/executor-outer-dispatch-dependency-gate.test.ts +++ b/packages/engine/src/__tests__/executor-outer-dispatch-dependency-gate.test.ts @@ -119,17 +119,14 @@ describe("executor outer dispatch dependency gate", () => { preserveResumeState: true, recoveryRehome: true, })); - expect(store.updateTask).toHaveBeenCalledWith( - child.id, - expect.objectContaining({ status: "queued", blockedBy: parent.id }), - undefined, - ); - expect(store.logEntry).toHaveBeenCalledWith( - child.id, - expect.stringContaining("queued — unmet dependencies: FN-PARENT"), - expect.stringContaining("dependency gate blocked"), - undefined, - ); + expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({ + signature: "dependency:FN-PARENT", + blockedBy: parent.id, + })); + expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({ + action: expect.stringContaining("queued — unmet dependencies: FN-PARENT"), + outcome: expect.stringContaining("dependency gate blocked"), + })); expect(graph).not.toHaveBeenCalled(); // 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. @@ -148,11 +145,10 @@ describe("executor outer dispatch dependency gate", () => { await executor.execute(child); - expect(store.updateTask).toHaveBeenCalledWith( - child.id, - expect.objectContaining({ status: "queued", blockedBy: parent.id }), - undefined, - ); + expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({ + signature: "dependency:FN-PARENT", + blockedBy: parent.id, + })); expect(graph).not.toHaveBeenCalled(); }); @@ -167,7 +163,7 @@ describe("executor outer dispatch dependency gate", () => { await executor.execute(child); 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 }); }); @@ -196,11 +192,10 @@ describe("executor outer dispatch dependency gate", () => { await executor.execute(child); expect(store.getCompletionHandoffAcceptedMarker).toHaveBeenCalledWith(parent.id); - expect(store.updateTask).toHaveBeenCalledWith( - child.id, - expect.objectContaining({ status: "queued", blockedBy: parent.id }), - undefined, - ); + expect(store.transitionQueuedEpisode).toHaveBeenCalledWith(child.id, expect.objectContaining({ + signature: "dependency:FN-PARENT", + blockedBy: parent.id, + })); expect(graph).not.toHaveBeenCalled(); }); diff --git a/packages/engine/src/__tests__/executor-test-helpers.ts b/packages/engine/src/__tests__/executor-test-helpers.ts index 26bdd318c1..7ead94617c 100644 --- a/packages/engine/src/__tests__/executor-test-helpers.ts +++ b/packages/engine/src/__tests__/executor-test-helpers.ts @@ -646,6 +646,13 @@ export function createMockStore() { updatedAt: new Date().toISOString(), })), 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), parseStepsFromPrompt: vi.fn().mockResolvedValue([]), parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]), diff --git a/packages/engine/src/__tests__/scheduler-overlap-starvation.test.ts b/packages/engine/src/__tests__/scheduler-overlap-starvation.test.ts index 4cbea214db..ca94621d67 100644 --- a/packages/engine/src/__tests__/scheduler-overlap-starvation.test.ts +++ b/packages/engine/src/__tests__/scheduler-overlap-starvation.test.ts @@ -57,6 +57,7 @@ function createStore(tasks: Task[], scopes: Record, settings: if (task) Object.assign(task, patch); return task as Task; }); + const logEntry = vi.fn(async () => undefined); const moveTask = vi.fn(async (id: string, column: Task["column"]) => { const task = tasks.find((candidate) => candidate.id === id); if (task) task.column = column; @@ -88,7 +89,22 @@ function createStore(tasks: Task[], scopes: Record, settings: moveTask, moveTaskIf, 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"), getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"), 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()); }); + 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).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 () => { const tasks = [ makeTask({ id: "FN-901", column: "todo", status: "queued", priority: "normal", overlapBlockedBy: "FN-MISSING" }), diff --git a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts index 9e4ae01b39..d022fe65ef 100644 --- a/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts +++ b/packages/engine/src/__tests__/scheduler-workflow-cutover.test.ts @@ -53,6 +53,12 @@ function storeWith( workflows: { selections?: Record; definitions?: Record } = {}, ): TaskStore { const byId = new Map(tasks.map((candidate) => [candidate.id, candidate])); + const updateTask = vi.fn(async (id: string, patch: Partial) => { + const current = byId.get(id); + if (current) Object.assign(current, patch); + return current as Task; + }); + const logEntry = vi.fn(async () => undefined); return { listTasks: vi.fn(async () => [...byId.values()]), getTask: vi.fn(async (id: string) => byId.get(id) ?? null), @@ -63,11 +69,7 @@ function storeWith( ...settings, })), updateSettings: vi.fn(async (patch: Record) => ({ ...settings, ...patch })), - updateTask: vi.fn(async (id: string, patch: Partial) => { - const current = byId.get(id); - if (current) Object.assign(current, patch); - return current as Task; - }), + updateTask, moveTask: vi.fn(async (id: string, column: Task["column"]) => { const current = byId.get(id); if (current) current.column = column; @@ -80,7 +82,22 @@ function storeWith( return { task: current, moved: true }; }), 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"), getTasksDir: vi.fn(() => "/tmp/project/.fusion/tasks"), on: vi.fn(), @@ -388,9 +405,11 @@ describe("Scheduler workflow cutover", () => { await scheduler.schedule(); - expect(store.updateTask).toHaveBeenCalledWith("FN-002", { - status: "queued", + expect(store.transitionQueuedEpisode).toHaveBeenCalledWith("FN-002", { + signature: "dependency: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(onBlocked).toHaveBeenCalledWith(expect.objectContaining({ id: "FN-002" }), ["FN-001"]); diff --git a/packages/engine/src/__tests__/self-healing-completion-fanout.test.ts b/packages/engine/src/__tests__/self-healing-completion-fanout.test.ts index 5f2447cd59..d5d28a14a8 100644 --- a/packages/engine/src/__tests__/self-healing-completion-fanout.test.ts +++ b/packages/engine/src/__tests__/self-healing-completion-fanout.test.ts @@ -70,6 +70,23 @@ function createStore(tasks: Task[], settings?: Partial): TaskStore & E map.set(id, { ...task, ...patch } as Task); 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"]) => { const task = map.get(id)!; 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 () => { (existsSyncMock as any).mockImplementation((p: string) => p === "/wt/fn-b"); const blocker = makeTask("FN-B", { column: "done", branch: "fusion/fn-b" }); diff --git a/packages/engine/src/__tests__/self-healing.test.ts b/packages/engine/src/__tests__/self-healing.test.ts index 2f31cdde19..281fb02764 100644 --- a/packages/engine/src/__tests__/self-healing.test.ts +++ b/packages/engine/src/__tests__/self-healing.test.ts @@ -204,6 +204,7 @@ function createMockStore(overrides: Record = {}): TaskStore & E } as unknown as Task), updateTask: vi.fn().mockResolvedValue({} as Task), logEntry: vi.fn().mockResolvedValue(undefined), + transitionQueuedEpisode: vi.fn().mockResolvedValue({ appended: true }), moveTask: vi.fn().mockResolvedValue(undefined), handoffToReview: 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"); return tasks; }); + (store.updateTask as ReturnType).mockImplementation(async (id: string, patch: Record) => { + const task = tasks.find((candidate) => candidate.id === id)!; + Object.assign(task, patch); + return task; + }); + (store.transitionQueuedEpisode as ReturnType).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; } @@ -9555,7 +9572,7 @@ describe("FN-4538 overlapBlockedBy self-healing", () => { 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 target = makeTask("FN-TARGET", { column: "todo", @@ -9579,11 +9596,11 @@ describe("FN-4538 overlapBlockedBy self-healing", () => { overlapBlocker.column = "in-progress"; await manager.clearStaleBlockedBy(); - expect((store.logEntry as ReturnType).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(2); + expect((store.logEntry as ReturnType).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(1); 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 target = makeTask("FN-TARGET", { column: "todo", @@ -9600,7 +9617,7 @@ describe("FN-4538 overlapBlockedBy self-healing", () => { manager.stop(); await manager.clearStaleBlockedBy(); - expect((store.logEntry as ReturnType).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(2); + expect((store.logEntry as ReturnType).mock.calls.filter((call) => call[0] === "FN-TARGET" && call[1] === message)).toHaveLength(1); manager.stop(); }); diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index e852a59917..6bca0413ab 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -12664,13 +12664,15 @@ export class TaskExecutor { recoveryRehome: true, }); } - await this.store.updateTask(liveTask.id, { status: "queued", blockedBy: unmetDeps[0] }, this.getRunContextFor(liveTask.id)); - await this.store.logEntry( - liveTask.id, - `queued — unmet dependencies: ${unmetDeps.join(", ")}`, - "Executor pre-dispatch dependency gate blocked workflow/authoritative execution.", - this.getRunContextFor(liveTask.id), - ); + const normalizedUnmetDeps = [...new Set(unmetDeps)].sort(); + await this.store.transitionQueuedEpisode(liveTask.id, { + signature: `dependency:${normalizedUnmetDeps.join(",")}`, + blockedBy: unmetDeps[0] ?? null, + overlapBlockedBy: liveTask.overlapBlockedBy ?? null, + 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(", ")}`); return true; } diff --git a/packages/engine/src/scheduler.ts b/packages/engine/src/scheduler.ts index 3b0ba8e00e..2fa47575c1 100644 --- a/packages/engine/src/scheduler.ts +++ b/packages/engine/src/scheduler.ts @@ -1617,6 +1617,13 @@ export class Scheduler { } } + private async transitionQueuedEpisode( + task: Task, + input: { signature: string; blockedBy: string | null; overlapBlockedBy: string | null; action: string }, + ): Promise { + return (await this.store.transitionQueuedEpisode(task.id, input)).appended; + } + private async logDispatchQueuedReason(taskId: string, reason: string, memoKey?: string): Promise { const key = `${taskId}:${memoKey ?? reason}`; if (this.wasDispatchQueuedReasonLogged.has(key)) { @@ -2398,11 +2405,13 @@ export class Scheduler { const unmetDeps = getUnmetSchedulingDependencies(task, tasks, schedulingDependencyOptions); if (unmetDeps.length > 0) { - await this.store.updateTask(task.id, { - status: "queued", - blockedBy: unmetDeps[0], + const normalizedUnmetDeps = [...new Set(unmetDeps)].sort(); + await this.transitionQueuedEpisode(task, { + 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); return null; } @@ -2876,16 +2885,13 @@ export class Scheduler { if (overlappingTaskId) { const activeLeaseColumn = activeScopeColumns.get(overlappingTaskId) ?? "in-progress"; - await this.store.updateTask(task.id, { - status: "queued", + await this.transitionQueuedEpisode(task, { + signature: `file-scope:${overlappingTaskId}`, blockedBy: null, overlapBlockedBy: overlappingTaskId, + action: `queued — blocked by active file-scope lease ${overlappingTaskId} (column=${activeLeaseColumn})`, }); await this.rollbackRunningAgentsForQueuedTodoTask(task.id); - await this.logDispatchQueuedReason( - task.id, - `queued — blocked by active file-scope lease ${overlappingTaskId} (column=${activeLeaseColumn})`, - ); return null; } diff --git a/packages/engine/src/self-healing.ts b/packages/engine/src/self-healing.ts index 3819670c69..a76105b6d0 100644 --- a/packages/engine/src/self-healing.ts +++ b/packages/engine/src/self-healing.ts @@ -896,7 +896,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence { private strandedHoldContinuationNoActionAudited = new Set(); /* 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 preservedQueuedOverlapLogged = new Map(); private maintenanceTickCounter = 0; private readonly taskLifecycleRetentionLastPrunedAt = new Map(); private readonly processBootStartedAt = Date.now(); @@ -1809,7 +1808,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence { this.finalizeUnprovenWarned.clear(); this.strandedCompletedFailureProvenanceWarned.clear(); - this.preservedQueuedOverlapLogged.clear(); log.debug("Stopped"); } @@ -5050,19 +5048,28 @@ export class SelfHealingManager extends SelfHealingGitEvidence { const hasActiveOverlapBlocker = await hasActiveFileScopeOverlapBlocker(dependent, overlapBlockedBy); 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) { const nextBlocker = unresolvedDeps[0]!; - await this.store.updateTask(dependent.id, { blockedBy: nextBlocker, overlapBlockedBy, status: "queued" }); - await this.store.logEntry( - dependent.id, - `Auto-recovered (FN-4523): cleared stale blockedBy — blocker ${taskId} is done; now blocked by ${nextBlocker}`, - ); + const normalizedUnresolvedDeps = [...new Set(unresolvedDeps)].sort(); + await this.store.transitionQueuedEpisode(dependent.id, { + signature: `dependency:${normalizedUnresolvedDeps.join(",")}`, + blockedBy: nextBlocker, + overlapBlockedBy, + action: `Auto-recovered (FN-4523): cleared stale blockedBy — blocker ${taskId} is done; now blocked by ${nextBlocker}`, + }); } else if (hasActiveOverlapBlocker) { - await this.store.updateTask(dependent.id, { blockedBy: null, overlapBlockedBy, status: "queued" }); - await this.store.logEntry( - dependent.id, - `Auto-recovered (FN-4523): preserved queued status — still blocked by file scope overlap with ${overlapBlockedBy}`, - ); + await this.store.transitionQueuedEpisode(dependent.id, { + signature: `file-scope:${overlapBlockedBy}`, + blockedBy: null, + overlapBlockedBy, + action: `Auto-recovered (FN-4523): preserved queued status — still blocked by file scope overlap with ${overlapBlockedBy}`, + }); } else { await this.store.updateTask(dependent.id, { blockedBy: null, overlapBlockedBy: null, status: null }); 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; } - 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 * 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) { - this.preservedQueuedOverlapLogged.clear(); return 0; } @@ -6105,50 +6099,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence { 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()) { const blockerId = task.blockedBy; @@ -6246,7 +6196,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence { let didRecover = false; if (todoTaskIds.has(task.id)) { if (unresolvedDeps.length > 0) { - this.clearPreservedQueuedOverlapMemo(task.id); const nextBlocker = unresolvedDeps[0]!; if (nextBlocker === blockerId) { 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}`); didRecover = true; } else if (hasActiveOverlapBlocker) { - await this.store.updateTask(task.id, { blockedBy: null, status: "queued" }); - if (this.shouldLogPreservedQueuedOverlap(task.id, 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}`); - didRecover = true; - } + const transition = await this.store.transitionQueuedEpisode(task.id, { + signature: `file-scope:${task.overlapBlockedBy}`, + blockedBy: null, + 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 { - this.clearPreservedQueuedOverlapMemo(task.id); 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}`); didRecover = true; } } else { - this.clearPreservedQueuedOverlapMemo(task.id); await this.store.updateTask(task.id, { blockedBy: null }); await this.store.logEntry(task.id, `Auto-recovered (FN-4091): cleared stale blockedBy — ${reason}`); didRecover = true; @@ -6289,13 +6238,14 @@ export class SelfHealingManager extends SelfHealingGitEvidence { if (queuedDependencyTaskIds.has(task.id)) { try { if (hasActiveOverlapBlocker) { - await this.store.updateTask(task.id, { blockedBy: null, status: "queued" }); - if (this.shouldLogPreservedQueuedOverlap(task.id, task.overlapBlockedBy)) { - await this.store.logEntry(task.id, `Auto-recovered: preserved queued status — still blocked by file scope overlap with ${task.overlapBlockedBy}`); - recovered++; - } + const transition = await this.store.transitionQueuedEpisode(task.id, { + signature: `file-scope:${task.overlapBlockedBy}`, + blockedBy: null, + overlapBlockedBy: task.overlapBlockedBy ?? null, + action: `Auto-recovered: preserved queued status — still blocked by file scope overlap with ${task.overlapBlockedBy}`, + }); + if (transition.appended) recovered++; } else { - this.clearPreservedQueuedOverlapMemo(task.id); // 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 }); } @@ -6307,7 +6257,6 @@ export class SelfHealingManager extends SelfHealingGitEvidence { continue; } - this.clearPreservedQueuedOverlapMemo(task.id); const nextBlocker = unresolvedDeps[0] ?? null; if (nextBlocker && task.blockedBy !== nextBlocker) { try {