FN-8684: persist task deletion events transactionally
Persist task-deletion lifecycle events through a transactional PostgreSQL outbox. - Add the task lifecycle outbox schema, migration, and atomic writer. - Route deletion notice persistence through the transaction and preserve non-blocking cleanup. - Cover outbox behavior, schema installation, and caller attribution with tests. Files changed: .changeset/fn-8684-task-deleted-outbox-writer.md | 7 + docs/architecture.md | 6 + docs/storage.md | 6 + .../src/__tests__/postgres/schema-applier.test.ts | 82 ++++++++- .../task-delete-caller-attribution.test.ts | 12 +- .../task-delete-nonblocking-cleanup.test.ts | 13 +- .../core/src/__tests__/task-delete-notice.test.ts | 7 +- .../0040_fn_8684_task_lifecycle_outbox.sql | 44 +++++ packages/core/src/postgres/schema-applier.ts | 15 +- packages/core/src/postgres/schema/project.ts | 27 +++ .../__tests__/lifecycle-outbox-writer.test.ts | 201 +++++++++++++++++++++ .../core/src/task-store/archive-lifecycle-2.ts | 93 ++++++++-- packages/core/src/task-store/async-persistence.ts | 26 ++- packages/core/src/task-store/lifecycle-outbox.ts | 46 +++++ 14 files changed, 557 insertions(+), 28 deletions(-) Fusion-Task-Id: FN-8684 Fusion-Task-Lineage: 1869d221-2d8a-48fb-b245-9cd1af56bda0 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8684-task-deleted-outbox-writer.md
Normal file
7
.changeset/fn-8684-task-deleted-outbox-writer.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Add durable transactional task-deletion lifecycle events for PostgreSQL projects.
|
||||||
|
category: feature
|
||||||
|
dev: Registers migration 0040, first-transition claim, and transactional writer seam.
|
||||||
@@ -2329,3 +2329,9 @@ The auto-recovery dispatcher at `packages/engine/src/auto-recovery.ts` (FN-4533)
|
|||||||
### Concurrent soft-delete heartbeat races (FN-8004)
|
### Concurrent soft-delete heartbeat races (FN-8004)
|
||||||
|
|
||||||
A heartbeat `moveTask` failure with the typed `TaskDeletedError` soft-delete message is a benign board miss: the durable agent stays active, clears `lastError` and heartbeat recovery metadata, and emits `agent:heartbeat-move-skipped-soft-delete`. Its audit metadata is structured only: `agentId`, `taskId`, `deletedAt`, `moveAttemptedAt`, and `source`.
|
A heartbeat `moveTask` failure with the typed `TaskDeletedError` soft-delete message is a benign board miss: the durable agent stays active, clears `lastError` and heartbeat recovery metadata, and emits `agent:heartbeat-move-skipped-soft-delete`. Its audit metadata is structured only: `agentId`, `taskId`, `deletedAt`, `moveAttemptedAt`, and `source`.
|
||||||
|
|
||||||
|
## PostgreSQL task-deletion lifecycle outbox
|
||||||
|
|
||||||
|
FN-8684 makes `task:deleted` observable across independently connected PostgreSQL processes. `deleteTaskBackendImpl` claims the first soft-delete transition with `deleted_at IS NULL` inside its transaction. The winner alone writes the run audit and `project.task_lifecycle_events` row; a concurrent loser re-reads the deleted row in the same project-scoped transaction and returns that truthful result without duplicate emit or mailbox effects.
|
||||||
|
|
||||||
|
The outbox insert is atomic with the task mutation. A transactional per-project counter row assigns the event sequence: its lock lasts to commit, allocation order equals commit order, and rollback restores the counter rather than burning a number. Event IDs are deterministic SHA-256-derived opaque IDs, and the fixed payload is IDs/outcomes-only. This writer is currently write-only; consumer cursors, delivery, catch-up, poison handling, and retention are FN-8685.
|
||||||
|
|||||||
@@ -804,3 +804,9 @@ Configuration changes are immutable `project.configuration_revisions` snapshots.
|
|||||||
User-global `~/.fusion/settings.json` history uses the reserved `__fusion_global_configuration__` owner identity rather than the project that initiated the write. Filesystem writes are serialized and stage a durable revision-intent file before replacing settings; a later mutation reconciles an interrupted intent by completing its journal append or restoring the old snapshot. When its revision append fails, the store restores the pre-write raw settings file before rejecting. `TaskStore.rollbackConfiguration()` exactly restores project/global/workflow snapshots; `RoutineStore` and `AutomationStore` expose the same rollback action for their stable-ID resources. Each rollback includes deletion/recreation semantics and appends exactly one forward revision marked `source: "rollback"`, rather than modifying history.
|
User-global `~/.fusion/settings.json` history uses the reserved `__fusion_global_configuration__` owner identity rather than the project that initiated the write. Filesystem writes are serialized and stage a durable revision-intent file before replacing settings; a later mutation reconciles an interrupted intent by completing its journal append or restoring the old snapshot. When its revision append fails, the store restores the pre-write raw settings file before rejecting. `TaskStore.rollbackConfiguration()` exactly restores project/global/workflow snapshots; `RoutineStore` and `AutomationStore` expose the same rollback action for their stable-ID resources. Each rollback includes deletion/recreation semantics and appends exactly one forward revision marked `source: "rollback"`, rather than modifying history.
|
||||||
|
|
||||||
Direct chat tags are stored in the project PostgreSQL schema as `chat_tags` and `chat_session_tags`. Tags are normalized and project-scoped; assignment cleanup does not delete chat sessions.
|
Direct chat tags are stored in the project PostgreSQL schema as `chat_tags` and `chat_session_tags`. Tags are normalized and project-scoped; assignment cleanup does not delete chat sessions.
|
||||||
|
|
||||||
|
## Task deletion lifecycle outbox
|
||||||
|
|
||||||
|
`project.task_lifecycle_events` is the write-only durable source for cross-process `task:deleted` observation. The delete transaction conditionally claims the first `deleted_at` transition; only its winner writes the audit record and one outbox event. A concurrent loser re-reads the project-scoped deleted task in that transaction and has no writer-side audit, event, emit, or mailbox effect.
|
||||||
|
|
||||||
|
Events use `(project_id, seq)` and a deterministic `evt_` SHA-256 identity over project, event type, task ID, and deletion timestamp. `project.task_lifecycle_event_seq` allocates that per-project sequence through a transactional counter upsert. Its lock is held until commit, so allocation order equals commit order, committed rows are in-order and gap-free, and rollback reverts the counter without consuming a number. Payloads contain only task IDs, previous lane/status, deletion timestamp, resurrection/issue action, and actor ID fields. Consumers, cursors, catch-up, retries, and retention are intentionally deferred to FN-8685.
|
||||||
|
|||||||
@@ -81,6 +81,8 @@ import {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
} from "../../postgres/schema-applier.js";
|
} from "../../postgres/schema-applier.js";
|
||||||
import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js";
|
import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js";
|
||||||
import type { PluginSchemaInitHook } from "../../postgres/plugin-schema-hook.js";
|
import type { PluginSchemaInitHook } from "../../postgres/plugin-schema-hook.js";
|
||||||
@@ -95,6 +97,11 @@ const PG_AVAILABLE =
|
|||||||
const pgDescribe = PG_AVAILABLE ? describe : describe.skip;
|
const pgDescribe = PG_AVAILABLE ? describe : describe.skip;
|
||||||
|
|
||||||
describe("schema-applier: immutable migration identities", () => {
|
describe("schema-applier: immutable migration identities", () => {
|
||||||
|
it("registers the task lifecycle outbox after credential selection", () => {
|
||||||
|
expect(TASK_LIFECYCLE_OUTBOX_VERSION).toBe("0040");
|
||||||
|
expect(SCHEMA_BASELINE_VERSION).toBe("0040");
|
||||||
|
});
|
||||||
|
|
||||||
it("keeps monitor and approval isolation assigned to version 0003", () => {
|
it("keeps monitor and approval isolation assigned to version 0003", () => {
|
||||||
expect(MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION).toBe("0003");
|
expect(MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION).toBe("0003");
|
||||||
expect(Number(SCHEMA_BASELINE_VERSION))
|
expect(Number(SCHEMA_BASELINE_VERSION))
|
||||||
@@ -362,6 +369,30 @@ The baseline declares symbol_locks but cannot attach its ownership trigger befor
|
|||||||
installation must therefore prove 0025 leaves the final table forced-RLS with
|
installation must therefore prove 0025 leaves the final table forced-RLS with
|
||||||
its policy and trigger, including actual second-project read/write isolation.
|
its policy and trigger, including actual second-project read/write isolation.
|
||||||
*/
|
*/
|
||||||
|
async function assertTaskLifecycleOutboxOwnershipContract(ctx: TestContext): Promise<void> {
|
||||||
|
const catalog = (await ctx.db.execute(sql`
|
||||||
|
SELECT c.relname AS table_name, c.relrowsecurity AS rls, c.relforcerowsecurity AS forced,
|
||||||
|
EXISTS (
|
||||||
|
SELECT 1 FROM pg_policies
|
||||||
|
WHERE schemaname = 'project' AND tablename = c.relname
|
||||||
|
AND policyname = 'fusion_project_isolation'
|
||||||
|
) AS policy,
|
||||||
|
EXISTS (
|
||||||
|
SELECT 1 FROM pg_trigger
|
||||||
|
WHERE tgrelid = c.oid AND tgname = 'fusion_assign_project_id' AND NOT tgisinternal
|
||||||
|
) AS trigger
|
||||||
|
FROM pg_class c
|
||||||
|
JOIN pg_namespace n ON n.oid = c.relnamespace
|
||||||
|
WHERE n.nspname = 'project'
|
||||||
|
AND c.relname IN ('task_lifecycle_events', 'task_lifecycle_event_seq')
|
||||||
|
ORDER BY c.relname
|
||||||
|
`)) as unknown as Array<{ table_name: string; rls: boolean; forced: boolean; policy: boolean; trigger: boolean }>;
|
||||||
|
expect(catalog).toEqual([
|
||||||
|
{ table_name: 'task_lifecycle_event_seq', rls: true, forced: true, policy: true, trigger: true },
|
||||||
|
{ table_name: 'task_lifecycle_events', rls: true, forced: true, policy: true, trigger: true },
|
||||||
|
]);
|
||||||
|
}
|
||||||
|
|
||||||
async function assertSymbolLocksOwnershipContract(ctx: TestContext): Promise<void> {
|
async function assertSymbolLocksOwnershipContract(ctx: TestContext): Promise<void> {
|
||||||
const catalog = (await ctx.db.execute(sql`
|
const catalog = (await ctx.db.execute(sql`
|
||||||
SELECT c.relrowsecurity AS rls, c.relforcerowsecurity AS forced,
|
SELECT c.relrowsecurity AS rls, c.relforcerowsecurity AS forced,
|
||||||
@@ -683,7 +714,7 @@ pgDescribe("schema-applier: VAL-SCHEMA-001 final-schema parity (table counts)",
|
|||||||
ctx = null;
|
ctx = null;
|
||||||
});
|
});
|
||||||
|
|
||||||
it("creates all 96 project tables, 17 central tables, 1 archive table", async () => {
|
it("creates all 100 project tables, 17 central tables, 1 archive table", async () => {
|
||||||
ctx = await setupFreshDb();
|
ctx = await setupFreshDb();
|
||||||
// FNXC:PostgresCutover 2026-07-05-15:55: apply the BASELINE only.
|
// FNXC:PostgresCutover 2026-07-05-15:55: apply the BASELINE only.
|
||||||
// applySchemaBaseline now runs the plugin schema-init hooks by default,
|
// applySchemaBaseline now runs the plugin schema-init hooks by default,
|
||||||
@@ -703,9 +734,10 @@ pgDescribe("schema-applier: VAL-SCHEMA-001 final-schema parity (table counts)",
|
|||||||
// + 1 configuration_revisions (FNXC:ConfigVersioning 2026-07-18-14:00)
|
// + 1 configuration_revisions (FNXC:ConfigVersioning 2026-07-18-14:00)
|
||||||
// + 2 ideation_sessions/ideation_candidates (FNXC:Ideation 2026-07-18-13:25 / FN-8295)
|
// + 2 ideation_sessions/ideation_candidates (FNXC:Ideation 2026-07-18-13:25 / FN-8295)
|
||||||
// + 1 task_verification_requests + 1 durable symbol_locks table (FN-8305)
|
// + 1 task_verification_requests + 1 durable symbol_locks table (FN-8305)
|
||||||
// + 1 mission_lineage_stops (FNXC:MissionLineageBudget FN-8543 / migration 0035).
|
// + 1 mission_lineage_stops (FNXC:MissionLineageBudget FN-8543 / migration 0035)
|
||||||
|
// + 2 task lifecycle outbox tables (FN-8684 migration 0040).
|
||||||
// Plugin tables are added separately by the hook.
|
// Plugin tables are added separately by the hook.
|
||||||
expect(bySchema.project).toBe(98);
|
expect(bySchema.project).toBe(100);
|
||||||
/*
|
/*
|
||||||
FNXC:CapacityModel 2026-07-29-08:10 (drop the cross-project cap — table half):
|
FNXC:CapacityModel 2026-07-29-08:10 (drop the cross-project cap — table half):
|
||||||
17, not 18: `central.global_concurrency` is dropped by migration 0037. A fresh
|
17, not 18: `central.global_concurrency` is dropped by migration 0037. A fresh
|
||||||
@@ -733,6 +765,40 @@ pgDescribe("schema-applier: VAL-SCHEMA-001 final-schema parity (table counts)",
|
|||||||
expect(second.applied).toBe(false);
|
expect(second.applied).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-11:02:
|
||||||
|
Migration discovery is intentionally disabled, so the writer tables require proof at
|
||||||
|
their fresh, upgrade, and manually-repaired re-execution surfaces. Both tables must
|
||||||
|
retain the ownership contract because lifecycle events cross process boundaries.
|
||||||
|
*/
|
||||||
|
it("installs lifecycle outbox ownership on fresh databases and re-executes its SQL safely", async () => {
|
||||||
|
ctx = await setupFreshDb();
|
||||||
|
await expect(applySchemaBaseline(ctx.db, { pluginHooks: [] })).resolves.toMatchObject({ applied: true });
|
||||||
|
await assertTaskLifecycleOutboxOwnershipContract(ctx);
|
||||||
|
|
||||||
|
const migrationSql = readFileSync(
|
||||||
|
fileURLToPath(new URL("../../postgres/migrations/0040_fn_8684_task_lifecycle_outbox.sql", import.meta.url)),
|
||||||
|
"utf8",
|
||||||
|
);
|
||||||
|
await expect(ctx.db.execute(sql.raw(migrationSql))).resolves.toBeDefined();
|
||||||
|
await assertTaskLifecycleOutboxOwnershipContract(ctx);
|
||||||
|
await expect(applySchemaBaseline(ctx.db, { pluginHooks: [] })).resolves.toEqual({ applied: false, pluginHooksRun: 0 });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("upgrades a database recorded through 0039 with both lifecycle outbox tables", async () => {
|
||||||
|
ctx = await setupFreshDb();
|
||||||
|
await applySchemaBaseline(ctx.db, { pluginHooks: [] });
|
||||||
|
await ctx.db.execute(sql.raw(`
|
||||||
|
DELETE FROM public.fusion_schema_migrations WHERE version = '0040';
|
||||||
|
DROP TABLE project.task_lifecycle_events;
|
||||||
|
DROP TABLE project.task_lifecycle_event_seq;
|
||||||
|
`));
|
||||||
|
|
||||||
|
await expect(applySchemaBaseline(ctx.db, { pluginHooks: [] })).resolves.toEqual({ applied: true, pluginHooksRun: 0 });
|
||||||
|
expect(await getAppliedMigrations(ctx.db)).toContain(TASK_LIFECYCLE_OUTBOX_VERSION);
|
||||||
|
await assertTaskLifecycleOutboxOwnershipContract(ctx);
|
||||||
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:PostgresMigrationColumnCoverage 2026-07-14-13:17:
|
FNXC:PostgresMigrationColumnCoverage 2026-07-14-13:17:
|
||||||
A cluster that already recorded migrations through 0006 must receive every late SQLite column before cutover retries. This is the production failure shape: the initial copy is blocked while the target schema is otherwise fully initialized.
|
A cluster that already recorded migrations through 0006 must receive every late SQLite column before cutover retries. This is the production failure shape: the initial copy is blocked while the target schema is otherwise fully initialized.
|
||||||
@@ -1679,6 +1745,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
]);
|
]);
|
||||||
expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false);
|
expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false);
|
||||||
});
|
});
|
||||||
@@ -1743,6 +1811,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -1940,6 +2010,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -2018,6 +2090,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -2096,6 +2170,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => {
|
|||||||
CHAT_SESSION_TAGS_VERSION,
|
CHAT_SESSION_TAGS_VERSION,
|
||||||
DROP_GLOBAL_CONCURRENCY_VERSION,
|
DROP_GLOBAL_CONCURRENCY_VERSION,
|
||||||
MISSION_TASK_PREFIX_VERSION,
|
MISSION_TASK_PREFIX_VERSION,
|
||||||
|
CREDENTIAL_INSTANCE_SELECTION_VERSION,
|
||||||
|
TASK_LIFECYCLE_OUTBOX_VERSION,
|
||||||
]);
|
]);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -39,13 +39,23 @@ let pgRow: TaskRowShape | null = null;
|
|||||||
|
|
||||||
vi.mock("../task-store/async-persistence.js", () => ({
|
vi.mock("../task-store/async-persistence.js", () => ({
|
||||||
readTaskRow: vi.fn(async () => pgRow),
|
readTaskRow: vi.fn(async () => pgRow),
|
||||||
softDeleteTaskRowInTransaction: vi.fn(async () => undefined),
|
readTaskRowInTransaction: vi.fn(async () => pgRow),
|
||||||
|
softDeleteTaskRowInTransaction: vi.fn(async () => true),
|
||||||
}));
|
}));
|
||||||
vi.mock("../task-store/async-lifecycle.js", () => ({
|
vi.mock("../task-store/async-lifecycle.js", () => ({
|
||||||
findLiveLineageChildren: vi.fn(async () => [] as string[]),
|
findLiveLineageChildren: vi.fn(async () => [] as string[]),
|
||||||
projectPartition: vi.fn(() => undefined),
|
projectPartition: vi.fn(() => undefined),
|
||||||
removeLineageReferences: vi.fn(async () => undefined),
|
removeLineageReferences: vi.fn(async () => undefined),
|
||||||
}));
|
}));
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-11:02:
|
||||||
|
This in-memory delete harness has no PostgreSQL transaction executor. Mock the
|
||||||
|
transaction-scoped writer at its module boundary so attribution tests retain their
|
||||||
|
narrow audit focus while production deletes always use the durable outbox writer.
|
||||||
|
*/
|
||||||
|
vi.mock("../task-store/lifecycle-outbox.js", () => ({
|
||||||
|
appendTaskLifecycleEventInTransaction: vi.fn(async () => ({ seq: "1", eventId: "test-event" })),
|
||||||
|
}));
|
||||||
vi.mock("../async-mission-store-queries.js", () => ({
|
vi.mock("../async-mission-store-queries.js", () => ({
|
||||||
getFeatureByTaskId: vi.fn(async () => null),
|
getFeatureByTaskId: vi.fn(async () => null),
|
||||||
unlinkFeatureFromTaskId: vi.fn(async () => undefined),
|
unlinkFeatureFromTaskId: vi.fn(async () => undefined),
|
||||||
|
|||||||
@@ -31,13 +31,23 @@ let lineageChildIds: string[] = [];
|
|||||||
|
|
||||||
vi.mock("../task-store/async-persistence.js", () => ({
|
vi.mock("../task-store/async-persistence.js", () => ({
|
||||||
readTaskRow: vi.fn(async () => pgRow),
|
readTaskRow: vi.fn(async () => pgRow),
|
||||||
softDeleteTaskRowInTransaction: vi.fn(async () => undefined),
|
readTaskRowInTransaction: vi.fn(async () => pgRow),
|
||||||
|
softDeleteTaskRowInTransaction: vi.fn(async () => true),
|
||||||
}));
|
}));
|
||||||
vi.mock("../task-store/async-lifecycle.js", () => ({
|
vi.mock("../task-store/async-lifecycle.js", () => ({
|
||||||
findLiveLineageChildren: vi.fn(async () => lineageChildIds),
|
findLiveLineageChildren: vi.fn(async () => lineageChildIds),
|
||||||
projectPartition: vi.fn(() => undefined),
|
projectPartition: vi.fn(() => undefined),
|
||||||
removeLineageReferences: vi.fn(async () => undefined),
|
removeLineageReferences: vi.fn(async () => undefined),
|
||||||
}));
|
}));
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-11:02:
|
||||||
|
This in-memory delete harness has no PostgreSQL transaction executor. Mock the
|
||||||
|
transaction-scoped writer at its module boundary so deletion-gate tests remain
|
||||||
|
focused while the real backend always persists the lifecycle event transactionally.
|
||||||
|
*/
|
||||||
|
vi.mock("../task-store/lifecycle-outbox.js", () => ({
|
||||||
|
appendTaskLifecycleEventInTransaction: vi.fn(async () => ({ seq: "1", eventId: "test-event" })),
|
||||||
|
}));
|
||||||
vi.mock("../async-mission-store-queries.js", () => ({
|
vi.mock("../async-mission-store-queries.js", () => ({
|
||||||
getFeatureByTaskId: vi.fn(async () => null),
|
getFeatureByTaskId: vi.fn(async () => null),
|
||||||
unlinkFeatureFromTaskId: vi.fn(async () => undefined),
|
unlinkFeatureFromTaskId: vi.fn(async () => undefined),
|
||||||
@@ -94,6 +104,7 @@ function makeDeleteStore(task: Task, children: string[] = []) {
|
|||||||
auditEvents.push(event);
|
auditEvents.push(event);
|
||||||
}),
|
}),
|
||||||
makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`),
|
makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`),
|
||||||
|
laneCache: { invalidate: vi.fn() },
|
||||||
withTaskLock: vi.fn(async (_id: string, fn: () => Promise<unknown>) => fn()),
|
withTaskLock: vi.fn(async (_id: string, fn: () => Promise<unknown>) => fn()),
|
||||||
cleanupBranchForTask: vi.fn(async () => [] as string[]),
|
cleanupBranchForTask: vi.fn(async () => [] as string[]),
|
||||||
clearNearDuplicateReferencesToFailSoft: vi.fn(async () => undefined),
|
clearNearDuplicateReferencesToFailSoft: vi.fn(async () => undefined),
|
||||||
|
|||||||
@@ -35,7 +35,11 @@ import { beforeEach, describe, expect, it, vi } from "vitest";
|
|||||||
|
|
||||||
vi.mock("../task-store/async-persistence.js", () => ({
|
vi.mock("../task-store/async-persistence.js", () => ({
|
||||||
readTaskRow: vi.fn(async () => pgRow),
|
readTaskRow: vi.fn(async () => pgRow),
|
||||||
softDeleteTaskRowInTransaction: vi.fn(async () => undefined),
|
readTaskRowInTransaction: vi.fn(async () => pgRow),
|
||||||
|
softDeleteTaskRowInTransaction: vi.fn(async () => true),
|
||||||
|
}));
|
||||||
|
vi.mock("../task-store/lifecycle-outbox.js", () => ({
|
||||||
|
appendTaskLifecycleEventInTransaction: vi.fn(async () => undefined),
|
||||||
}));
|
}));
|
||||||
vi.mock("../task-store/async-lifecycle.js", () => ({
|
vi.mock("../task-store/async-lifecycle.js", () => ({
|
||||||
findLiveLineageChildren: vi.fn(async () => [] as string[]),
|
findLiveLineageChildren: vi.fn(async () => [] as string[]),
|
||||||
@@ -120,6 +124,7 @@ function makePgStore(task: Task) {
|
|||||||
recordRunAuditEventBackend: vi.fn(async () => undefined),
|
recordRunAuditEventBackend: vi.fn(async () => undefined),
|
||||||
makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`),
|
makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`),
|
||||||
emit: vi.fn(),
|
emit: vi.fn(),
|
||||||
|
laneCache: { invalidate: vi.fn() },
|
||||||
/* `deleteTaskIf` wraps the conditional delete in the per-task lock; the fake runs
|
/* `deleteTaskIf` wraps the conditional delete in the per-task lock; the fake runs
|
||||||
the body inline so the predicate/short-circuit paths are exercised for real. */
|
the body inline so the predicate/short-circuit paths are exercised for real. */
|
||||||
withTaskLock: vi.fn(async (_id: string, fn: () => Promise<unknown>) => fn()),
|
withTaskLock: vi.fn(async (_id: string, fn: () => Promise<unknown>) => fn()),
|
||||||
|
|||||||
@@ -0,0 +1,44 @@
|
|||||||
|
-- FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
-- PostgreSQL needs a durable, cross-process observation record for task:deleted after
|
||||||
|
-- FN-8683 removed the unreachable SQLite polling replica. The event and delete state
|
||||||
|
-- must commit together; the transactional counter avoids MAX(seq)+1 collisions and,
|
||||||
|
-- unlike a SEQUENCE, rolls back with a failed delete. Schema-applier versions run once,
|
||||||
|
-- while the policy/trigger guards keep repaired manual re-runs executable.
|
||||||
|
CREATE TABLE IF NOT EXISTS project.task_lifecycle_events (
|
||||||
|
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
|
||||||
|
seq bigint NOT NULL,
|
||||||
|
event_id text NOT NULL,
|
||||||
|
event_type text NOT NULL,
|
||||||
|
task_id text NOT NULL,
|
||||||
|
occurred_at text NOT NULL,
|
||||||
|
created_at text NOT NULL,
|
||||||
|
payload jsonb NOT NULL,
|
||||||
|
PRIMARY KEY (project_id, seq),
|
||||||
|
UNIQUE (project_id, event_id),
|
||||||
|
CONSTRAINT task_lifecycle_events_type_check CHECK (event_type IN ('task:deleted'))
|
||||||
|
);
|
||||||
|
CREATE INDEX IF NOT EXISTS "idxTaskLifecycleEventsTask" ON project.task_lifecycle_events(project_id, task_id);
|
||||||
|
ALTER TABLE project.task_lifecycle_events ENABLE ROW LEVEL SECURITY;
|
||||||
|
ALTER TABLE project.task_lifecycle_events FORCE ROW LEVEL SECURITY;
|
||||||
|
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_events;
|
||||||
|
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_events
|
||||||
|
USING (current_setting('fusion.project_bypass', true) = 'on' OR project_id = current_setting('fusion.project_id', true))
|
||||||
|
WITH CHECK (current_setting('fusion.project_bypass', true) = 'on' OR project_id = current_setting('fusion.project_id', true));
|
||||||
|
DROP TRIGGER IF EXISTS fusion_assign_project_id ON project.task_lifecycle_events;
|
||||||
|
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_events
|
||||||
|
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS project.task_lifecycle_event_seq (
|
||||||
|
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
|
||||||
|
last_seq bigint NOT NULL DEFAULT 0,
|
||||||
|
PRIMARY KEY (project_id)
|
||||||
|
);
|
||||||
|
ALTER TABLE project.task_lifecycle_event_seq ENABLE ROW LEVEL SECURITY;
|
||||||
|
ALTER TABLE project.task_lifecycle_event_seq FORCE ROW LEVEL SECURITY;
|
||||||
|
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_event_seq;
|
||||||
|
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_event_seq
|
||||||
|
USING (current_setting('fusion.project_bypass', true) = 'on' OR project_id = current_setting('fusion.project_id', true))
|
||||||
|
WITH CHECK (current_setting('fusion.project_bypass', true) = 'on' OR project_id = current_setting('fusion.project_id', true));
|
||||||
|
DROP TRIGGER IF EXISTS fusion_assign_project_id ON project.task_lifecycle_event_seq;
|
||||||
|
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_event_seq
|
||||||
|
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();
|
||||||
@@ -54,7 +54,8 @@ FNXC:MissionTaskPrefix 2026-07-30-21:10 (rebase onto migrated main):
|
|||||||
SCHEMA_BASELINE_VERSION advances to 0038 for optional per-mission task_prefix — 0037 is the
|
SCHEMA_BASELINE_VERSION advances to 0038 for optional per-mission task_prefix — 0037 is the
|
||||||
capacity-model table drop that landed while this PR was open.
|
capacity-model table drop that landed while this PR was open.
|
||||||
*/
|
*/
|
||||||
export const SCHEMA_BASELINE_VERSION = "0039";
|
/* FNXC:LifecycleOutbox 2026-08-01-10:33: advance the explicit schema ceiling so both lifecycle events and their transactional counter exist before delete writers run. */
|
||||||
|
export const SCHEMA_BASELINE_VERSION = "0040";
|
||||||
/** FNXC:SymbolLock 2026-07-20-10:00: upgrades need durable task declarations before admission resolves symbols. */
|
/** FNXC:SymbolLock 2026-07-20-10:00: upgrades need durable task declarations before admission resolves symbols. */
|
||||||
export const TASK_DECLARED_SYMBOLS_VERSION = "0028";
|
export const TASK_DECLARED_SYMBOLS_VERSION = "0028";
|
||||||
const INITIAL_SCHEMA_VERSION = "0000";
|
const INITIAL_SCHEMA_VERSION = "0000";
|
||||||
@@ -175,6 +176,8 @@ is landed on real databases and its identity is immutable, per the MONITOR_APPRO
|
|||||||
export const MISSION_TASK_PREFIX_VERSION = "0038";
|
export const MISSION_TASK_PREFIX_VERSION = "0038";
|
||||||
/** FNXC:CredentialInstanceSelection 2026-08-01-05:43: explicit registration prevents migration 0039 from being silently skipped on upgraded PostgreSQL databases. */
|
/** FNXC:CredentialInstanceSelection 2026-08-01-05:43: explicit registration prevents migration 0039 from being silently skipped on upgraded PostgreSQL databases. */
|
||||||
export const CREDENTIAL_INSTANCE_SELECTION_VERSION = "0039";
|
export const CREDENTIAL_INSTANCE_SELECTION_VERSION = "0039";
|
||||||
|
/** FNXC:LifecycleOutbox 2026-08-01-10:33: migrations are explicitly registered; 0040 installs both writer tables for fresh and upgraded projects. */
|
||||||
|
export const TASK_LIFECYCLE_OUTBOX_VERSION = "0040";
|
||||||
|
|
||||||
/** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */
|
/** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */
|
||||||
export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained";
|
export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained";
|
||||||
@@ -388,6 +391,7 @@ const CHAT_SESSION_TAGS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0036_chat_session
|
|||||||
const DROP_GLOBAL_CONCURRENCY_MIGRATION_PATH = join(MIGRATIONS_DIR, "0037_drop_global_concurrency.sql");
|
const DROP_GLOBAL_CONCURRENCY_MIGRATION_PATH = join(MIGRATIONS_DIR, "0037_drop_global_concurrency.sql");
|
||||||
const MISSION_TASK_PREFIX_MIGRATION_PATH = join(MIGRATIONS_DIR, "0038_mission_task_prefix.sql");
|
const MISSION_TASK_PREFIX_MIGRATION_PATH = join(MIGRATIONS_DIR, "0038_mission_task_prefix.sql");
|
||||||
const CREDENTIAL_INSTANCE_SELECTION_MIGRATION_PATH = join(MIGRATIONS_DIR, "0039_fn_8660_credential_instance_selection.sql");
|
const CREDENTIAL_INSTANCE_SELECTION_MIGRATION_PATH = join(MIGRATIONS_DIR, "0039_fn_8660_credential_instance_selection.sql");
|
||||||
|
const TASK_LIFECYCLE_OUTBOX_MIGRATION_PATH = join(MIGRATIONS_DIR, "0040_fn_8684_task_lifecycle_outbox.sql");
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Ensure the migration bookkeeping table exists. Lives in the public schema so
|
* Ensure the migration bookkeeping table exists. Lives in the public schema so
|
||||||
@@ -497,6 +501,7 @@ export async function applySchemaBaseline(
|
|||||||
const dropGlobalConcurrencyAlreadyApplied = applied.includes(DROP_GLOBAL_CONCURRENCY_VERSION);
|
const dropGlobalConcurrencyAlreadyApplied = applied.includes(DROP_GLOBAL_CONCURRENCY_VERSION);
|
||||||
const missionTaskPrefixAlreadyApplied = applied.includes(MISSION_TASK_PREFIX_VERSION);
|
const missionTaskPrefixAlreadyApplied = applied.includes(MISSION_TASK_PREFIX_VERSION);
|
||||||
const credentialInstanceSelectionAlreadyApplied = applied.includes(CREDENTIAL_INSTANCE_SELECTION_VERSION);
|
const credentialInstanceSelectionAlreadyApplied = applied.includes(CREDENTIAL_INSTANCE_SELECTION_VERSION);
|
||||||
|
const taskLifecycleOutboxAlreadyApplied = applied.includes(TASK_LIFECYCLE_OUTBOX_VERSION);
|
||||||
assertBinaryNotOlderThanDatabase(applied);
|
assertBinaryNotOlderThanDatabase(applied);
|
||||||
let schemaChanged = false;
|
let schemaChanged = false;
|
||||||
|
|
||||||
@@ -1054,6 +1059,14 @@ export async function applySchemaBaseline(
|
|||||||
schemaChanged = true;
|
schemaChanged = true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* FNXC:LifecycleOutbox 2026-08-01-10:33: explicit application is required because migration discovery is intentionally disabled; the events table and counter must arrive together. */
|
||||||
|
if (!taskLifecycleOutboxAlreadyApplied) {
|
||||||
|
const migrationSql = await readFile(TASK_LIFECYCLE_OUTBOX_MIGRATION_PATH, "utf8");
|
||||||
|
await tx.execute(sql.raw(migrationSql));
|
||||||
|
await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_LIFECYCLE_OUTBOX_VERSION}) ON CONFLICT (version) DO NOTHING`);
|
||||||
|
schemaChanged = true;
|
||||||
|
}
|
||||||
|
|
||||||
return { applied: schemaChanged, pluginHooksRun: pluginHooks.length };
|
return { applied: schemaChanged, pluginHooksRun: pluginHooks.length };
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -564,6 +564,33 @@ export const symbolLocks = projectSchema.table("symbol_locks", {
|
|||||||
index("idxSymbolLocksExpiry").on(t.status, t.expiresAt),
|
index("idxSymbolLocksExpiry").on(t.status, t.expiresAt),
|
||||||
]);
|
]);
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
These project-scoped rows make task:deleted observable across PostgreSQL processes after
|
||||||
|
FN-8683 removed unreachable SQLite polling. The writer inserts them with the soft-delete;
|
||||||
|
the counter row serializes allocation without MAX(seq)+1 races and rolls back on failure.
|
||||||
|
*/
|
||||||
|
export const taskLifecycleEvents = projectSchema.table("task_lifecycle_events", {
|
||||||
|
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
|
||||||
|
seq: bigint("seq", { mode: "bigint" }).notNull(),
|
||||||
|
eventId: text("event_id").notNull(),
|
||||||
|
eventType: text("event_type").notNull(),
|
||||||
|
taskId: text("task_id").notNull(),
|
||||||
|
occurredAt: text("occurred_at").notNull(),
|
||||||
|
createdAt: text("created_at").notNull(),
|
||||||
|
payload: jsonb("payload").notNull(),
|
||||||
|
}, (t) => [
|
||||||
|
primaryKey({ columns: [t.projectId, t.seq] }),
|
||||||
|
unique("task_lifecycle_events_project_event_unique").on(t.projectId, t.eventId),
|
||||||
|
check("task_lifecycle_events_type_check", sql`${t.eventType} IN ('task:deleted')`),
|
||||||
|
index("idxTaskLifecycleEventsTask").on(t.projectId, t.taskId),
|
||||||
|
]);
|
||||||
|
|
||||||
|
export const taskLifecycleEventSeq = projectSchema.table("task_lifecycle_event_seq", {
|
||||||
|
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
|
||||||
|
lastSeq: bigint("last_seq", { mode: "bigint" }).notNull().default(sql`0`),
|
||||||
|
}, (t) => [primaryKey({ columns: [t.projectId] })]);
|
||||||
|
|
||||||
// ── Workflow step definitions ────────────────────────────────────────
|
// ── Workflow step definitions ────────────────────────────────────────
|
||||||
export const workflowSteps = projectSchema.table("workflow_steps", {
|
export const workflowSteps = projectSchema.table("workflow_steps", {
|
||||||
id: text("id").primaryKey(),
|
id: text("id").primaryKey(),
|
||||||
|
|||||||
@@ -0,0 +1,201 @@
|
|||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:51:
|
||||||
|
The lifecycle outbox is a PostgreSQL cross-process contract, so these tests use real independent
|
||||||
|
connections rather than mocks. In particular the barrier holds two live snapshots before either
|
||||||
|
conditional delete claim, which proves the old TOCTOU duplication cannot return unnoticed.
|
||||||
|
*/
|
||||||
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import { and, eq } from "drizzle-orm";
|
||||||
|
import { readFile } from "node:fs/promises";
|
||||||
|
import {
|
||||||
|
createTaskStoreForTest,
|
||||||
|
pgDescribe,
|
||||||
|
type PgTestHarness,
|
||||||
|
} from "../../__test-utils__/pg-test-harness.js";
|
||||||
|
import { TaskStore } from "../../store.js";
|
||||||
|
import { createAsyncDataLayer, type AsyncDataLayer } from "../../postgres/data-layer.js";
|
||||||
|
import { createConnectionSetFromUrl } from "../../postgres/connection.js";
|
||||||
|
import type { ResolvedBackend } from "../../postgres/backend-resolver.js";
|
||||||
|
import * as schema from "../../postgres/schema/index.js";
|
||||||
|
import { countRunAuditEvents } from "../async-audit.js";
|
||||||
|
import { softDeleteTaskRowInTransaction } from "../async-persistence.js";
|
||||||
|
import {
|
||||||
|
makeTaskLifecycleEventId,
|
||||||
|
type TaskDeletedLifecyclePayload,
|
||||||
|
} from "../lifecycle-outbox.js";
|
||||||
|
import { registerTaskDeleteNoticeMailbox } from "../../task-delete-notice.js";
|
||||||
|
|
||||||
|
const pgTest = pgDescribe;
|
||||||
|
const PROJECT_ID = "__legacy_unscoped__";
|
||||||
|
|
||||||
|
type DeleteTestStore = TaskStore & {
|
||||||
|
__beforeDeleteClaimForTest?: (taskId: string) => void | Promise<void>;
|
||||||
|
__afterLifecycleOutboxWriteForTest?: () => void | Promise<void>;
|
||||||
|
};
|
||||||
|
|
||||||
|
function deferred<T = void>() {
|
||||||
|
let resolve!: (value: T | PromiseLike<T>) => void;
|
||||||
|
const promise = new Promise<T>((resolvePromise) => { resolve = resolvePromise; });
|
||||||
|
return { promise, resolve };
|
||||||
|
}
|
||||||
|
|
||||||
|
function lifecycleRows(h: PgTestHarness, taskId: string) {
|
||||||
|
return h.layer.db.select().from(schema.project.taskLifecycleEvents).where(and(
|
||||||
|
eq(schema.project.taskLifecycleEvents.projectId, PROJECT_ID),
|
||||||
|
eq(schema.project.taskLifecycleEvents.taskId, taskId),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
async function newIndependentStore(h: PgTestHarness): Promise<{ store: TaskStore; layer: AsyncDataLayer }> {
|
||||||
|
const backend: ResolvedBackend = {
|
||||||
|
mode: "external", runtimeUrl: h.testUrl, migrationUrl: h.testUrl, migrationUrlOverridden: false,
|
||||||
|
};
|
||||||
|
const connections = await createConnectionSetFromUrl(backend, { poolMax: 1, connectTimeoutSeconds: 5 });
|
||||||
|
const layer = createAsyncDataLayer(connections);
|
||||||
|
return { store: new TaskStore(h.rootDir, undefined, { asyncLayer: layer }), layer };
|
||||||
|
}
|
||||||
|
|
||||||
|
function attachMailbox(store: TaskStore, notices: string[]): void {
|
||||||
|
registerTaskDeleteNoticeMailbox(store, {
|
||||||
|
sendMessageOnce: vi.fn(async (_input, key) => { notices.push(key); return { inserted: true }; }),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
const payload: TaskDeletedLifecyclePayload = {
|
||||||
|
taskId: "FN-1", previousColumn: "todo", previousStatus: null,
|
||||||
|
deletedAt: "2026-08-01T10:33:00.000Z", allowResurrection: false,
|
||||||
|
githubIssueAction: null, deletedBy: null,
|
||||||
|
};
|
||||||
|
|
||||||
|
describe("task lifecycle outbox identity", () => {
|
||||||
|
it("uses a deterministic opaque project-scoped identity and fixed ids-only payload", () => {
|
||||||
|
const input = ["project-a", "task:deleted", payload.taskId, payload.deletedAt] as const;
|
||||||
|
expect(makeTaskLifecycleEventId(...input)).toBe(makeTaskLifecycleEventId(...input));
|
||||||
|
expect(makeTaskLifecycleEventId(...input)).not.toBe(makeTaskLifecycleEventId("project-b", ...input.slice(1)));
|
||||||
|
expect(Object.keys(payload).sort()).toEqual([
|
||||||
|
"allowResurrection", "deletedAt", "deletedBy", "githubIssueAction",
|
||||||
|
"previousColumn", "previousStatus", "taskId",
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
pgTest("transactional task:deleted lifecycle outbox writer (PostgreSQL)", () => {
|
||||||
|
let h: PgTestHarness | undefined;
|
||||||
|
let second: Awaited<ReturnType<typeof newIndependentStore>> | undefined;
|
||||||
|
|
||||||
|
afterEach(async () => {
|
||||||
|
await second?.store.close().catch(() => {});
|
||||||
|
await second?.layer.close().catch(() => {});
|
||||||
|
second = undefined;
|
||||||
|
await h?.teardown();
|
||||||
|
h = undefined;
|
||||||
|
});
|
||||||
|
|
||||||
|
it("wins exactly one deterministic two-store same-task race and returns the re-read deletion", async () => {
|
||||||
|
h = await createTaskStoreForTest({ prefix: "lifecycle_outbox_race", copyFromGolden: true });
|
||||||
|
second = await newIndependentStore(h);
|
||||||
|
const task = await h.store.createTask({ description: "same-task race" });
|
||||||
|
const notices: string[] = [];
|
||||||
|
attachMailbox(h.store, notices);
|
||||||
|
attachMailbox(second.store, notices);
|
||||||
|
const arrivedA = deferred();
|
||||||
|
const arrivedB = deferred();
|
||||||
|
const release = deferred();
|
||||||
|
(h.store as DeleteTestStore).__beforeDeleteClaimForTest = async () => { arrivedA.resolve(); await release.promise; };
|
||||||
|
(second.store as DeleteTestStore).__beforeDeleteClaimForTest = async () => { arrivedB.resolve(); await release.promise; };
|
||||||
|
const emitA = vi.spyOn(h.store, "emit");
|
||||||
|
const emitB = vi.spyOn(second.store, "emit");
|
||||||
|
|
||||||
|
const deletingA = h.store.deleteTask(task.id);
|
||||||
|
const deletingB = second.store.deleteTask(task.id);
|
||||||
|
await Promise.all([arrivedA.promise, arrivedB.promise]);
|
||||||
|
release.resolve();
|
||||||
|
const [resultA, resultB] = await Promise.all([deletingA, deletingB]);
|
||||||
|
|
||||||
|
const rows = await lifecycleRows(h, task.id);
|
||||||
|
expect(rows).toHaveLength(1);
|
||||||
|
expect(await countRunAuditEvents(h.layer.db, { taskId: task.id, mutationType: "task:deleted" })).toBe(1);
|
||||||
|
expect(notices).toHaveLength(1);
|
||||||
|
expect([emitA, emitB].flatMap((spy) => spy.mock.calls).filter(([event]) => event === "task:deleted")).toHaveLength(1);
|
||||||
|
expect(resultA.deletedAt).toBeTruthy();
|
||||||
|
expect(resultB.deletedAt).toBe(resultA.deletedAt);
|
||||||
|
await expect(h.layer.transactionImmediate((tx) => softDeleteTaskRowInTransaction(tx, task.id, new Date().toISOString(), false, PROJECT_ID, true))).resolves.toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reports only the claim winner as deleted in a deterministic two-store deleteTaskIf race", async () => {
|
||||||
|
h = await createTaskStoreForTest({ prefix: "lifecycle_outbox_conditional_race", copyFromGolden: true });
|
||||||
|
second = await newIndependentStore(h);
|
||||||
|
const task = await h.store.createTask({ description: "same-task conditional race" });
|
||||||
|
const notices: string[] = [];
|
||||||
|
attachMailbox(h.store, notices);
|
||||||
|
attachMailbox(second.store, notices);
|
||||||
|
const arrivedA = deferred();
|
||||||
|
const arrivedB = deferred();
|
||||||
|
const release = deferred();
|
||||||
|
(h.store as DeleteTestStore).__beforeDeleteClaimForTest = async () => { arrivedA.resolve(); await release.promise; };
|
||||||
|
(second.store as DeleteTestStore).__beforeDeleteClaimForTest = async () => { arrivedB.resolve(); await release.promise; };
|
||||||
|
const emitA = vi.spyOn(h.store, "emit");
|
||||||
|
const emitB = vi.spyOn(second.store, "emit");
|
||||||
|
|
||||||
|
const deletingA = h.store.deleteTaskIf(task.id, () => true);
|
||||||
|
const deletingB = second.store.deleteTaskIf(task.id, () => true);
|
||||||
|
await Promise.all([arrivedA.promise, arrivedB.promise]);
|
||||||
|
release.resolve();
|
||||||
|
const [resultA, resultB] = await Promise.all([deletingA, deletingB]);
|
||||||
|
|
||||||
|
expect([resultA.deleted, resultB.deleted].filter(Boolean)).toHaveLength(1);
|
||||||
|
expect(resultA.task.deletedAt).toBeTruthy();
|
||||||
|
expect(resultB.task.deletedAt).toBe(resultA.task.deletedAt);
|
||||||
|
expect(await lifecycleRows(h, task.id)).toHaveLength(1);
|
||||||
|
expect(await countRunAuditEvents(h.layer.db, { taskId: task.id, mutationType: "task:deleted" })).toBe(1);
|
||||||
|
expect(notices).toHaveLength(1);
|
||||||
|
expect([emitA, emitB].flatMap((spy) => spy.mock.calls).filter(([event]) => event === "task:deleted")).toHaveLength(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("rolls back the soft-delete, audit, outbox, and transactional counter together", async () => {
|
||||||
|
h = await createTaskStoreForTest({ prefix: "lifecycle_outbox_atomic", copyFromGolden: true });
|
||||||
|
const task = await h.store.createTask({ description: "rollback target" });
|
||||||
|
const notices: string[] = [];
|
||||||
|
attachMailbox(h.store, notices);
|
||||||
|
(h.store as DeleteTestStore).__afterLifecycleOutboxWriteForTest = () => { throw new Error("inject rollback"); };
|
||||||
|
|
||||||
|
await expect(h.store.deleteTask(task.id)).rejects.toThrow("inject rollback");
|
||||||
|
expect(await lifecycleRows(h, task.id)).toHaveLength(0);
|
||||||
|
expect(await countRunAuditEvents(h.layer.db, { taskId: task.id, mutationType: "task:deleted" })).toBe(0);
|
||||||
|
expect((await h.store.getTask(task.id, { includeDeleted: true })).deletedAt).toBeFalsy();
|
||||||
|
expect(notices).toHaveLength(0);
|
||||||
|
expect(await h.layer.db.select().from(schema.project.taskLifecycleEventSeq)).toHaveLength(0);
|
||||||
|
|
||||||
|
delete (h.store as DeleteTestStore).__afterLifecycleOutboxWriteForTest;
|
||||||
|
await h.store.deleteTask(task.id);
|
||||||
|
const rows = await lifecycleRows(h, task.id);
|
||||||
|
expect(rows).toHaveLength(1);
|
||||||
|
expect(rows[0]?.seq).toBe(1n);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("serializes distinct deletes without collisions while archive and no-op paths write nothing", async () => {
|
||||||
|
h = await createTaskStoreForTest({ prefix: "lifecycle_outbox_distinct", copyFromGolden: true });
|
||||||
|
const tasks = await Promise.all(Array.from({ length: 4 }, (_, index) => h!.store.createTask({ description: `distinct ${index}` })));
|
||||||
|
await Promise.all(tasks.map((task) => h!.store.deleteTask(task.id)));
|
||||||
|
const rows = await h.layer.db.select().from(schema.project.taskLifecycleEvents).orderBy(schema.project.taskLifecycleEvents.seq);
|
||||||
|
expect(rows.map((row) => row.seq)).toEqual([1n, 2n, 3n, 4n]);
|
||||||
|
expect((await Promise.all(tasks.map((task) => h!.store.getTask(task.id, { includeDeleted: true })))).every((task) => Boolean(task.deletedAt))).toBe(true);
|
||||||
|
|
||||||
|
const archived = await h.store.createTask({ description: "archive is not delete" });
|
||||||
|
await h.store.archiveTask(archived.id, { cleanup: false });
|
||||||
|
expect(await lifecycleRows(h, archived.id)).toHaveLength(0);
|
||||||
|
await h.store.deleteTask(tasks[0]!.id);
|
||||||
|
expect(await lifecycleRows(h, tasks[0]!.id)).toHaveLength(1);
|
||||||
|
await expect(h.store.deleteTask("FN-DOES-NOT-EXIST")).rejects.toThrow();
|
||||||
|
expect(await h.layer.db.select().from(schema.project.taskLifecycleEvents)).toHaveLength(4);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("keeps private deterministic fault hooks absent from ordinary stores and production source", async () => {
|
||||||
|
h = await createTaskStoreForTest({ prefix: "lifecycle_outbox_hook", copyFromGolden: true });
|
||||||
|
expect((h.store as DeleteTestStore).__beforeDeleteClaimForTest).toBeUndefined();
|
||||||
|
expect((h.store as DeleteTestStore).__afterLifecycleOutboxWriteForTest).toBeUndefined();
|
||||||
|
const source = await readFile(new URL("../archive-lifecycle-2.ts", import.meta.url), "utf8");
|
||||||
|
expect(source.match(/__beforeDeleteClaimForTest\s*=/g) ?? []).toHaveLength(0);
|
||||||
|
expect(source.match(/__afterLifecycleOutboxWriteForTest\s*=/g) ?? []).toHaveLength(0);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -24,7 +24,8 @@ import {normalizeTaskPriority} from "../task-priority.js";
|
|||||||
import {generateTaskLineageId} from "../task-lineage.js";
|
import {generateTaskLineageId} from "../task-lineage.js";
|
||||||
import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
|
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} from "../task-store/async-persistence.js";
|
import {softDeleteTaskRowInTransaction, readTaskRow as readTaskRowAsync, readTaskRowInTransaction} from "../task-store/async-persistence.js";
|
||||||
|
import {appendTaskLifecycleEventInTransaction} from "../task-store/lifecycle-outbox.js";
|
||||||
import {findLiveLineageChildren as findLiveLineageChildrenAsync, projectPartition, removeLineageReferences} from "../task-store/async-lifecycle.js";
|
import {findLiveLineageChildren as findLiveLineageChildrenAsync, projectPartition, removeLineageReferences} from "../task-store/async-lifecycle.js";
|
||||||
import { resolveProjectColumnsForRoles } from "../project-lane-vocabulary.js";
|
import { resolveProjectColumnsForRoles } from "../project-lane-vocabulary.js";
|
||||||
import {archiveParentTaskWithLineageGate, findArchivedTaskEntry, deleteArchivedTaskEntry, restoreTaskFromArchive} from "../task-store/async-archive-lineage.js";
|
import {archiveParentTaskWithLineageGate, findArchivedTaskEntry, deleteArchivedTaskEntry, restoreTaskFromArchive} from "../task-store/async-archive-lineage.js";
|
||||||
@@ -137,7 +138,16 @@ export async function taskToArchiveEntryImpl(store: TaskStore, task: Task, archi
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function deleteTaskBackendImpl(store: TaskStore, id: string, options?: { removeDependencyReferences?: boolean; removeLineageReferences?: boolean; allowResurrection?: boolean; githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext; auditContext?: TaskDeleteAuditContext; },): Promise<Task> {
|
type DeleteTaskBackendOptions = { removeDependencyReferences?: boolean; removeLineageReferences?: boolean; allowResurrection?: boolean; githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext; auditContext?: TaskDeleteAuditContext; };
|
||||||
|
type DeleteTaskClaimResult = { task: Task; claimed: boolean };
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-11:12:
|
||||||
|
The internal result preserves whether this caller won the conditional first-transition claim.
|
||||||
|
`deleteTaskIf` exposes that fact as `deleted`, so a cross-process loser cannot report that it
|
||||||
|
performed a deletion merely because its predicate ran against a stale live snapshot.
|
||||||
|
*/
|
||||||
|
async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string, options?: DeleteTaskBackendOptions): Promise<DeleteTaskClaimResult> {
|
||||||
/*
|
/*
|
||||||
FNXC:TaskDeletion 2026-07-01-00:00:
|
FNXC:TaskDeletion 2026-07-01-00:00:
|
||||||
Task-bound runtime callers may never soft-delete the task they are executing; this guard is the PostgreSQL-backend mirror of the SQLite-path guard in deleteTaskImpl so direct callers of deleteTaskBackend inherit the same invariant before any mutation or audit.
|
Task-bound runtime callers may never soft-delete the task they are executing; this guard is the PostgreSQL-backend mirror of the SQLite-path guard in deleteTaskImpl so direct callers of deleteTaskBackend inherit the same invariant before any mutation or audit.
|
||||||
@@ -157,7 +167,7 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option
|
|||||||
|
|
||||||
// Idempotent: already soft-deleted is a no-op.
|
// Idempotent: already soft-deleted is a no-op.
|
||||||
if (task.deletedAt) {
|
if (task.deletedAt) {
|
||||||
return task;
|
return { task, claimed: false };
|
||||||
}
|
}
|
||||||
|
|
||||||
// Lineage-integrity gate (VAL-DATA-010).
|
// Lineage-integrity gate (VAL-DATA-010).
|
||||||
@@ -171,9 +181,28 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option
|
|||||||
|
|
||||||
const deletedAt = new Date().toISOString();
|
const deletedAt = new Date().toISOString();
|
||||||
const allowResurrection = options?.allowResurrection === true;
|
const allowResurrection = options?.allowResurrection === true;
|
||||||
|
const projectId = layer.projectId?.trim() || "__legacy_unscoped__";
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
Test-only barrier: production never assigns this private store property. It makes the
|
||||||
|
cross-process pre-claim TOCTOU deterministic instead of relying on scheduler timing.
|
||||||
|
*/
|
||||||
|
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).
|
||||||
await layer.transactionImmediate(async (tx) => {
|
const deletion = await layer.transactionImmediate(async (tx) => {
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
The pre-transaction deletedAt read is a cross-process TOCTOU window. A conditional
|
||||||
|
claim makes one transition own all side effects; a loser re-reads on this transaction
|
||||||
|
because returning its captured live snapshot would lie about deletedAt.
|
||||||
|
*/
|
||||||
|
const claimed = await softDeleteTaskRowInTransaction(tx, id, deletedAt, allowResurrection, projectId, true);
|
||||||
|
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.
|
// Clear lineage references on live children so the parent can be deleted.
|
||||||
if (lineageChildIds.length > 0) {
|
if (lineageChildIds.length > 0) {
|
||||||
await removeLineageReferences(tx, id, lineageChildIds, deletedAt, layer.projectId);
|
await removeLineageReferences(tx, id, lineageChildIds, deletedAt, layer.projectId);
|
||||||
@@ -198,8 +227,6 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option
|
|||||||
await recordGeneratedFixOperatorStop(tx, linkedFeature, "task-delete");
|
await recordGeneratedFixOperatorStop(tx, linkedFeature, "task-delete");
|
||||||
await unlinkMissionFeatureFromTaskId(tx, linkedFeature.id);
|
await unlinkMissionFeatureFromTaskId(tx, linkedFeature.id);
|
||||||
}
|
}
|
||||||
// Soft-delete the task row.
|
|
||||||
await softDeleteTaskRowInTransaction(tx, id, deletedAt, allowResurrection, layer.projectId);
|
|
||||||
// Record the audit event.
|
// Record the audit event.
|
||||||
await store.recordRunAuditEventBackend(tx, {
|
await store.recordRunAuditEventBackend(tx, {
|
||||||
domain: "database",
|
domain: "database",
|
||||||
@@ -222,8 +249,42 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option
|
|||||||
...buildDeleteClosureAuditFields(options?.closureContext),
|
...buildDeleteClosureAuditFields(options?.closureContext),
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
This stays inside transactionImmediate, unlike mailbox delivery: state and durable
|
||||||
|
observation must commit or roll back together, which is the outbox's purpose.
|
||||||
|
*/
|
||||||
|
await appendTaskLifecycleEventInTransaction(tx, {
|
||||||
|
projectId,
|
||||||
|
eventType: "task:deleted",
|
||||||
|
taskId: id,
|
||||||
|
occurredAt: deletedAt,
|
||||||
|
payload: {
|
||||||
|
taskId: id,
|
||||||
|
previousColumn: task.column ?? "unknown",
|
||||||
|
previousStatus: task.status ?? null,
|
||||||
|
deletedAt,
|
||||||
|
allowResurrection,
|
||||||
|
githubIssueAction: options?.githubIssueAction ?? null,
|
||||||
|
deletedBy: options?.auditContext?.agentId ?? null,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-10:51:
|
||||||
|
This private test seam injects a failure after every durable delete write, proving the
|
||||||
|
outbox, counter, audit, and soft-delete share one transaction. Production construction
|
||||||
|
never assigns it; it exists instead of timing-dependent fault injection.
|
||||||
|
*/
|
||||||
|
await (store as unknown as { __afterLifecycleOutboxWriteForTest?: () => void | Promise<void> }).__afterLifecycleOutboxWriteForTest?.();
|
||||||
|
// FNXC:LifecycleOutbox 2026-08-01-10:33: return the persisted transition to both
|
||||||
|
// 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)) };
|
||||||
});
|
});
|
||||||
|
|
||||||
|
if (!deletion.claimed) return deletion;
|
||||||
|
|
||||||
// Emit lifecycle event (best-effort, outside the transaction).
|
// Emit lifecycle event (best-effort, outside the transaction).
|
||||||
store.laneCache.invalidate(task.id);
|
store.laneCache.invalidate(task.id);
|
||||||
store.emit("task:deleted", task, {
|
store.emit("task:deleted", task, {
|
||||||
@@ -243,9 +304,13 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option
|
|||||||
{ id: task.id, title: task.title, previousColumn: task.column, previousStatus: task.status ?? null },
|
{ id: task.id, title: task.title, previousColumn: task.column, previousStatus: task.status ?? null },
|
||||||
options?.auditContext,
|
options?.auditContext,
|
||||||
);
|
);
|
||||||
return task;
|
return deletion;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export async function deleteTaskBackendImpl(store: TaskStore, id: string, options?: DeleteTaskBackendOptions): Promise<Task> {
|
||||||
|
return (await deleteTaskBackendWithClaimResultImpl(store, id, options)).task;
|
||||||
|
}
|
||||||
|
|
||||||
/** PostgreSQL mirror of deleteTaskIfImpl: predicate and deletion share one task lock. */
|
/** PostgreSQL mirror of deleteTaskIfImpl: predicate and deletion share one task lock. */
|
||||||
export async function deleteTaskIfBackendImpl(
|
export async function deleteTaskIfBackendImpl(
|
||||||
store: TaskStore,
|
store: TaskStore,
|
||||||
@@ -271,16 +336,14 @@ export async function deleteTaskIfBackendImpl(
|
|||||||
throw new TaskHasLineageChildrenError(id, lineageChildIds);
|
throw new TaskHasLineageChildrenError(id, lineageChildIds);
|
||||||
}
|
}
|
||||||
if (!await predicate(live)) return { task: live, deleted: false };
|
if (!await predicate(live)) return { task: live, deleted: false };
|
||||||
await deleteTaskBackendImpl(store, id, options);
|
const deletion = await deleteTaskBackendWithClaimResultImpl(store, id, options);
|
||||||
/*
|
/*
|
||||||
FNXC:TaskDeletion 2026-07-29-18:45:
|
FNXC:LifecycleOutbox 2026-08-01-11:12:
|
||||||
FN-8361 exposes `{ task, deleted }` as the authoritative conditional-delete
|
`deleted` means this caller won the first-transition claim, not merely that the predicate
|
||||||
result. Return the persisted archived row, not deleteTaskBackendImpl's
|
observed a live row. The loser carries the transaction-scoped re-read deleted task while
|
||||||
pre-delete audit snapshot, so callers can safely inspect an applied result.
|
returning false, preventing downstream conditional-delete callers from duplicating work.
|
||||||
*/
|
*/
|
||||||
const deletedRow = await readTaskRowAsync(layer, id, { includeDeleted: true });
|
return { task: deletion.task, deleted: deletion.claimed };
|
||||||
if (!deletedRow) throw new Error(`Task ${id} disappeared after soft-delete`);
|
|
||||||
return { task: store.rowToTask(store.pgRowToTaskRow(deletedRow)), deleted: true };
|
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -257,6 +257,9 @@ export async function resolveActiveTaskWedgeEpisodeRow(
|
|||||||
* @param id The task id to soft-delete.
|
* @param id The task id to soft-delete.
|
||||||
* @param deletedAt The deletion timestamp (ISO-8601).
|
* @param deletedAt The deletion timestamp (ISO-8601).
|
||||||
* @param allowResurrection Whether the task may be resurrected (1/0).
|
* @param allowResurrection Whether the task may be resurrected (1/0).
|
||||||
|
* @param firstTransitionOnly Require `deleted_at IS NULL` and report whether this
|
||||||
|
* transaction won the first-transition claim. Archive callers retain their
|
||||||
|
* established unconditional soft-delete behavior.
|
||||||
*/
|
*/
|
||||||
export async function softDeleteTaskRowInTransaction(
|
export async function softDeleteTaskRowInTransaction(
|
||||||
tx: DbTransaction,
|
tx: DbTransaction,
|
||||||
@@ -264,12 +267,24 @@ export async function softDeleteTaskRowInTransaction(
|
|||||||
deletedAt: string,
|
deletedAt: string,
|
||||||
allowResurrection = false,
|
allowResurrection = false,
|
||||||
projectId?: string,
|
projectId?: string,
|
||||||
): Promise<void> {
|
firstTransitionOnly = false,
|
||||||
|
): Promise<boolean> {
|
||||||
/*
|
/*
|
||||||
FNXC:ArchiveProjectIsolation 2026-07-14-16:20:
|
FNXC:ArchiveProjectIsolation 2026-07-14-16:20:
|
||||||
Transactional archive/delete helpers receive the owning project explicitly because task IDs repeat across projects. The composite predicate is required for atomicity to protect the intended row instead of whichever same-ID row PostgreSQL returns first.
|
Transactional archive/delete helpers receive the owning project explicitly because task IDs repeat across projects. The composite predicate is required for atomicity to protect the intended row instead of whichever same-ID row PostgreSQL returns first.
|
||||||
*/
|
*/
|
||||||
await tx
|
const predicates = [
|
||||||
|
eq(schema.project.tasks.projectId, projectId?.trim() || "__legacy_unscoped__"),
|
||||||
|
eq(schema.project.tasks.id, id),
|
||||||
|
];
|
||||||
|
/*
|
||||||
|
FNXC:LifecycleOutbox 2026-08-01-11:01:
|
||||||
|
Delete alone needs a first-transition claim for its durable side effects. Keep
|
||||||
|
archive's existing unconditional mutation separate: applying the claim to it
|
||||||
|
would let an archive snapshot commit after a concurrent delete won the row.
|
||||||
|
*/
|
||||||
|
if (firstTransitionOnly) predicates.push(isNull(schema.project.tasks.deletedAt));
|
||||||
|
const claimed = await tx
|
||||||
.update(schema.project.tasks)
|
.update(schema.project.tasks)
|
||||||
.set({
|
.set({
|
||||||
column: "archived",
|
column: "archived",
|
||||||
@@ -277,10 +292,9 @@ export async function softDeleteTaskRowInTransaction(
|
|||||||
allowResurrection: allowResurrection ? 1 : 0,
|
allowResurrection: allowResurrection ? 1 : 0,
|
||||||
updatedAt: deletedAt,
|
updatedAt: deletedAt,
|
||||||
})
|
})
|
||||||
.where(and(
|
.where(and(...predicates))
|
||||||
eq(schema.project.tasks.projectId, projectId?.trim() || "__legacy_unscoped__"),
|
.returning({ id: schema.project.tasks.id });
|
||||||
eq(schema.project.tasks.id, id),
|
return claimed.length === 1;
|
||||||
));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
46
packages/core/src/task-store/lifecycle-outbox.ts
Normal file
46
packages/core/src/task-store/lifecycle-outbox.ts
Normal file
@@ -0,0 +1,46 @@
|
|||||||
|
import { createHash } from "node:crypto";
|
||||||
|
import { sql } from "drizzle-orm";
|
||||||
|
import type { DbTransaction } from "../postgres/data-layer.js";
|
||||||
|
|
||||||
|
export type TaskDeletedLifecyclePayload = {
|
||||||
|
taskId: string;
|
||||||
|
previousColumn: string;
|
||||||
|
previousStatus: string | null;
|
||||||
|
deletedAt: string;
|
||||||
|
allowResurrection: boolean;
|
||||||
|
githubIssueAction: string | null;
|
||||||
|
deletedBy: string | null;
|
||||||
|
};
|
||||||
|
|
||||||
|
export function makeTaskLifecycleEventId(projectId: string, eventType: string, taskId: string, occurredAt: string): string {
|
||||||
|
return `evt_${createHash("sha256").update(`${projectId}\0${eventType}\0${taskId}\0${occurredAt}`).digest("hex").slice(0, 32)}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* FNXC:LifecycleOutbox 2026-08-01-10:33:
|
||||||
|
* Allocation occurs in the delete transaction. Each project counter row remains locked to
|
||||||
|
* commit, so allocation order is commit order; rollback reverts the counter and consumes no
|
||||||
|
* sequence. This avoids both cross-project contention and MAX(seq)+1 collision aborts.
|
||||||
|
*/
|
||||||
|
export async function appendTaskLifecycleEventInTransaction(
|
||||||
|
tx: DbTransaction,
|
||||||
|
input: { projectId: string; eventType: "task:deleted"; taskId: string; occurredAt: string; payload: TaskDeletedLifecyclePayload },
|
||||||
|
): Promise<{ seq: string; eventId: string }> {
|
||||||
|
const sequenceRows = await tx.execute(sql`
|
||||||
|
INSERT INTO project.task_lifecycle_event_seq (project_id, last_seq)
|
||||||
|
VALUES (${input.projectId}, 1)
|
||||||
|
ON CONFLICT (project_id)
|
||||||
|
DO UPDATE SET last_seq = project.task_lifecycle_event_seq.last_seq + 1
|
||||||
|
RETURNING last_seq
|
||||||
|
`) as unknown as Array<{ last_seq: number | string }>;
|
||||||
|
// FNXC:LifecycleOutbox 2026-08-01-10:33: PostgreSQL bigint values exceed
|
||||||
|
// Number's exact range; preserve the returned decimal sequence for the INSERT.
|
||||||
|
const seq = String(sequenceRows[0]!.last_seq);
|
||||||
|
const eventId = makeTaskLifecycleEventId(input.projectId, input.eventType, input.taskId, input.occurredAt);
|
||||||
|
await tx.execute(sql`
|
||||||
|
INSERT INTO project.task_lifecycle_events
|
||||||
|
(project_id, seq, event_id, event_type, task_id, occurred_at, created_at, payload)
|
||||||
|
VALUES (${input.projectId}, ${seq}, ${eventId}, ${input.eventType}, ${input.taskId}, ${input.occurredAt}, ${input.occurredAt}, ${JSON.stringify(input.payload)}::jsonb)
|
||||||
|
`);
|
||||||
|
return { seq, eventId };
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user