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:
7
.changeset/fn-8945-lineage-approval-invalidation.md
Normal file
7
.changeset/fn-8945-lineage-approval-invalidation.md
Normal 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.
|
||||||
@@ -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.
|
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.
|
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
|
## Durable intake executor ownership
|
||||||
|
|||||||
@@ -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");
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -7,7 +7,9 @@ import {
|
|||||||
} from "../../__test-utils__/pg-test-harness.js";
|
} from "../../__test-utils__/pg-test-harness.js";
|
||||||
import {
|
import {
|
||||||
PlanningLifecycleLockTransportError,
|
PlanningLifecycleLockTransportError,
|
||||||
|
planningLifecycleLockTransportAvailability,
|
||||||
withPlanningLifecycleAdvisoryLock,
|
withPlanningLifecycleAdvisoryLock,
|
||||||
|
withPlanningLifecycleAdvisoryLocks,
|
||||||
} from "../../postgres/advisory-locks.js";
|
} from "../../postgres/advisory-locks.js";
|
||||||
import { resolveBackendWithOptions } from "../../postgres/backend-resolver.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"]);
|
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 () => {
|
it("does not contend across project or task keys", async () => {
|
||||||
let releaseFirst!: () => void;
|
let releaseFirst!: () => void;
|
||||||
const firstCanFinish = new Promise<void>((resolve) => { releaseFirst = resolve; });
|
const firstCanFinish = new Promise<void>((resolve) => { releaseFirst = resolve; });
|
||||||
|
|||||||
@@ -269,6 +269,7 @@ pgDescribe("runtime-lifecycle-async: deleteTask lineage gate (PostgreSQL)", () =
|
|||||||
createdAt: new Date().toISOString(),
|
createdAt: new Date().toISOString(),
|
||||||
updatedAt: new Date().toISOString(),
|
updatedAt: new Date().toISOString(),
|
||||||
sourceParentTaskId: "FN-PARENT2",
|
sourceParentTaskId: "FN-PARENT2",
|
||||||
|
approvedPlanFingerprint: "approved-lineage-scope",
|
||||||
status: null,
|
status: null,
|
||||||
} as never,
|
} as never,
|
||||||
{ lineageId: "test" },
|
{ lineageId: "test" },
|
||||||
@@ -284,5 +285,8 @@ pgDescribe("runtime-lifecycle-async: deleteTask lineage gate (PostgreSQL)", () =
|
|||||||
.where(eq(schema.project.tasks.id, "FN-PARENT2"));
|
.where(eq(schema.project.tasks.id, "FN-PARENT2"));
|
||||||
expect(rows.length).toBe(1);
|
expect(rows.length).toBe(1);
|
||||||
expect(rows[0].deletedAt).not.toBeNull();
|
expect(rows[0].deletedAt).not.toBeNull();
|
||||||
|
const child = await h.store.getTask("FN-CHILD2");
|
||||||
|
expect(child.sourceParentTaskId).toBeUndefined();
|
||||||
|
expect(child.approvedPlanFingerprint).toBeUndefined();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -219,7 +219,8 @@ pgDescribe("U13 taskstore-lifecycle (PostgreSQL)", () => {
|
|||||||
const childIds = await findLiveLineageChildren(tx, "KB-PARENT");
|
const childIds = await findLiveLineageChildren(tx, "KB-PARENT");
|
||||||
expect(childIds).toEqual(["KB-CHILD"]);
|
expect(childIds).toEqual(["KB-CHILD"]);
|
||||||
const cleared = await removeLineageReferences(tx, "KB-PARENT", childIds, nowIso);
|
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).
|
// After: gate passes (no live children).
|
||||||
|
|||||||
@@ -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 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 = {
|
type PostgresBackendIdentity = {
|
||||||
database: string;
|
database: string;
|
||||||
host: string | null;
|
host: string | null;
|
||||||
@@ -85,53 +116,17 @@ async function runBoundedTransportPhase<T>(
|
|||||||
* lock and unlock and silently defeat session advisory locking.
|
* lock and unlock and silently defeat session advisory locking.
|
||||||
*/
|
*/
|
||||||
export async function withPlanningLifecycleAdvisoryLock<T>(
|
export async function withPlanningLifecycleAdvisoryLock<T>(
|
||||||
input: {
|
input: PlanningLifecycleLockInput,
|
||||||
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;
|
|
||||||
},
|
|
||||||
callback: () => Promise<T>,
|
callback: () => Promise<T>,
|
||||||
): Promise<T> {
|
): Promise<T> {
|
||||||
const directUrl = input.directSessionUrl;
|
const availability = planningLifecycleLockTransportAvailability(input);
|
||||||
const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS);
|
if (!availability.available) {
|
||||||
if (!directUrl || !input.provenance || looksLikePoolerUrl(directUrl)) {
|
throw new PlanningLifecycleLockTransportError(`Planning lifecycle lock transport unavailable: ${availability.reason}`);
|
||||||
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 directUrl = input.directSessionUrl!;
|
||||||
|
const directDatabase = decodeURIComponent(new URL(directUrl).pathname.replace(/^\//, ""));
|
||||||
const operationalUrl = input.runtimeUrl ?? input.migrationUrl;
|
const operationalUrl = input.runtimeUrl ?? input.migrationUrl;
|
||||||
if (operationalUrl) {
|
const timeoutMs = Math.max(1, input.timeoutMs ?? DEFAULT_PLANNING_LIFECYCLE_LOCK_TIMEOUT_MS);
|
||||||
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 client = postgres(directUrl, {
|
const client = postgres(directUrl, {
|
||||||
max: 1,
|
max: 1,
|
||||||
@@ -144,7 +139,7 @@ export async function withPlanningLifecycleAdvisoryLock<T>(
|
|||||||
}
|
}
|
||||||
: undefined,
|
: undefined,
|
||||||
});
|
});
|
||||||
const key = `fusion:planning-lifecycle:${input.projectId}:${input.taskId}`;
|
const key = planningLifecycleLockKey(input.projectId, input.taskId);
|
||||||
let acquired = false;
|
let acquired = false;
|
||||||
try {
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -78,7 +78,7 @@ import { GlobalSettingsStore } from "./config/global-settings.js";
|
|||||||
import { Database } from "./db/db.js";
|
import { Database } from "./db/db.js";
|
||||||
import { ArchiveDatabase } from "./db/archive-db.js";
|
import { ArchiveDatabase } from "./db/archive-db.js";
|
||||||
import { projectScopeFor, type AsyncDataLayer, type DbTransaction } from "./postgres/data-layer.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 { MissionStore } from "./missions/mission-store.js";
|
||||||
import { AsyncMissionStore } from "./async-stores/async-mission-store.js";
|
import { AsyncMissionStore } from "./async-stores/async-mission-store.js";
|
||||||
import { AsyncIdeationStore } from "./async-stores/async-ideation-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);
|
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 {
|
public getTaskIdFromDir(dir: string): string {
|
||||||
return getTaskIdFromDirImpl(this, dir);
|
return getTaskIdFromDirImpl(this, dir);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -26,7 +26,8 @@ import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
|
|||||||
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
|
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
|
||||||
import {softDeleteTaskRowInTransaction, readTaskRow as readTaskRowAsync, readTaskRowInTransaction} from "../task-store/async/async-persistence.js";
|
import {softDeleteTaskRowInTransaction, readTaskRow as readTaskRowAsync, readTaskRowInTransaction} from "../task-store/async/async-persistence.js";
|
||||||
import {appendTaskLifecycleEventInTransaction} from "../task-store/lifecycle-outbox.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 { resolveProjectColumnsForRoles } from "../project-lane-vocabulary.js";
|
||||||
import {archiveParentTaskWithLineageGate, findArchivedTaskEntry, deleteArchivedTaskEntry, restoreTaskFromArchive} from "../task-store/async/async-archive-lineage.js";
|
import {archiveParentTaskWithLineageGate, findArchivedTaskEntry, deleteArchivedTaskEntry, restoreTaskFromArchive} from "../task-store/async/async-archive-lineage.js";
|
||||||
import {getArchivedRowCount, listArchivedTaskEntriesPage} from "../async-stores/async-archive-db.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);
|
await (store as unknown as { __beforeDeleteClaimForTest?: (taskId: string) => void | Promise<void> }).__beforeDeleteClaimForTest?.(id);
|
||||||
|
|
||||||
// Soft-delete + lineage clear + mission unlink + audit in one transaction (atomicity).
|
// 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:
|
FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
The pre-transaction deletedAt read is a cross-process TOCTOU window. A conditional
|
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) {
|
if (claimed === false) {
|
||||||
const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId);
|
const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId);
|
||||||
if (!reloaded) throw new TaskNotFoundError(id);
|
if (!reloaded) throw new TaskNotFoundError(id);
|
||||||
return { claimed: false, task: store.rowToTask(store.pgRowToTaskRow(reloaded)) };
|
return { claimed: false, task: store.rowToTask(store.pgRowToTaskRow(reloaded)), lineageOutcome: { clearedChildIds: [] as string[], evidenceVersionByChild: new Map<string, number>(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 } };
|
||||||
}
|
|
||||||
// Clear lineage references on live children so the parent can be deleted.
|
|
||||||
if (lineageChildIds.length > 0) {
|
|
||||||
await removeLineageReferences(tx, id, lineageChildIds, deletedAt, layer.projectId);
|
|
||||||
}
|
}
|
||||||
|
// 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:
|
FNXC:MissionStore 2026-07-17-17:40:
|
||||||
Clear any mission feature→task link IN THIS TRANSACTION so it commits (or rolls
|
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.
|
// callers so neither receives the pre-claim live snapshot after a successful delete.
|
||||||
const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId);
|
const reloaded = await readTaskRowInTransaction(tx, id, { includeDeleted: true }, projectId);
|
||||||
if (!reloaded) throw new TaskNotFoundError(id);
|
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;
|
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);
|
const entry = await store.taskToArchiveEntry(task, archivedAt);
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowLifecycle 2026-07-16-15:30:
|
FNXC:SpecLockLineageInvalidation 2026-08-10-14:33:
|
||||||
Backend archive persists cold storage before its cleanup phase. Hold the
|
Archive keeps its legacy archived-lane semantics (undefined -> "archived") but threads that
|
||||||
per-repository reservations across that transaction so another process sees
|
one value through pre-read and gate. The workspace reservation is created inside the locked body.
|
||||||
the path as unavailable until the awaited workspace disposer has removed it.
|
|
||||||
*/
|
*/
|
||||||
const preparedWorkspace = cleanup ? await prepareArchivedWorkspaceWorktrees(store, task) : undefined;
|
const archiveLineageArchivedLanes: ReadonlySet<string> | undefined = undefined;
|
||||||
let result;
|
const archiveRun = async (context?: { candidateIds: string[]; promptByChildId: ReadonlyMap<string, string>; locksHeld: boolean; attempt: number }) => {
|
||||||
try {
|
const preparedWorkspace = cleanup ? await prepareArchivedWorkspaceWorktrees(store, task) : undefined;
|
||||||
// Lineage gate + archive in one transaction.
|
try {
|
||||||
result = await archiveParentTaskWithLineageGate(layer, id, entry, {
|
const result = await archiveParentTaskWithLineageGate(layer, id, entry, {
|
||||||
removeLineageReferences: removeLineageRefs,
|
removeLineageReferences: removeLineageRefs,
|
||||||
now: archivedAt,
|
now: archivedAt,
|
||||||
beforeArchive: async (tx) => {
|
archivedColumns: archiveLineageArchivedLanes,
|
||||||
const linkedFeature = await getMissionFeatureByTaskId(tx, id);
|
...(context ? {
|
||||||
if (linkedFeature) {
|
revalidateAgainst: context.candidateIds,
|
||||||
await recordGeneratedFixOperatorStop(tx, linkedFeature, "task-archive");
|
promptByChildId: context.promptByChildId,
|
||||||
await unlinkMissionFeatureFromTaskId(tx, linkedFeature.id);
|
evidenceTargetVersionForTest: lineageEvidenceTargetVersionForTest(store),
|
||||||
}
|
beforeLineageGate: (store as unknown as { __beforeArchiveLineageGateForTest?: () => void | Promise<void> }).__beforeArchiveLineageGateForTest,
|
||||||
},
|
} : {}),
|
||||||
});
|
beforeArchive: async (tx) => {
|
||||||
} catch (error) {
|
const linkedFeature = await getMissionFeatureByTaskId(tx, id);
|
||||||
if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace);
|
if (linkedFeature) {
|
||||||
throw error;
|
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 (!result.archived) {
|
||||||
if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace);
|
if (preparedWorkspace) await releasePreparedWorkspaceArchiveDisposal(preparedWorkspace);
|
||||||
throw new TaskHasLineageChildrenError(id, result.liveChildIds);
|
throw new TaskHasLineageChildrenError(id, result.liveChildIds);
|
||||||
|
|||||||
@@ -32,7 +32,8 @@ import { and, desc, eq, inArray, sql } from "drizzle-orm";
|
|||||||
import * as schema from "../../postgres/schema/index.js";
|
import * as schema from "../../postgres/schema/index.js";
|
||||||
import type { AsyncDataLayer, DbTransaction } from "../../postgres/data-layer.js";
|
import type { AsyncDataLayer, DbTransaction } from "../../postgres/data-layer.js";
|
||||||
import { ACTIVE_TASK_FILTER } from "./async-persistence.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 {
|
import {
|
||||||
softDeleteTaskRowInTransaction,
|
softDeleteTaskRowInTransaction,
|
||||||
readTaskRowInTransaction,
|
readTaskRowInTransaction,
|
||||||
@@ -236,21 +237,25 @@ export async function archiveParentTaskWithLineageGate(
|
|||||||
layer: AsyncDataLayer,
|
layer: AsyncDataLayer,
|
||||||
taskId: string,
|
taskId: string,
|
||||||
entry: ArchivedTaskEntry,
|
entry: ArchivedTaskEntry,
|
||||||
options: { removeLineageReferences?: boolean; now?: string; beforeArchive?: (tx: DbTransaction) => Promise<void> } = {},
|
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 } | { archived: false; liveChildIds: string[] }> {
|
|
||||||
|
): Promise<{ archived: true; lineageOutcome?: LineageRemovalOutcome } | { archived: false; liveChildIds: string[] }> {
|
||||||
const now = options.now ?? new Date().toISOString();
|
const now = options.now ?? new Date().toISOString();
|
||||||
|
|
||||||
return layer.transactionImmediate(async (tx) => {
|
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.
|
// 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) {
|
if (liveChildIds.length > 0 && !options.removeLineageReferences) {
|
||||||
return { archived: false as const, liveChildIds };
|
return { archived: false as const, liveChildIds };
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. Lineage clear (if requested and there are live children).
|
if (options.removeLineageReferences && options.revalidateAgainst !== undefined) assertLineageCandidatesUnchanged(liveChildIds, options.revalidateAgainst);
|
||||||
if (liveChildIds.length > 0 && options.removeLineageReferences) {
|
// 2. Lineage clear only for candidates whose planning locks this attempt holds.
|
||||||
await removeLineageReferences(tx, taskId, liveChildIds, now, layer.projectId);
|
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:
|
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).
|
// breaking atomicity (a later rollback left the soft-delete persisted).
|
||||||
await softDeleteTaskRowInTransaction(tx, taskId, now, false, layer.projectId);
|
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 };
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -27,10 +27,12 @@
|
|||||||
* integration tests consume. They program against the stable `AsyncDataLayer`
|
* integration tests consume. They program against the stable `AsyncDataLayer`
|
||||||
* interface (U4), not the underlying driver.
|
* 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 * 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 { 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:
|
* FNXC:TaskStoreLifecycle 2026-06-24-04:35:
|
||||||
@@ -152,38 +154,71 @@ export async function findLiveLineageChildren(
|
|||||||
* @param nowIso The timestamp to stamp on `updated_at`.
|
* @param nowIso The timestamp to stamp on `updated_at`.
|
||||||
* @returns The number of child rows actually updated (cleared).
|
* @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(
|
export async function removeLineageReferences(
|
||||||
tx: DbTransaction,
|
tx: DbTransaction,
|
||||||
parentId: string,
|
parentId: string,
|
||||||
childIds: readonly string[],
|
childIds: readonly string[],
|
||||||
nowIso: string,
|
nowIso: string,
|
||||||
projectId?: string,
|
projectId?: string,
|
||||||
): Promise<number> {
|
promptByChildId: ReadonlyMap<string, string> = new Map(),
|
||||||
// FNXC:TaskStoreLifecycle 2026-06-24-06:05:
|
evidenceTargetVersionForTest?: (childId: string, computed: number, attempt: number) => number,
|
||||||
// A single bulk UPDATE clears all children that still point at this parent.
|
): Promise<LineageRemovalOutcome> {
|
||||||
// The WHERE guards on BOTH id (in the child set) AND source_parent_task_id
|
if (childIds.length === 0) return { clearedChildIds: [], evidenceVersionByChild: new Map(), evidenceUnavailableChildIds: [], evidenceInsertAttempts: 0 };
|
||||||
// (still pointing at this parent), so a child that was reparented elsewhere
|
const returned = await tx.update(schema.project.tasks).set({ sourceParentTaskId: null, approvedPlanFingerprint: null, awaitingApprovalReason: null, updatedAt: nowIso })
|
||||||
// is left untouched. Using an IN-list keeps this to one round-trip regardless
|
.where(and(eq(schema.project.tasks.projectId, projectPartition(projectId)), sql`${schema.project.tasks.id} IN ${childIds}`, eq(schema.project.tasks.sourceParentTaskId, parentId)))
|
||||||
// of child count. We count affected rows via a RETURNING read so the count is
|
.returning({ id: schema.project.tasks.id, dependencies: schema.project.tasks.dependencies, missionId: schema.project.tasks.missionId, sliceId: schema.project.tasks.sliceId });
|
||||||
// accurate regardless of how the driver exposes rowCount.
|
const evidenceVersionByChild = new Map<string, number>();
|
||||||
if (childIds.length === 0) {
|
const evidenceUnavailableChildIds: string[] = [];
|
||||||
return 0;
|
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
|
return { clearedChildIds: returned.map((row) => row.id), evidenceVersionByChild, evidenceUnavailableChildIds, evidenceInsertAttempts };
|
||||||
.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;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
168
packages/core/src/task-store/lineage-approval-invalidation.ts
Normal file
168
packages/core/src/task-store/lineage-approval-invalidation.ts
Normal 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)}`);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user