FN-8945: invalidate approved lineage fingerprints

Invalidate stale approval evidence atomically when a parent task leaves lineage.

- Lock and revalidate lineage children before parent deletion or archival.
- Clear affected approvals, append invalidation evidence, and reconcile cleared children.
- Cover lifecycle races, evidence collisions, and async archive/delete paths.

Files changed:
 .changeset/fn-8945-lineage-approval-invalidation.md       |   7 +
 docs/architecture.md                               |   6 +
 .../lineage-approval-invalidation.pg.test.ts       | 379 +++++++++++++++++++++
 .../planning-lifecycle-advisory-lock.pg.test.ts    |  58 ++++
 .../postgres/runtime-lifecycle-async.test.ts       |   4 +
 .../__tests__/postgres/taskstore-lifecycle.test.ts |   3 +-
 packages/core/src/postgres/advisory-locks.ts       | 164 ++++++---
 packages/core/src/store.ts                         |  26 +-
 .../core/src/task-store/archive-lifecycle-2.ts     | 123 +++++--
 .../src/task-store/async/async-archive-lineage.ts  |  27 +-
 .../core/src/task-store/async/async-lifecycle.ts   |  89 +--
 .../task-store/lineage-approval-invalidation.ts    | 168 +++++++++
 12 files changed, 939 insertions(+), 115 deletions(-)

Fusion-Task-Id: FN-8945

Fusion-Task-Lineage: be7dd6c8-5411-4b0e-8734-9fe63b9d1d9f

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-10 08:48:47 -07:00
parent 0fc6f3d849
commit 82b5d78bf8
12 changed files with 939 additions and 115 deletions

View File

@@ -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.

View File

@@ -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

View File

@@ -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<LineageEvidenceAppendError>);
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<LineageEvidenceAppendError>);
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: <T>(ids: readonly string[], callback: () => Promise<T>) => Promise<T>;
onEvent: (event: LineageInvalidationTestEvent) => void;
};
__beforeDeleteClaimForTest?: () => Promise<void>;
}).__lineageInvalidationForTest = {
availability: () => ({ available: true }),
async withLocks(_ids, callback) {
return await callback();
},
onEvent: (event) => events.push(event),
};
(store as unknown as { __beforeDeleteClaimForTest?: () => Promise<void> }).__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<void>;
}).__lineageInvalidationForTest = { onEvent: (event) => events.push(event) };
(store as unknown as { __beforeArchiveLineageGateForTest?: () => Promise<void> }).__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");
});
});

View File

@@ -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<void>((resolve) => { releaseMulti = resolve; });
let multiEntered!: () => void;
const multiIsHolding = new Promise<void>((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<void>((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<void>((resolve) => { releaseFirst = resolve; });

View File

@@ -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();
});
});

View File

@@ -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).

View File

@@ -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<typeof postgres>;
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<PlanningLifecycleLockInput, "directSessionUrl" | "provenance" | "runtimeUrl" | "migrationUrl">): { 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<T>(
* lock and unlock and silently defeat session advisory locking.
*/
export async function withPlanningLifecycleAdvisoryLock<T>(
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<T>,
): Promise<T> {
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<T>(
}
: 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<T>(
}
}
}
/**
* 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<T>(
input: Omit<PlanningLifecycleLockInput, "taskId"> & { taskIds: readonly string[] },
callback: () => Promise<T>,
): Promise<T> {
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);
}
}

View File

@@ -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<TaskStoreEvents> {
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<T>(ids: readonly string[], fn: () => Promise<T>): Promise<T> {
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<T> => 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);
}

View File

@@ -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<void> }).__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<string, string>; 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<string, number>(), 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<string, number>(), 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<string> | undefined = undefined;
const archiveRun = async (context?: { candidateIds: string[]; promptByChildId: ReadonlyMap<string, string>; 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<void> }).__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<string, number>(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 }
: { clearedChildIds: [], evidenceVersionByChild: new Map<string, number>(), 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);

View File

@@ -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<void> } = {},
): Promise<{ archived: true } | { archived: false; liveChildIds: string[] }> {
options: { removeLineageReferences?: boolean; now?: string; beforeArchive?: (tx: DbTransaction) => Promise<void>; beforeLineageGate?: () => void | Promise<void>; archivedColumns?: ReadonlySet<string>; revalidateAgainst?: readonly string[]; promptByChildId?: ReadonlyMap<string, string>; 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<string, number>(), 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 };
});
}

View File

@@ -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<string, number>; 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<number> {
// 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<string, string> = new Map(),
evidenceTargetVersionForTest?: (childId: string, computed: number, attempt: number) => number,
): Promise<LineageRemovalOutcome> {
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<string, number>();
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 };
}
/**

View File

@@ -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<string, number>;
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?: <T>(ids: readonly string[], callback: () => Promise<T>) => Promise<T>;
reconcile?: (childId: string) => Promise<void>;
/** 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<string> | undefined): Promise<string[]> {
return (await findLiveLineageChildren(layer.db, parentId, layer.projectId, archivedColumns)).sort();
}
export async function readLineagePromptMap(store: TaskStore, childIds: readonly string[]): Promise<Map<string, string>> {
const prompts = new Map<string, string>();
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<string> | undefined, candidateIds: readonly string[]): Promise<string[]> {
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<T>(
store: TaskStore,
parentId: string,
input: { archivedColumns: ReadonlySet<string> | undefined; initialCandidateIds?: readonly string[] },
run: (context: { candidateIds: string[]; promptByChildId: ReadonlyMap<string, string>; locksHeld: boolean; attempt: number }) => Promise<T>,
): Promise<T> {
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<string, string>): Promise<T> => {
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<void> {
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)}`);
});
}
}