diff --git a/.changeset/fn-8945-lineage-approval-invalidation.md b/.changeset/fn-8945-lineage-approval-invalidation.md new file mode 100644 index 0000000000..0068c0d867 --- /dev/null +++ b/.changeset/fn-8945-lineage-approval-invalidation.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Keep approved plans accurate when parent task lineage is removed. +category: fix +dev: Parent delete and archive now invalidate approved lineage evidence atomically. diff --git a/docs/architecture.md b/docs/architecture.md index 62be8c05cf..de7f829daa 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -2385,6 +2385,12 @@ An accepted `PROMPT.md` is preserved as an immutable, project-scoped spec-lock. Each authoritative prompt write reaches `PROMPT.md` before its task-row approval invalidation, then serializes evidence capture and reconciliation under the planning lifecycle lock. Current-plan evidence structurally binds the live task dependency set plus mission, slice, and parent-task lineage IDs as well as the prompt sections, so a dependency or lineage mutation appends comparable evidence even when `PROMPT.md` bytes are unchanged. Manual approval and a successful workflow-graph Plan Review create/reuse the lock before publishing acceptance evidence consumed by scheduling. If evidence persistence is interrupted after the file/row write, the inactive task remains safely unreleased; startup or event reconciliation re-reads `PROMPT.md`, appends the missing comparable revision, and reports drift. Re-approval reuses identical content or appends a lock with a prior-version diff; a clean later lock after an earlier retained divergence reports `diverged-relocked-approved` rather than losing that history. Drift reports compare the latest lock with current evidence and execution file paths, producing only machine-observable `scope-creep`, `silent-expansion`, or `plan-deviation` findings. Added File Scope boundaries, dependencies, lineage, durable steps, and acceptance criteria are silent expansion; reordered steps and changed Mission hashes are plan deviation. Reports are retained independently of live task rows, so archive, cleanup, and unarchive retain operator history. +Parent delete/archive with `removeLineageReferences:true` uses the same contract. It enters the invalidation boundary based on that flag even when its pre-read set is empty: the empty fast path skips prompt reads, probing, and locks but still revalidates inside the transaction, so a raced-in child aborts and retries once under its planning lock. One archive-lane value is resolved per operation and threaded through candidate reads and revalidation; delete retains resolved project archive lanes while archive retains its legacy `archived` lane semantics. + +Each attempt orders candidates → prompt map → structural transport probe → child locks → (archive workspace reservation → transaction → outcome → reconciliation). The transaction clears each affected child's parent binding and approved fingerprint, then resolves evidence from the pre-read prompt before committing. Evidence conflicts are classified by the matching `source_hash` row at any version (never only the latest row), and a reused row is accepted only after its canonical lineage is verified not to name the removed parent; an unresolvable or untruthful append raises `LineageEvidenceAppendError` and rolls the deletion/archive back. The only evidence-less commit is a child without a readable `PROMPT.md`; append-only history remains, so its report fences pre-existing, possibly stale evidence but cannot claim approved alignment after the fingerprint clear. The archive path pre-resolves prompts because filesystem reads never occur in its transaction. + +A session-scoped multi-key advisory lock covers every affected child, acquired in sorted order and nested under the parent task lock. Its callback publishes post-commit drift reconciliation before releasing those locks; lock acquisition failure occurs before archive workspace reservation, while an in-body failure releases its prepared reservation. Candidate revalidation is one retry only, then raises `LineageInvalidationCandidateRaceError`; no-clear gate rejection or lost delete claim publishes nothing. Structural pooler-only deployments may degrade serialization only after a pre-flight probe, but lock contention never degrades into an unlocked clear. Both channels reconcile through lock-free `reconcileSpecDriftWhilePlanningLocked`; report publication is the sole best-effort part. The transaction's resolved evidence version can differ from the report's `evaluateSpecDrift`-fenced latest evidence version; cleared approval, rather than evidence history removal, makes the report truthful. + The evaluator fences the latest lock version, current-plan evidence version/hash, approval fingerprint, and execution fingerprint before insertion. A stale candidate is re-evaluated from a fresh snapshot; project-engine task create, update, and move events enqueue a coalesced fresh comparison, while failed report persistence schedules a retry and shutdown cancels queued work and timers. Read APIs derive `activeLock` from the live fingerprint and current-plan hash, and return a current report only when its lock/current-plan/approval identities match the latest snapshots; older reports remain history rather than falsely showing `on-plan` after a dependency or lineage invalidation. Reports and run-audit events store hashes, IDs, and fixed outcomes only, never prompt or Mission prose. Drift is neither an admission nor quality/completion decision and does not move or fail work. ## Durable intake executor ownership diff --git a/packages/core/src/__tests__/postgres/lineage-approval-invalidation.pg.test.ts b/packages/core/src/__tests__/postgres/lineage-approval-invalidation.pg.test.ts new file mode 100644 index 0000000000..25e0b5cd64 --- /dev/null +++ b/packages/core/src/__tests__/postgres/lineage-approval-invalidation.pg.test.ts @@ -0,0 +1,379 @@ +/* +FNXC:SpecLockLineageInvalidation 2026-08-10-14:47: +Parent removal changes a child's canonical lineage binding. The regression covers both lifecycle +surfaces so neither may retain an active approval or report approved alignment after clearing it. +`evaluateSpecDrift` fences reports to the latest current-plan evidence; source hashes include the +canonical plan bindings, and listCurrentPlanEvidence is the row-enumeration API for this harness. + +FNXC:SpecLockLineageInvalidation 2026-08-10-15:37: +Evidence-conflict coverage must exercise both legal idempotence shapes: a matching hash may be +latest or historical, so classification must query that hash rather than inspect the latest row. +A deliberately malformed matching snapshot also proves that hash identity never substitutes for +verifying the stored lineage before committing a parent removal. +*/ +import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest"; +import { rm, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { + createSharedPgTaskStoreTestHarness, + pgDescribe, + type SharedPgTaskStoreHarness, +} from "../../__test-utils__/pg-test-harness.js"; +import { storeLog, type TaskStore } from "../../store.js"; +import { PlanningLifecycleLockTransportError } from "../../postgres/advisory-locks.js"; +import { LineageEvidenceAppendError, type LineageInvalidationTestEvent } from "../../task-store/lineage-approval-invalidation.js"; +import { createCurrentPlanEvidence } from "../../planner/spec-lock.js"; + +const SPEC_LOCK_PROMPT = `# Task + +## Mission + +Keep lineage scope observable. + +## File Scope + +- packages/core/src/task-store/async/async-lifecycle.ts + +## Steps + +1. Preserve evidence + +## Completion Criteria + +- [ ] Evidence is retained + +## Do NOT + +- Hide lineage changes + +## Dependencies + +- None +`; + +pgDescribe("lineage approval invalidation", () => { + const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_lineage_approval" }); + let store: TaskStore; + + beforeAll(h.beforeAll); + afterAll(h.afterAll); + beforeEach(async () => { + await h.beforeEach(); + store = h.store(); + }); + afterEach(h.afterEach); + + async function createApprovedLineagePair() { + const parent = await store.createTask({ description: "parent scope" }); + const child = await store.createTask({ + description: "approved child scope", + source: { sourceType: "api", sourceParentTaskId: parent.id }, + }); + await writeFile(join(store.taskDir(child.id), "PROMPT.md"), SPEC_LOCK_PROMPT); + const lock = await store.lockCurrentPlan(child.id, "approved-lineage-scope", SPEC_LOCK_PROMPT); + await store.updateTask(child.id, { approvedPlanFingerprint: "approved-lineage-scope" }); + return { parent, child, lock }; + } + + it("invalidates approval and publishes a divergent report when parent lineage is removed by delete or archive", async () => { + for (const operation of ["delete", "archive"] as const) { + const { parent, child, lock } = await createApprovedLineagePair(); + expect(lock.plan.sections.lineage.canonical).toContain(`parent-task:${parent.id}`); + expect(await store.getActiveSpecLock(child.id)).toBeDefined(); + + if (operation === "delete") { + await store.deleteTask(parent.id, { removeLineageReferences: true }); + } else { + await store.archiveTask(parent.id, { cleanup: false, removeLineageReferences: true }); + } + + const [updated, activeLock, latestLock, evidence, report] = await Promise.all([ + store.getTask(child.id), + store.getActiveSpecLock(child.id), + store.getLatestSpecLock(child.id), + store.getLatestCurrentPlanEvidence(child.id), + store.getLatestSpecDriftReport(child.id), + ]); + expect(updated?.sourceParentTaskId).toBeUndefined(); + expect(updated?.approvedPlanFingerprint).toBeUndefined(); + expect(activeLock).toBeUndefined(); + expect(latestLock?.version).toBe(lock.version); + expect(evidence?.version).toBeGreaterThan(lock.currentPlanVersion); + expect(evidence?.plan.sections.lineage.canonical).not.toContain(`parent-task:${parent.id}`); + expect(report).toEqual(expect.objectContaining({ + alignment: "diverged-needs-review", + currentPlanVersion: evidence?.version, + })); + expect(report?.approvedPlanFingerprint).toBeUndefined(); + expect(report?.alignment).not.toBe("on-plan"); + } + }); + + it("warns when a missing prompt is the sanctioned evidence-unavailable input", async () => { + const { parent, child } = await createApprovedLineagePair(); + const events: LineageInvalidationTestEvent[] = []; + const warn = vi.spyOn(storeLog, "warn").mockImplementation(() => undefined); + (store as unknown as { __lineageInvalidationForTest?: { onEvent: (event: LineageInvalidationTestEvent) => void } }).__lineageInvalidationForTest = { + onEvent: (event) => events.push(event), + }; + await rm(join(store.taskDir(child.id), "PROMPT.md")); + + await store.deleteTask(parent.id, { removeLineageReferences: true }); + + expect((await store.getTask(child.id))?.approvedPlanFingerprint).toBeUndefined(); + expect(events.find((event) => event.kind === "outcome")).toMatchObject({ + kind: "outcome", outcome: { evidenceUnavailableChildIds: [child.id], evidenceInsertAttempts: 0 }, + }); + expect(warn).toHaveBeenCalledWith(expect.stringContaining(child.id)); + warn.mockRestore(); + }); + + async function appendPostClearEvidence(childId: string, version: number, prompt = SPEC_LOCK_PROMPT) { + const evidence = createCurrentPlanEvidence({ + version, + sourceRevision: Date.now(), + capturedAt: new Date().toISOString(), + prompt, + }); + return await store.appendCurrentPlanEvidence(childId, evidence); + } + + it("reuses a verified latest matching source-hash evidence row without duplicating it", async () => { + const { parent, child } = await createApprovedLineagePair(); + const events: LineageInvalidationTestEvent[] = []; + const existing = await store.listCurrentPlanEvidence(child.id); + const matching = await appendPostClearEvidence(child.id, existing.at(-1)!.version + 1); + const before = await store.listCurrentPlanEvidence(child.id); + expect(matching.plan.sections.lineage.canonical).not.toContain(`parent-task:${parent.id}`); + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + onEvent: (event: LineageInvalidationTestEvent) => events.push(event), + }; + + await store.deleteTask(parent.id, { removeLineageReferences: true }); + + expect(await store.listCurrentPlanEvidence(child.id)).toEqual(before); + expect(events.find((event) => event.kind === "outcome")).toMatchObject({ + kind: "outcome", outcome: { evidenceVersionByChild: new Map([[child.id, matching.version]]), evidenceInsertAttempts: 1 }, + }); + }); + + it("resolves a verified stale matching hash instead of comparing only the latest evidence", async () => { + const { parent, child } = await createApprovedLineagePair(); + const events: LineageInvalidationTestEvent[] = []; + const original = await store.listCurrentPlanEvidence(child.id); + const matching = await appendPostClearEvidence(child.id, original.at(-1)!.version + 1); + const newer = await appendPostClearEvidence(child.id, matching.version + 1, `${SPEC_LOCK_PROMPT}\n\n## Later revision\nDifferent evidence.`); + const before = await store.listCurrentPlanEvidence(child.id); + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + onEvent: (event: LineageInvalidationTestEvent) => events.push(event), + }; + + await store.deleteTask(parent.id, { removeLineageReferences: true }); + + expect(await store.listCurrentPlanEvidence(child.id)).toEqual(before); + const outcome = events.find((event) => event.kind === "outcome"); + expect(outcome).toMatchObject({ + kind: "outcome", outcome: { evidenceVersionByChild: new Map([[child.id, matching.version]]), evidenceInsertAttempts: 1 }, + }); + const report = await store.getLatestSpecDriftReport(child.id); + // The report fences to the latest row; the durable idempotence row may be older. + expect(report?.currentPlanVersion).toBe(newer.version); + expect(report?.alignment).toBe("diverged-needs-review"); + }); + + it("rolls back an untruthful matching-hash row instead of accepting it by hash alone", async () => { + const { parent, child } = await createApprovedLineagePair(); + const existing = await store.listCurrentPlanEvidence(child.id); + const postClearPrompt = `${SPEC_LOCK_PROMPT}\n\n## Truthfulness fixture\nUnique source identity.`; + await writeFile(join(store.taskDir(child.id), "PROMPT.md"), postClearPrompt); + const truthful = createCurrentPlanEvidence({ + version: existing.at(-1)!.version + 1, + sourceRevision: Date.now(), + capturedAt: new Date().toISOString(), + prompt: postClearPrompt, + }); + const untruthful = { + ...truthful, + plan: { + ...truthful.plan, + sections: { + ...truthful.plan.sections, + lineage: { canonical: `parent-task:${parent.id}` }, + }, + }, + }; + await store.appendCurrentPlanEvidence(child.id, untruthful); + const before = await store.listCurrentPlanEvidence(child.id); + + await expect(store.deleteTask(parent.id, { removeLineageReferences: true })) + .rejects.toMatchObject({ name: "LineageEvidenceAppendError", reason: "matched-row-not-truthful" } satisfies Partial); + expect((await store.getTask(parent.id))?.deletedAt).toBeUndefined(); + expect((await store.getTask(child.id))?.sourceParentTaskId).toBe(parent.id); + expect((await store.getTask(child.id))?.approvedPlanFingerprint).toBe("approved-lineage-scope"); + expect(await store.listCurrentPlanEvidence(child.id)).toEqual(before); + }); + + it("retries a forced evidence version collision and rolls back an unresolvable evidence append", async () => { + const collisionVersion = 99; + const events: LineageInvalidationTestEvent[] = []; + const { parent, child } = await createApprovedLineagePair(); + // A disk-only prompt revision supplies a post-clear hash absent from approval-time evidence. + await writeFile(join(store.taskDir(child.id), "PROMPT.md"), `${SPEC_LOCK_PROMPT}\n\n## Collision target\nFresh lineage evidence.`); + await store.appendCurrentPlanEvidence(child.id, createCurrentPlanEvidence({ + version: collisionVersion, + sourceRevision: Date.now(), + capturedAt: new Date().toISOString(), + prompt: `${SPEC_LOCK_PROMPT}\n\n## Collision fixture\nDifferent source hash.`, + })); + const beforeRetry = await store.listCurrentPlanEvidence(child.id); + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + evidenceTargetVersionForTest: (_childId: string, computed: number, attempt: number) => attempt === 0 ? collisionVersion : computed, + onEvent: (event: LineageInvalidationTestEvent) => events.push(event), + }; + + await store.deleteTask(parent.id, { removeLineageReferences: true }); + + const afterRetry = await store.listCurrentPlanEvidence(child.id); + expect(afterRetry).toHaveLength(beforeRetry.length + 1); + expect(afterRetry.at(-1)?.version).toBeGreaterThan(collisionVersion); + expect(events.find((event) => event.kind === "outcome")).toMatchObject({ + kind: "outcome", outcome: { evidenceInsertAttempts: 2 }, + }); + + const abortPair = await createApprovedLineagePair(); + await writeFile(join(store.taskDir(abortPair.child.id), "PROMPT.md"), `${SPEC_LOCK_PROMPT}\n\n## Abort target\nFresh lineage evidence.`); + await store.appendCurrentPlanEvidence(abortPair.child.id, createCurrentPlanEvidence({ + version: collisionVersion, + sourceRevision: Date.now(), + capturedAt: new Date().toISOString(), + prompt: `${SPEC_LOCK_PROMPT}\n\n## Permanent collision fixture\nDifferent source hash.`, + })); + const beforeAbort = await store.listCurrentPlanEvidence(abortPair.child.id); + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + evidenceTargetVersionForTest: () => collisionVersion, + }; + + await expect(store.archiveTask(abortPair.parent.id, { cleanup: false, removeLineageReferences: true })) + .rejects.toMatchObject({ name: "LineageEvidenceAppendError", reason: "no-durable-version" } satisfies Partial); + expect((await store.getTask(abortPair.parent.id))?.deletedAt).toBeUndefined(); + expect((await store.getTask(abortPair.child.id))?.sourceParentTaskId).toBe(abortPair.parent.id); + expect((await store.getTask(abortPair.child.id))?.approvedPlanFingerprint).toBe("approved-lineage-scope"); + expect(await store.listCurrentPlanEvidence(abortPair.child.id)).toEqual(beforeAbort); + }); + + it("retries a raced-in delete child and publishes before releasing all child locks", async () => { + const { parent, child } = await createApprovedLineagePair(); + const events: LineageInvalidationTestEvent[] = []; + let secondChildId: string | undefined; + (store as unknown as { + __lineageInvalidationForTest?: { + availability: () => { available: true }; + withLocks: (ids: readonly string[], callback: () => Promise) => Promise; + onEvent: (event: LineageInvalidationTestEvent) => void; + }; + __beforeDeleteClaimForTest?: () => Promise; + }).__lineageInvalidationForTest = { + availability: () => ({ available: true }), + async withLocks(_ids, callback) { + return await callback(); + }, + onEvent: (event) => events.push(event), + }; + (store as unknown as { __beforeDeleteClaimForTest?: () => Promise }).__beforeDeleteClaimForTest = async () => { + const second = await store.createTask({ + description: "raced-in approved child", + source: { sourceType: "api", sourceParentTaskId: parent.id }, + }); + secondChildId = second.id; + await writeFile(join(store.taskDir(second.id), "PROMPT.md"), SPEC_LOCK_PROMPT); + await store.lockCurrentPlan(second.id, "raced-in-lineage-scope", SPEC_LOCK_PROMPT); + await store.updateTask(second.id, { approvedPlanFingerprint: "raced-in-lineage-scope" }); + }; + + await store.deleteTask(parent.id, { removeLineageReferences: true }); + + expect(secondChildId).toBeDefined(); + for (const id of [child.id, secondChildId!]) { + expect((await store.getTask(id))?.sourceParentTaskId).toBeUndefined(); + expect((await store.getTask(id))?.approvedPlanFingerprint).toBeUndefined(); + expect(await store.getActiveSpecLock(id)).toBeUndefined(); + } + const acquisitions = events.filter((event) => event.kind === "acquire"); + expect(acquisitions).toHaveLength(2); + expect(acquisitions[1]?.childIds).toEqual([child.id, secondChildId].sort()); + const reconciles = events.filter((event) => event.kind === "reconcile"); + expect(reconciles.map((event) => event.kind === "reconcile" ? event.childId : undefined).sort()) + .toEqual([child.id, secondChildId].sort()); + }); + + it("retries an archive child discovered inside its transaction and publishes before releasing its lock", async () => { + const parent = await store.createTask({ description: "archive lineage parent" }); + const events: LineageInvalidationTestEvent[] = []; + let childId: string | undefined; + let injected = false; + (store as unknown as { + __lineageInvalidationForTest?: { onEvent: (event: LineageInvalidationTestEvent) => void }; + __beforeArchiveLineageGateForTest?: () => Promise; + }).__lineageInvalidationForTest = { onEvent: (event) => events.push(event) }; + (store as unknown as { __beforeArchiveLineageGateForTest?: () => Promise }).__beforeArchiveLineageGateForTest = async () => { + if (injected) return; + injected = true; + const child = await store.createTask({ + description: "archive raced-in approved child", + source: { sourceType: "api", sourceParentTaskId: parent.id }, + }); + childId = child.id; + await writeFile(join(store.taskDir(child.id), "PROMPT.md"), SPEC_LOCK_PROMPT); + await store.lockCurrentPlan(child.id, "archive-raced-lineage-scope", SPEC_LOCK_PROMPT); + await store.updateTask(child.id, { approvedPlanFingerprint: "archive-raced-lineage-scope" }); + }; + + await store.archiveTask(parent.id, { cleanup: false, removeLineageReferences: true }); + + expect(childId).toBeDefined(); + const child = await store.getTask(childId!); + expect(child?.sourceParentTaskId).toBeUndefined(); + expect(child?.approvedPlanFingerprint).toBeUndefined(); + expect(await store.getActiveSpecLock(childId!)).toBeUndefined(); + const outcomes = events.filter((event) => event.kind === "outcome"); + expect(outcomes).toHaveLength(2); + expect(outcomes[0]).toMatchObject({ kind: "outcome", outcome: { candidateIds: [], clearedChildIds: [], evidenceInsertAttempts: 0 } }); + expect(outcomes[1]).toMatchObject({ kind: "outcome", outcome: { candidateIds: [childId], clearedChildIds: [childId] } }); + const acquire = events.findIndex((event) => event.kind === "acquire" && event.childIds.includes(childId!)); + const reconcile = events.findIndex((event) => event.kind === "reconcile" && event.childId === childId); + const release = events.findIndex((event) => event.kind === "release" && event.childIds.includes(childId!)); + expect(acquire).toBeGreaterThan(-1); + expect(reconcile).toBeGreaterThan(acquire); + expect(release).toBeGreaterThan(reconcile); + }); + + it("records lifecycle outcomes and uses structural degradation but never turns a lock acquisition failure into an unlocked delete", async () => { + const { parent, child } = await createApprovedLineagePair(); + const degradedEvents: LineageInvalidationTestEvent[] = []; + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + availability: () => ({ available: false, reason: "direct-session-unavailable" }), + onEvent: (event: LineageInvalidationTestEvent) => degradedEvents.push(event), + }; + await store.deleteTask(parent.id, { removeLineageReferences: true }); + expect((await store.getTask(child.id))?.approvedPlanFingerprint).toBeUndefined(); + expect(degradedEvents.some((event) => event.kind === "run" && event.degraded)).toBe(true); + expect(degradedEvents.some((event) => event.kind === "acquire")).toBe(false); + expect(degradedEvents.some((event) => event.kind === "reconcile" && event.childId === child.id)).toBe(true); + + await h.beforeEach(); + store = h.store(); + const pair = await createApprovedLineagePair(); + let ranBody = false; + (store as unknown as { __lineageInvalidationForTest?: unknown }).__lineageInvalidationForTest = { + availability: () => ({ available: true }), + withLocks: async () => { throw new PlanningLifecycleLockTransportError("contended"); }, + onEvent: (event: LineageInvalidationTestEvent) => { if (event.kind === "run") ranBody = true; }, + }; + await expect(store.deleteTask(pair.parent.id, { removeLineageReferences: true })) + .rejects.toBeInstanceOf(PlanningLifecycleLockTransportError); + expect(ranBody).toBe(false); + expect((await store.getTask(pair.parent.id))?.deletedAt).toBeUndefined(); + expect((await store.getTask(pair.child.id))?.sourceParentTaskId).toBe(pair.parent.id); + expect((await store.getTask(pair.child.id))?.approvedPlanFingerprint).toBe("approved-lineage-scope"); + }); +}); diff --git a/packages/core/src/__tests__/postgres/planning-lifecycle-advisory-lock.pg.test.ts b/packages/core/src/__tests__/postgres/planning-lifecycle-advisory-lock.pg.test.ts index a463b69772..bd9e29fce1 100644 --- a/packages/core/src/__tests__/postgres/planning-lifecycle-advisory-lock.pg.test.ts +++ b/packages/core/src/__tests__/postgres/planning-lifecycle-advisory-lock.pg.test.ts @@ -7,7 +7,9 @@ import { } from "../../__test-utils__/pg-test-harness.js"; import { PlanningLifecycleLockTransportError, + planningLifecycleLockTransportAvailability, withPlanningLifecycleAdvisoryLock, + withPlanningLifecycleAdvisoryLocks, } from "../../postgres/advisory-locks.js"; import { resolveBackendWithOptions } from "../../postgres/backend-resolver.js"; @@ -69,6 +71,62 @@ pgDescribe("planning lifecycle advisory lock", () => { expect(order).toEqual(["first-enter", "first-exit", "second-enter"]); }); + it("shares key space with a single-key holder across a sorted multi-key session", async () => { + let releaseMulti!: () => void; + const multiCanFinish = new Promise((resolve) => { releaseMulti = resolve; }); + let multiEntered!: () => void; + const multiIsHolding = new Promise((resolve) => { multiEntered = resolve; }); + const multi = withPlanningLifecycleAdvisoryLocks({ + projectId: "project-a", + taskIds: ["FN-2", "FN-1", "FN-2"], + directSessionUrl: h.testUrl(), + provenance: "migration-override", + runtimeUrl: h.testUrl(), + migrationUrl: h.testUrl(), + timeoutMs: 1_000, + }, async () => { + multiEntered(); + await multiCanFinish; + }); + await multiIsHolding; + + let singleAttempted!: () => void; + const singleDispatched = new Promise((resolve) => { singleAttempted = resolve; }); + const single = lock("project-a", "FN-1", async () => {}, 1_000, singleAttempted); + await singleDispatched; + releaseMulti(); + await Promise.all([multi, single]); + }); + + it("releases every multi-key lock when its callback fails", async () => { + await expect(withPlanningLifecycleAdvisoryLocks({ + projectId: "project-a", + taskIds: ["FN-1", "FN-2"], + directSessionUrl: h.testUrl(), + provenance: "migration-override", + runtimeUrl: h.testUrl(), + migrationUrl: h.testUrl(), + }, async () => { throw new Error("callback failed"); })).rejects.toThrow("callback failed"); + + await Promise.all([ + lock("project-a", "FN-1", async () => {}), + lock("project-a", "FN-2", async () => {}), + ]); + }); + + it("shares the structural availability predicate used by acquisition", () => { + expect(planningLifecycleLockTransportAvailability({ + directSessionUrl: null, + provenance: null, + })).toEqual({ available: false, reason: "direct-session-unavailable" }); + expect(planningLifecycleLockTransportAvailability({ + directSessionUrl: h.testUrl(), + provenance: "migration-override", + runtimeUrl: h.testUrl(), + migrationUrl: h.testUrl(), + })).toEqual({ available: true }); + }); + it("does not contend across project or task keys", async () => { let releaseFirst!: () => void; const firstCanFinish = new Promise((resolve) => { releaseFirst = resolve; }); diff --git a/packages/core/src/__tests__/postgres/runtime-lifecycle-async.test.ts b/packages/core/src/__tests__/postgres/runtime-lifecycle-async.test.ts index d7030d1d75..bdaaade61c 100644 --- a/packages/core/src/__tests__/postgres/runtime-lifecycle-async.test.ts +++ b/packages/core/src/__tests__/postgres/runtime-lifecycle-async.test.ts @@ -269,6 +269,7 @@ pgDescribe("runtime-lifecycle-async: deleteTask lineage gate (PostgreSQL)", () = createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(), sourceParentTaskId: "FN-PARENT2", + approvedPlanFingerprint: "approved-lineage-scope", status: null, } as never, { lineageId: "test" }, @@ -284,5 +285,8 @@ pgDescribe("runtime-lifecycle-async: deleteTask lineage gate (PostgreSQL)", () = .where(eq(schema.project.tasks.id, "FN-PARENT2")); expect(rows.length).toBe(1); expect(rows[0].deletedAt).not.toBeNull(); + const child = await h.store.getTask("FN-CHILD2"); + expect(child.sourceParentTaskId).toBeUndefined(); + expect(child.approvedPlanFingerprint).toBeUndefined(); }); }); diff --git a/packages/core/src/__tests__/postgres/taskstore-lifecycle.test.ts b/packages/core/src/__tests__/postgres/taskstore-lifecycle.test.ts index 67a900b164..1b09288f45 100644 --- a/packages/core/src/__tests__/postgres/taskstore-lifecycle.test.ts +++ b/packages/core/src/__tests__/postgres/taskstore-lifecycle.test.ts @@ -219,7 +219,8 @@ pgDescribe("U13 taskstore-lifecycle (PostgreSQL)", () => { const childIds = await findLiveLineageChildren(tx, "KB-PARENT"); expect(childIds).toEqual(["KB-CHILD"]); const cleared = await removeLineageReferences(tx, "KB-PARENT", childIds, nowIso); - expect(cleared).toBe(1); + expect(cleared.clearedChildIds).toEqual(["KB-CHILD"]); + expect(cleared.evidenceUnavailableChildIds).toEqual(["KB-CHILD"]); }); // After: gate passes (no live children). diff --git a/packages/core/src/postgres/advisory-locks.ts b/packages/core/src/postgres/advisory-locks.ts index 6909943706..0de4658213 100644 --- a/packages/core/src/postgres/advisory-locks.ts +++ b/packages/core/src/postgres/advisory-locks.ts @@ -30,9 +30,40 @@ export class PlanningLifecycleLockTransportError extends Error { } } -const DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS = 5_000; +export const DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS = 5_000; type DedicatedPostgresClient = ReturnType; +type PlanningLifecycleLockInput = { + projectId: string; + taskId: string; + directSessionUrl: string | null; + provenance: "embedded-lifecycle" | "migration-override" | "runtime-direct" | null; + runtimeUrl?: string | null; + migrationUrl?: string | null; + timeoutMs?: number; + onLockAcquisitionAttempt?: () => void; +}; + +/** Side-effect-free structural check shared by the lineage degrade path and lock acquisition. */ +export function planningLifecycleLockTransportAvailability(input: Pick): { available: true } | { available: false; reason: string } { + const directUrl = input.directSessionUrl; + if (!directUrl || !input.provenance || looksLikePoolerUrl(directUrl)) return { available: false, reason: "direct-session-unavailable" }; + let directDatabase: string; + try { directDatabase = decodeURIComponent(new URL(directUrl).pathname.replace(/^\//, "")); } catch { return { available: false, reason: "direct-session-invalid" }; } + if (!directDatabase) return { available: false, reason: "direct-session-database-missing" }; + const descriptorEndpoint = input.provenance === "embedded-lifecycle" || input.provenance === "runtime-direct" ? input.runtimeUrl : input.migrationUrl; + if (!descriptorEndpoint || descriptorEndpoint !== directUrl) return { available: false, reason: "backend-descriptor-mismatch" }; + const operationalUrl = input.runtimeUrl ?? input.migrationUrl; + if (operationalUrl) { + try { + const operationalDatabase = decodeURIComponent(new URL(operationalUrl).pathname.replace(/^\//, "")); + if (!operationalDatabase || operationalDatabase !== directDatabase) return { available: false, reason: "backend-database-mismatch" }; + } catch { return { available: false, reason: "backend-endpoint-invalid" }; } + } + return { available: true }; +} + +const planningLifecycleLockKey = (projectId: string, taskId: string): string => `fusion:planning-lifecycle:${projectId}:${taskId}`; type PostgresBackendIdentity = { database: string; host: string | null; @@ -85,53 +116,17 @@ async function runBoundedTransportPhase( * lock and unlock and silently defeat session advisory locking. */ export async function withPlanningLifecycleAdvisoryLock( - input: { - projectId: string; - taskId: string; - directSessionUrl: string | null; - provenance: "embedded-lifecycle" | "migration-override" | "runtime-direct" | null; - runtimeUrl?: string | null; - migrationUrl?: string | null; - /** Bounds dedicated-session setup and lock acquisition, not callback work. */ - timeoutMs?: number; - /** @internal Test seam fired when the driver dispatches the lock query. */ - onLockAcquisitionAttempt?: () => void; - }, + input: PlanningLifecycleLockInput, callback: () => Promise, ): Promise { - const directUrl = input.directSessionUrl; - const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS); - if (!directUrl || !input.provenance || looksLikePoolerUrl(directUrl)) { - throw new PlanningLifecycleLockTransportError("Planning lifecycle lock requires a direct PostgreSQL session endpoint; set DATABASE_MIGRATION_URL to a direct, non-pooled connection"); - } - - let directDatabase: string; - try { - directDatabase = decodeURIComponent(new URL(directUrl).pathname.replace(/^\//, "")); - } catch { - throw new PlanningLifecycleLockTransportError("Planning lifecycle lock has an invalid direct PostgreSQL session endpoint"); - } - if (!directDatabase) { - throw new PlanningLifecycleLockTransportError("Planning lifecycle lock direct endpoint must select a database"); - } - const descriptorEndpoint = input.provenance === "embedded-lifecycle" || input.provenance === "runtime-direct" - ? input.runtimeUrl - : input.migrationUrl; - if (!descriptorEndpoint || descriptorEndpoint !== directUrl) { - throw new PlanningLifecycleLockTransportError("Planning lifecycle lock endpoint does not match the resolved backend descriptor"); + const availability = planningLifecycleLockTransportAvailability(input); + if (!availability.available) { + throw new PlanningLifecycleLockTransportError(`Planning lifecycle lock transport unavailable: ${availability.reason}`); } + const directUrl = input.directSessionUrl!; + const directDatabase = decodeURIComponent(new URL(directUrl).pathname.replace(/^\//, "")); const operationalUrl = input.runtimeUrl ?? input.migrationUrl; - if (operationalUrl) { - try { - const operationalDatabase = decodeURIComponent(new URL(operationalUrl).pathname.replace(/^\//, "")); - if (!operationalDatabase || operationalDatabase !== directDatabase) { - throw new PlanningLifecycleLockTransportError("Planning lifecycle lock endpoint selects a different database than the resolved backend"); - } - } catch (error) { - if (error instanceof PlanningLifecycleLockTransportError) throw error; - throw new PlanningLifecycleLockTransportError("Resolved backend has an invalid PostgreSQL endpoint"); - } - } + const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS); const client = postgres(directUrl, { max: 1, @@ -144,7 +139,7 @@ export async function withPlanningLifecycleAdvisoryLock( } : undefined, }); - const key = `fusion:planning-lifecycle:${input.projectId}:${input.taskId}`; + const key = planningLifecycleLockKey(input.projectId, input.taskId); let acquired = false; try { /* @@ -231,3 +226,82 @@ export async function withPlanningLifecycleAdvisoryLock( } } } + +/** + * FNXC:SpecLockLineageInvalidation 2026-08-10-14:33: + * Parent removal can affect many children, but no child may be cleared outside its planning lock. + * Advisory locks are session-scoped, so this acquires every sorted task key on one dedicated session + * instead of trading correctness for a fan-out cap or opening one connection per child. + */ +export async function withPlanningLifecycleAdvisoryLocks( + input: Omit & { taskIds: readonly string[] }, + callback: () => Promise, +): Promise { + const taskIds = [...new Set(input.taskIds)].sort(); + if (taskIds.length === 0) return callback(); + const availability = planningLifecycleLockTransportAvailability(input); + if (!availability.available) throw new PlanningLifecycleLockTransportError(`Planning lifecycle lock transport unavailable: ${availability.reason}`); + const directUrl = input.directSessionUrl!; + const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS); + const directDatabase = decodeURIComponent(new URL(directUrl).pathname.replace(/^\//, "")); + const operationalUrl = input.runtimeUrl ?? input.migrationUrl; + const client = postgres(directUrl, { + max: 1, connect_timeout: Math.max(1, Math.ceil(timeoutMs / 1_000)), prepare: false, onnotice: () => {}, + debug: input.onLockAcquisitionAttempt ? (_connection, query) => { if (query.includes("pg_advisory_lock(")) input.onLockAcquisitionAttempt?.(); } : undefined, + }); + const keys = taskIds.map((taskId) => planningLifecycleLockKey(input.projectId, taskId)); + const acquired: string[] = []; + try { + const identity = await runBoundedTransportPhase(client, timeoutMs, `Planning lifecycle lock session setup timed out after ${timeoutMs}ms`, "Planning lifecycle lock could not establish its dedicated session", async () => { + const rows = await readPostgresBackendIdentity(client); + await client`SELECT set_config('lock_timeout', ${`${timeoutMs}ms`}, false)`; + return rows; + }); + if (identity[0]?.database !== directDatabase) throw new PlanningLifecycleLockTransportError("Planning lifecycle lock session selected an unexpected database"); + + /* + FNXC:SpecLockLineageInvalidation 2026-08-10-14:47: + The multi-key session must prove the same resolved backend identity as the single-key lock. + Matching advisory key strings on a different PostgreSQL server would serialize nothing. + */ + if (operationalUrl && operationalUrl !== directUrl) { + const operationalClient = postgres(operationalUrl, { + max: 1, + connect_timeout: Math.max(1, Math.ceil(timeoutMs / 1_000)), + prepare: false, + onnotice: () => {}, + }); + try { + const operationalIdentity = await runBoundedTransportPhase( + operationalClient, + timeoutMs, + `Planning lifecycle lock backend identity check timed out after ${timeoutMs}ms`, + "Planning lifecycle lock could not verify the resolved backend server identity", + () => readPostgresBackendIdentity(operationalClient), + ); + const expected = operationalIdentity[0]; + const actual = identity[0]; + if (!expected || !actual + || expected.database !== actual.database + || expected.host !== actual.host + || expected.port !== actual.port + || (expected.cluster && actual.cluster && expected.cluster !== actual.cluster)) { + throw new PlanningLifecycleLockTransportError("Planning lifecycle lock endpoint targets a different PostgreSQL server than the resolved backend"); + } + } catch (error) { + if (error instanceof PlanningLifecycleLockTransportError) throw error; + throw new PlanningLifecycleLockTransportError("Planning lifecycle lock could not verify the resolved backend server identity"); + } finally { + await operationalClient.end({ timeout: 5 }).catch(() => undefined); + } + } + for (const key of keys) { + await runBoundedTransportPhase(client, timeoutMs, `Planning lifecycle lock acquisition timed out after ${timeoutMs}ms`, "Planning lifecycle lock acquisition failed", () => client`SELECT pg_advisory_lock(hashtext(${key}))`); + acquired.push(key); + } + return await callback(); + } finally { + for (const key of acquired.reverse()) await runBoundedTransportPhase(client, timeoutMs, `Planning lifecycle lock cleanup timed out after ${timeoutMs}ms`, "Planning lifecycle lock cleanup failed", () => client`SELECT pg_advisory_unlock(hashtext(${key}))`).catch(() => undefined); + await client.end({ timeout: 5 }).catch(() => undefined); + } +} diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index e2699cb813..0d5cf6a46d 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -78,7 +78,7 @@ import { GlobalSettingsStore } from "./config/global-settings.js"; import { Database } from "./db/db.js"; import { ArchiveDatabase } from "./db/archive-db.js"; import { projectScopeFor, type AsyncDataLayer, type DbTransaction } from "./postgres/data-layer.js"; -import { withPlanningLifecycleAdvisoryLock } from "./postgres/advisory-locks.js"; +import { planningLifecycleLockTransportAvailability, withPlanningLifecycleAdvisoryLock, withPlanningLifecycleAdvisoryLocks } from "./postgres/advisory-locks.js"; import { MissionStore } from "./missions/mission-store.js"; import { AsyncMissionStore } from "./async-stores/async-mission-store.js"; import { AsyncIdeationStore } from "./async-stores/async-ideation-store.js"; @@ -891,6 +891,30 @@ export class TaskStore extends EventEmitter { if (planningLifecycleLocks.get(key) === current) planningLifecycleLocks.delete(key); } } + /** Acquire all child planning keys in deterministic order without a fan-out connection cost. */ + public async withPlanningLifecycleLocks(ids: readonly string[], fn: () => Promise): Promise { + const sorted = [...new Set(ids)].sort(); + if (sorted.length === 0) return fn(); + const backend = this.asyncLayer?.backend; + if (this.asyncLayer && backend) { + return withPlanningLifecycleAdvisoryLocks({ + projectId: this.asyncLayer.projectId ?? this.rootDir, + taskIds: sorted, + directSessionUrl: backend.directSessionUrl ?? null, + provenance: backend.directSessionProvenance ?? null, + runtimeUrl: backend.runtimeUrl, + migrationUrl: backend.migrationUrl, + }, fn); + } + const acquire = async (index: number): Promise => index === sorted.length ? fn() : this.withPlanningLifecycleLock(sorted[index]!, () => acquire(index + 1)); + return acquire(0); + } + public planningLifecycleLockTransportAvailability(): { available: true } | { available: false; reason: string } { + const backend = this.asyncLayer?.backend; + if (!this.asyncLayer || !backend) return { available: true }; + return planningLifecycleLockTransportAvailability({ directSessionUrl: backend.directSessionUrl ?? null, provenance: backend.directSessionProvenance ?? null, runtimeUrl: backend.runtimeUrl, migrationUrl: backend.migrationUrl }); + } + public getTaskIdFromDir(dir: string): string { return getTaskIdFromDirImpl(this, dir); } diff --git a/packages/core/src/task-store/archive-lifecycle-2.ts b/packages/core/src/task-store/archive-lifecycle-2.ts index b06610c019..22bfccadde 100644 --- a/packages/core/src/task-store/archive-lifecycle-2.ts +++ b/packages/core/src/task-store/archive-lifecycle-2.ts @@ -26,7 +26,8 @@ import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {softDeleteTaskRowInTransaction, readTaskRow as readTaskRowAsync, readTaskRowInTransaction} from "../task-store/async/async-persistence.js"; import {appendTaskLifecycleEventInTransaction} from "../task-store/lifecycle-outbox.js"; -import {findLiveLineageChildren as findLiveLineageChildrenAsync, projectPartition, removeLineageReferences} from "../task-store/async/async-lifecycle.js"; +import {findLiveLineageChildren as findLiveLineageChildrenAsync, projectPartition, removeLineageReferences, type LineageRemovalOutcome} from "../task-store/async/async-lifecycle.js"; +import { classifyLineageInvalidationOutcomeError, lineageEvidenceTargetVersionForTest, recordLineageInvalidationOutcome, reconcileClearedLineageChildren, resolveAndAssertLineageCandidatesUnchanged, runLineageInvalidation } from "../task-store/lineage-approval-invalidation.js"; import { resolveProjectColumnsForRoles } from "../project-lane-vocabulary.js"; import {archiveParentTaskWithLineageGate, findArchivedTaskEntry, deleteArchivedTaskEntry, restoreTaskFromArchive} from "../task-store/async/async-archive-lineage.js"; import {getArchivedRowCount, listArchivedTaskEntriesPage} from "../async-stores/async-archive-db.js"; @@ -190,7 +191,11 @@ async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string await (store as unknown as { __beforeDeleteClaimForTest?: (taskId: string) => void | Promise }).__beforeDeleteClaimForTest?.(id); // Soft-delete + lineage clear + mission unlink + audit in one transaction (atomicity). - const deletion = await layer.transactionImmediate(async (tx) => { + const executeDelete = async (context?: { candidateIds: string[]; promptByChildId: ReadonlyMap; locksHeld: boolean; attempt: number }) => { + let deletion: DeleteTaskClaimResult & { lineageOutcome: LineageRemovalOutcome }; + try { + deletion = await layer.transactionImmediate(async (tx) => { + if (context) await resolveAndAssertLineageCandidatesUnchanged(tx, id, layer.projectId, lineageArchivedLanes, context.candidateIds); /* FNXC:LifecycleOutbox 2026-08-01-10:33: The pre-transaction deletedAt read is a cross-process TOCTOU window. A conditional @@ -201,12 +206,12 @@ async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string if (claimed === false) { const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId); if (!reloaded) throw new TaskNotFoundError(id); - return { claimed: false, task: store.rowToTask(store.pgRowToTaskRow(reloaded)) }; - } - // Clear lineage references on live children so the parent can be deleted. - if (lineageChildIds.length > 0) { - await removeLineageReferences(tx, id, lineageChildIds, deletedAt, layer.projectId); + return { claimed: false, task: store.rowToTask(store.pgRowToTaskRow(reloaded)), lineageOutcome: { clearedChildIds: [] as string[], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 } }; } + // Clear lineage references and approval only after locked candidates were revalidated. + const lineageOutcome = context + ? await removeLineageReferences(tx, id, context.candidateIds, deletedAt, layer.projectId, context.promptByChildId, lineageEvidenceTargetVersionForTest(store)) + : { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 }; /* FNXC:MissionStore 2026-07-17-17:40: Clear any mission feature→task link IN THIS TRANSACTION so it commits (or rolls @@ -281,8 +286,32 @@ async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string // callers so neither receives the pre-claim live snapshot after a successful delete. const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId); if (!reloaded) throw new TaskNotFoundError(id); - return { claimed: true, task: store.rowToTask(store.pgRowToTaskRow(reloaded)) }; + return { claimed: true, task: store.rowToTask(store.pgRowToTaskRow(reloaded)), lineageOutcome }; }); + } catch (error) { + if (context) recordLineageInvalidationOutcome(store, { + attempt: context.attempt, locksHeld: context.locksHeld, degraded: !context.locksHeld, + candidateIds: context.candidateIds, clearedChildIds: [], evidenceVersionByChild: new Map(), + evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0, + error: classifyLineageInvalidationOutcomeError(error), + }); + throw error; + } + if (context) { + recordLineageInvalidationOutcome(store, { + attempt: context.attempt, locksHeld: context.locksHeld, degraded: !context.locksHeld, + candidateIds: context.candidateIds, clearedChildIds: deletion.lineageOutcome.clearedChildIds, + evidenceVersionByChild: deletion.lineageOutcome.evidenceVersionByChild, + evidenceUnavailableChildIds: deletion.lineageOutcome.evidenceUnavailableChildIds, + evidenceInsertAttempts: deletion.lineageOutcome.evidenceInsertAttempts, + }); + await reconcileClearedLineageChildren(store, deletion.lineageOutcome.clearedChildIds, { locksHeld: context.locksHeld }); + } + return deletion; + }; + const deletion = options?.removeLineageReferences + ? await runLineageInvalidation(store, id, { archivedColumns: lineageArchivedLanes, initialCandidateIds: lineageChildIds }, executeDelete) + : await executeDelete(); if (!deletion.claimed) return deletion; @@ -396,31 +425,61 @@ export async function archiveTaskBackendImpl(store: TaskStore, id: string, optio const entry = await store.taskToArchiveEntry(task, archivedAt); /* - FNXC:WorkflowLifecycle 2026-07-16-15:30: - Backend archive persists cold storage before its cleanup phase. Hold the - per-repository reservations across that transaction so another process sees - the path as unavailable until the awaited workspace disposer has removed it. + FNXC:SpecLockLineageInvalidation 2026-08-10-14:33: + Archive keeps its legacy archived-lane semantics (undefined -> "archived") but threads that + one value through pre-read and gate. The workspace reservation is created inside the locked body. */ - const preparedWorkspace = cleanup ? await prepareArchivedWorkspaceWorktrees(store, task) : undefined; - let result; - try { - // Lineage gate + archive in one transaction. - result = await archiveParentTaskWithLineageGate(layer, id, entry, { - removeLineageReferences: removeLineageRefs, - now: archivedAt, - beforeArchive: async (tx) => { - const linkedFeature = await getMissionFeatureByTaskId(tx, id); - if (linkedFeature) { - await recordGeneratedFixOperatorStop(tx, linkedFeature, "task-archive"); - await unlinkMissionFeatureFromTaskId(tx, linkedFeature.id); - } - }, - }); - } catch (error) { - if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace); - throw error; - } - + const archiveLineageArchivedLanes: ReadonlySet | undefined = undefined; + const archiveRun = async (context?: { candidateIds: string[]; promptByChildId: ReadonlyMap; locksHeld: boolean; attempt: number }) => { + const preparedWorkspace = cleanup ? await prepareArchivedWorkspaceWorktrees(store, task) : undefined; + try { + const result = await archiveParentTaskWithLineageGate(layer, id, entry, { + removeLineageReferences: removeLineageRefs, + now: archivedAt, + archivedColumns: archiveLineageArchivedLanes, + ...(context ? { + revalidateAgainst: context.candidateIds, + promptByChildId: context.promptByChildId, + evidenceTargetVersionForTest: lineageEvidenceTargetVersionForTest(store), + beforeLineageGate: (store as unknown as { __beforeArchiveLineageGateForTest?: () => void | Promise }).__beforeArchiveLineageGateForTest, + } : {}), + beforeArchive: async (tx) => { + const linkedFeature = await getMissionFeatureByTaskId(tx, id); + if (linkedFeature) { + await recordGeneratedFixOperatorStop(tx, linkedFeature, "task-archive"); + await unlinkMissionFeatureFromTaskId(tx, linkedFeature.id); + } + }, + }); + // Reconcile only rows the guarded UPDATE actually cleared; candidates can reparent away. + if (context) { + const lineageOutcome = result.archived + ? result.lineageOutcome ?? { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 } + : { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 }; + recordLineageInvalidationOutcome(store, { + attempt: context.attempt, locksHeld: context.locksHeld, degraded: !context.locksHeld, + candidateIds: context.candidateIds, ...lineageOutcome, + ...(result.archived ? {} : { error: "gate-rejected" as const }), + }); + if (result.archived) await reconcileClearedLineageChildren(store, lineageOutcome.clearedChildIds, { locksHeld: context.locksHeld }); + } + return { result, preparedWorkspace }; + } catch (error) { + if (context) recordLineageInvalidationOutcome(store, { + attempt: context.attempt, locksHeld: context.locksHeld, degraded: !context.locksHeld, + candidateIds: context.candidateIds, clearedChildIds: [], evidenceVersionByChild: new Map(), + evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0, + error: classifyLineageInvalidationOutcomeError(error), + }); + if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace); + throw error; + } + }; + const archiveExecution = removeLineageRefs + ? await runLineageInvalidation(store, id, { archivedColumns: archiveLineageArchivedLanes }, archiveRun) + : await archiveRun(); + const result = archiveExecution.result; + const preparedWorkspace = archiveExecution.preparedWorkspace; if (!result.archived) { if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace); throw new TaskHasLineageChildrenError(id, result.liveChildIds); diff --git a/packages/core/src/task-store/async/async-archive-lineage.ts b/packages/core/src/task-store/async/async-archive-lineage.ts index 9aaf4210f8..9eff02fbf8 100644 --- a/packages/core/src/task-store/async/async-archive-lineage.ts +++ b/packages/core/src/task-store/async/async-archive-lineage.ts @@ -32,7 +32,8 @@ import { and, desc, eq, inArray, sql } from "drizzle-orm"; import * as schema from "../../postgres/schema/index.js"; import type { AsyncDataLayer, DbTransaction } from "../../postgres/data-layer.js"; import { ACTIVE_TASK_FILTER } from "./async-persistence.js"; -import { findLiveLineageChildren, projectPartition, removeLineageReferences } from "./async-lifecycle.js"; +import { findLiveLineageChildren, projectPartition, removeLineageReferences, type LineageRemovalOutcome } from "./async-lifecycle.js"; +import { assertLineageCandidatesUnchanged } from "../lineage-approval-invalidation.js"; import { softDeleteTaskRowInTransaction, readTaskRowInTransaction, @@ -236,21 +237,25 @@ export async function archiveParentTaskWithLineageGate( layer: AsyncDataLayer, taskId: string, entry: ArchivedTaskEntry, - options: { removeLineageReferences?: boolean; now?: string; beforeArchive?: (tx: DbTransaction) => Promise } = {}, -): Promise<{ archived: true } | { archived: false; liveChildIds: string[] }> { + options: { removeLineageReferences?: boolean; now?: string; beforeArchive?: (tx: DbTransaction) => Promise; beforeLineageGate?: () => void | Promise; archivedColumns?: ReadonlySet; revalidateAgainst?: readonly string[]; promptByChildId?: ReadonlyMap; evidenceTargetVersionForTest?: (childId: string, computed: number, attempt: number) => number } = {}, + +): Promise<{ archived: true; lineageOutcome?: LineageRemovalOutcome } | { archived: false; liveChildIds: string[] }> { const now = options.now ?? new Date().toISOString(); return layer.transactionImmediate(async (tx) => { + // Test-only barrier is before this operation's single in-transaction lineage read. + await options.beforeLineageGate?.(); // 1. Lineage gate — check for live children inside the transaction. - const liveChildIds = await findLiveLineageChildren(tx, taskId, layer.projectId); + const liveChildIds = await findLiveLineageChildren(tx, taskId, layer.projectId, options.archivedColumns); if (liveChildIds.length > 0 && !options.removeLineageReferences) { return { archived: false as const, liveChildIds }; } - // 2. Lineage clear (if requested and there are live children). - if (liveChildIds.length > 0 && options.removeLineageReferences) { - await removeLineageReferences(tx, taskId, liveChildIds, now, layer.projectId); - } + if (options.removeLineageReferences && options.revalidateAgainst !== undefined) assertLineageCandidatesUnchanged(liveChildIds, options.revalidateAgainst); + // 2. Lineage clear only for candidates whose planning locks this attempt holds. + const lineageOutcome = options.removeLineageReferences + ? await removeLineageReferences(tx, taskId, options.revalidateAgainst ?? liveChildIds, now, layer.projectId, options.promptByChildId, options.evidenceTargetVersionForTest) + : { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 }; /* FNXC:MissionLineageBudget 2026-07-22-15:00: @@ -275,7 +280,11 @@ export async function archiveParentTaskWithLineageGate( // breaking atomicity (a later rollback left the soft-delete persisted). await softDeleteTaskRowInTransaction(tx, taskId, now, false, layer.projectId); - return { archived: true as const }; + // Preserve the public legacy success shape for direct callers; only the serialized boundary + // supplies revalidation and needs the actual-clear outcome for post-commit reconciliation. + return options.revalidateAgainst === undefined + ? { archived: true as const } + : { archived: true as const, lineageOutcome }; }); } diff --git a/packages/core/src/task-store/async/async-lifecycle.ts b/packages/core/src/task-store/async/async-lifecycle.ts index c0220ce6f1..1fb63cade4 100644 --- a/packages/core/src/task-store/async/async-lifecycle.ts +++ b/packages/core/src/task-store/async/async-lifecycle.ts @@ -27,10 +27,12 @@ * integration tests consume. They program against the stable `AsyncDataLayer` * interface (U4), not the underlying driver. */ -import { and, eq, ne, sql } from "drizzle-orm"; +import { and, desc, eq, ne, sql } from "drizzle-orm"; import * as schema from "../../postgres/schema/index.js"; -import type { AsyncDataLayer, DbTransaction } from "../../postgres/data-layer.js"; +import { projectScopeFor, type AsyncDataLayer, type DbTransaction } from "../../postgres/data-layer.js"; import { ACTIVE_TASK_FILTER } from "./async-persistence.js"; +import { createCurrentPlanEvidence, type CurrentPlanEvidence } from "../../planner/spec-lock.js"; +import { storeLog } from "../../store.js"; /** * FNXC:TaskStoreLifecycle 2026-06-24-04:35: @@ -152,38 +154,71 @@ export async function findLiveLineageChildren( * @param nowIso The timestamp to stamp on `updated_at`. * @returns The number of child rows actually updated (cleared). */ +export class LineageEvidenceAppendError extends Error { + constructor(public readonly childId: string, public readonly attempts: number, public readonly reason: "no-durable-version" | "matched-row-not-truthful") { + super(`Could not resolve truthful lineage evidence for ${childId}: ${reason}`); + this.name = "LineageEvidenceAppendError"; + } +} + +export type LineageRemovalOutcome = { clearedChildIds: string[]; evidenceVersionByChild: Map; evidenceUnavailableChildIds: string[]; evidenceInsertAttempts: number }; + +/** + * FNXC:SpecLockLineageInvalidation 2026-08-10-14:33: + * The lineage edge, approval fingerprint, and successor evidence are one transaction: a cleared + * parent binding can never retain approval or commit without truthful evidence. Report publication + * is deliberately outside this helper because it must evaluate committed state and is retryable. + */ export async function removeLineageReferences( tx: DbTransaction, parentId: string, childIds: readonly string[], nowIso: string, projectId?: string, -): Promise { - // FNXC:TaskStoreLifecycle 2026-06-24-06:05: - // A single bulk UPDATE clears all children that still point at this parent. - // The WHERE guards on BOTH id (in the child set) AND source_parent_task_id - // (still pointing at this parent), so a child that was reparented elsewhere - // is left untouched. Using an IN-list keeps this to one round-trip regardless - // of child count. We count affected rows via a RETURNING read so the count is - // accurate regardless of how the driver exposes rowCount. - if (childIds.length === 0) { - return 0; + promptByChildId: ReadonlyMap = new Map(), + evidenceTargetVersionForTest?: (childId: string, computed: number, attempt: number) => number, +): Promise { + if (childIds.length === 0) return { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 }; + const returned = await tx.update(schema.project.tasks).set({ sourceParentTaskId: null, approvedPlanFingerprint: null, awaitingApprovalReason: null, updatedAt: nowIso }) + .where(and(eq(schema.project.tasks.projectId, projectPartition(projectId)), sql`${schema.project.tasks.id} IN ${childIds}`, eq(schema.project.tasks.sourceParentTaskId, parentId))) + .returning({ id: schema.project.tasks.id, dependencies: schema.project.tasks.dependencies, missionId: schema.project.tasks.missionId, sliceId: schema.project.tasks.sliceId }); + const evidenceVersionByChild = new Map(); + const evidenceUnavailableChildIds: string[] = []; + let evidenceInsertAttempts = 0; + /* + FNXC:SpecLockLineageInvalidation 2026-08-10-14:47: + Evidence inserts use the trigger-compatible blank value, while reads preserve an unbound layer's + project-agnostic scope instead of filtering for the blank value the trigger rewrites. + */ + const evidenceProjectId = projectId ?? ""; + for (const child of returned) { + const prompt = promptByChildId.get(child.id); + if (prompt === undefined) { + evidenceUnavailableChildIds.push(child.id); + // A missing plan is the sole evidence-less commit: make the repairable input absence visible. + storeLog.warn(`[spec-lock] lineage evidence unavailable: missing PROMPT.md for ${child.id}`); + continue; + } + let resolved = false; + for (let attempt = 0; attempt < 3 && !resolved; attempt += 1) { + const latest = await tx.select({ version: schema.project.currentPlanEvidence.version }).from(schema.project.currentPlanEvidence) + .where(and(projectScopeFor(schema.project.currentPlanEvidence.projectId, projectId), eq(schema.project.currentPlanEvidence.taskId, child.id))).orderBy(desc(schema.project.currentPlanEvidence.version)).limit(1); + const computed = (latest[0]?.version ?? 0) + 1; + const evidence = createCurrentPlanEvidence({ version: evidenceTargetVersionForTest?.(child.id, computed, attempt) ?? computed, sourceRevision: Date.now(), capturedAt: nowIso, prompt, bindings: { dependencies: (child.dependencies as string[] | undefined) ?? [], missionId: child.missionId ?? undefined, sliceId: child.sliceId ?? undefined } }); + evidenceInsertAttempts += 1; + const inserted = await tx.insert(schema.project.currentPlanEvidence).values({ projectId: evidenceProjectId, taskId: child.id, version: evidence.version, sourceRevision: evidence.sourceRevision, sourceHash: evidence.sourceHash, capturedAt: evidence.capturedAt, snapshot: evidence }).onConflictDoNothing().returning({ version: schema.project.currentPlanEvidence.version }); + if (inserted[0]) { evidenceVersionByChild.set(child.id, inserted[0].version); resolved = true; break; } + const matching = await tx.select({ version: schema.project.currentPlanEvidence.version, snapshot: schema.project.currentPlanEvidence.snapshot }).from(schema.project.currentPlanEvidence) + .where(and(projectScopeFor(schema.project.currentPlanEvidence.projectId, projectId), eq(schema.project.currentPlanEvidence.taskId, child.id), eq(schema.project.currentPlanEvidence.sourceHash, evidence.sourceHash))).limit(1); + if (matching[0]) { + const snapshot = matching[0].snapshot as CurrentPlanEvidence; + if (snapshot.plan.sections.lineage.canonical.split("\n").includes(`parent-task:${parentId}`)) throw new LineageEvidenceAppendError(child.id, attempt + 1, "matched-row-not-truthful"); + evidenceVersionByChild.set(child.id, matching[0].version); resolved = true; break; + } + } + if (!resolved) throw new LineageEvidenceAppendError(child.id, 3, "no-durable-version"); } - const returned = await tx - .update(schema.project.tasks) - .set({ - sourceParentTaskId: null, - updatedAt: nowIso, - }) - .where( - and( - eq(schema.project.tasks.projectId, projectPartition(projectId)), - sql`${schema.project.tasks.id} IN ${childIds}`, - eq(schema.project.tasks.sourceParentTaskId, parentId), - ), - ) - .returning({ id: schema.project.tasks.id }); - return returned.length; + return { clearedChildIds: returned.map((row) => row.id), evidenceVersionByChild, evidenceUnavailableChildIds, evidenceInsertAttempts }; } /** diff --git a/packages/core/src/task-store/lineage-approval-invalidation.ts b/packages/core/src/task-store/lineage-approval-invalidation.ts new file mode 100644 index 0000000000..97619bf4c5 --- /dev/null +++ b/packages/core/src/task-store/lineage-approval-invalidation.ts @@ -0,0 +1,168 @@ +import { readFile } from "node:fs/promises"; +import { join } from "node:path"; +import { storeLog, type TaskStore } from "../store.js"; +import type { AsyncDataLayer, DbTransaction } from "../postgres/data-layer.js"; +import { findLiveLineageChildren, LineageEvidenceAppendError } from "./async/async-lifecycle.js"; + +/** A raced-in child must be locked on a fresh attempt rather than cleared unlocked. */ +class CandidateChanged extends Error {} +export class LineageInvalidationCandidateRaceError extends Error { + constructor(readonly parentId: string) { super(`Lineage children changed while invalidating ${parentId}`); this.name = "LineageInvalidationCandidateRaceError"; } +} +export { LineageEvidenceAppendError }; + +export type LineageInvalidationOutcome = { + attempt: number; + locksHeld: boolean; + degraded: boolean; + candidateIds: readonly string[]; + clearedChildIds: readonly string[]; + evidenceVersionByChild: ReadonlyMap; + evidenceUnavailableChildIds: readonly string[]; + evidenceInsertAttempts: number; + error?: "candidate-race-sentinel" | "evidence-append" | "gate-rejected" | "other"; +}; + +export type LineageInvalidationTestEvent = + | { kind: "probe"; attempt: number } + | { kind: "acquire"; attempt: number; childIds: readonly string[] } + | { kind: "release"; attempt: number; childIds: readonly string[] } + | { kind: "run"; attempt: number; candidateIds: readonly string[]; locksHeld: boolean; degraded: boolean } + | { kind: "outcome"; outcome: LineageInvalidationOutcome } + | { kind: "run-error"; attempt: number; error: unknown } + | { kind: "reconcile"; childId: string }; + +type LineageInvalidationTestSeam = { + availability?: () => { available: true } | { available: false; reason: string }; + withLocks?: (ids: readonly string[], callback: () => Promise) => Promise; + reconcile?: (childId: string) => Promise; + /** Test-only deterministic target collision injection for transactional evidence coverage. */ + evidenceTargetVersionForTest?: (childId: string, computed: number, attempt: number) => number; + onEvent?: (event: LineageInvalidationTestEvent) => void; +}; + +function testSeam(store: TaskStore): LineageInvalidationTestSeam | undefined { + return (store as unknown as { __lineageInvalidationForTest?: LineageInvalidationTestSeam }).__lineageInvalidationForTest; +} + +/* +FNXC:SpecLockLineageInvalidation 2026-08-10-15:35: +Evidence version collisions cannot be produced by pre-seeding the next version because the +transaction derives its target from MAX(version). Expose only this inert deterministic seam so the +retry and rollback terminals are regression-tested without weakening production conflict handling. +*/ +export function lineageEvidenceTargetVersionForTest(store: TaskStore): + | ((childId: string, computed: number, attempt: number) => number) + | undefined { + return testSeam(store)?.evidenceTargetVersionForTest; +} + +/* +FNXC:SpecLockLineageInvalidation 2026-08-10-15:15: +Delete/archive return only their parent result, so the per-attempt child clear, durable-evidence, +and error facts would otherwise be untestable. Publish this inert test seam after the transaction +and before reconciliation, preserving the required acquire → outcome → reconcile → release order. +*/ +export function recordLineageInvalidationOutcome(store: TaskStore, outcome: LineageInvalidationOutcome): void { + testSeam(store)?.onEvent?.({ kind: "outcome", outcome }); +} + +export function classifyLineageInvalidationOutcomeError(error: unknown): LineageInvalidationOutcome["error"] { + if (error instanceof CandidateChanged) return "candidate-race-sentinel"; + if (error instanceof LineageEvidenceAppendError) return "evidence-append"; + return "other"; +} + +export async function resolveLineageInvalidationCandidates(layer: AsyncDataLayer, parentId: string, archivedColumns: ReadonlySet | undefined): Promise { + return (await findLiveLineageChildren(layer.db, parentId, layer.projectId, archivedColumns)).sort(); +} +export async function readLineagePromptMap(store: TaskStore, childIds: readonly string[]): Promise> { + const prompts = new Map(); + await Promise.all(childIds.map(async (childId) => { + try { prompts.set(childId, await readFile(join(store.taskDir(childId), "PROMPT.md"), "utf8")); } catch { /* Missing plans are the explicitly supported evidence-unavailable case. */ } + })); + return prompts; +} +export function assertLineageCandidatesUnchanged(inTransactionChildIds: readonly string[], candidateIds: readonly string[]): void { + const candidates = new Set(candidateIds); + if (inTransactionChildIds.some((id) => !candidates.has(id))) throw new CandidateChanged(); +} +export async function resolveAndAssertLineageCandidatesUnchanged(tx: DbTransaction, parentId: string, projectId: string | undefined, archivedColumns: ReadonlySet | undefined, candidateIds: readonly string[]): Promise { + const live = await findLiveLineageChildren(tx, parentId, projectId, archivedColumns); + assertLineageCandidatesUnchanged(live, candidateIds); + return live; +} + +/** + * FNXC:SpecLockLineageInvalidation 2026-08-10-14:33: + * Enter on the remove-lineage flag even with no candidates: archive discovers children in its + * transaction, so an empty pre-read still needs revalidation. A structural pooler-only transport + * may omit serialization, but a runtime lock failure propagates before the mutation body. + */ +export async function runLineageInvalidation( + store: TaskStore, + parentId: string, + input: { archivedColumns: ReadonlySet | undefined; initialCandidateIds?: readonly string[] }, + run: (context: { candidateIds: string[]; promptByChildId: ReadonlyMap; locksHeld: boolean; attempt: number }) => Promise, +): Promise { + const layer = store.asyncLayer!; + const seam = testSeam(store); + for (let attempt = 0; attempt < 2; attempt += 1) { + const candidateIds = (attempt === 0 && input.initialCandidateIds !== undefined ? [...input.initialCandidateIds] : await resolveLineageInvalidationCandidates(layer, parentId, input.archivedColumns)).sort(); + const invoke = async (locksHeld: boolean, degraded: boolean, promptByChildId: ReadonlyMap): Promise => { + seam?.onEvent?.({ kind: "run", attempt, candidateIds, locksHeld, degraded }); + try { + return await run({ candidateIds, promptByChildId, locksHeld, attempt }); + } catch (error) { + seam?.onEvent?.({ kind: "run-error", attempt, error }); + throw error; + } + }; + try { + if (candidateIds.length === 0) return await invoke(false, false, new Map()); + const promptByChildId = await readLineagePromptMap(store, candidateIds); + seam?.onEvent?.({ kind: "probe", attempt }); + const availability = seam?.availability?.() ?? store.planningLifecycleLockTransportAvailability(); + if (!availability.available) { + /* + FNXC:SpecLockLineageInvalidation 2026-08-10-15:22: + A pooler-only deployment deliberately degrades serialization, not durable invalidation. + Warn once per operation so operators can distinguish that sanctioned path from a contention + failure, which must still reject before the mutation body. + */ + storeLog.warn(`[spec-lock] lineage invalidation for ${parentId} runs without planning lifecycle locks for ${candidateIds.length} child(ren): ${availability.reason}`); + return await invoke(false, true, promptByChildId); + } + seam?.onEvent?.({ kind: "acquire", attempt, childIds: candidateIds }); + const withLocks = seam?.withLocks ?? ((ids, callback) => store.withPlanningLifecycleLocks(ids, callback)); + let callbackEntered = false; + try { + return await withLocks(candidateIds, () => { + callbackEntered = true; + return invoke(true, false, promptByChildId); + }); + } finally { + // Acquisition failures run no body and therefore own no observable callback lock window. + if (callbackEntered) seam?.onEvent?.({ kind: "release", attempt, childIds: candidateIds }); + } + } catch (error) { + if (!(error instanceof CandidateChanged)) throw error; + if (attempt === 1) throw new LineageInvalidationCandidateRaceError(parentId); + } + } + throw new LineageInvalidationCandidateRaceError(parentId); +} + +/** Publish after commit but before the multi-key callback releases child planning locks. */ +export async function reconcileClearedLineageChildren(store: TaskStore, childIds: readonly string[], _context: { locksHeld: boolean }): Promise { + const seam = testSeam(store); + for (const childId of [...new Set(childIds)].sort()) { + seam?.onEvent?.({ kind: "reconcile", childId }); + await (seam?.reconcile?.(childId) + ?? store.reconcileSpecDriftWhilePlanningLocked({ id: childId, approvedPlanFingerprint: undefined, modifiedFiles: [] })) + .catch((error: unknown) => { + // Report publication is retryable telemetry, but an operator still needs the failed child id. + storeLog.warn(`[spec-lock] deferred lineage drift reconciliation for ${childId}: ${error instanceof Error ? error.message : String(error)}`); + }); + } +}