From 8c9346ee946e2324d9748f9aca39d200dc78d5be Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sat, 1 Aug 2026 04:23:51 -0700 Subject: [PATCH] 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) --- .../fn-8684-task-deleted-outbox-writer.md | 7 + docs/architecture.md | 6 + docs/storage.md | 6 + .../__tests__/postgres/schema-applier.test.ts | 82 ++++++- .../task-delete-caller-attribution.test.ts | 12 +- .../task-delete-nonblocking-cleanup.test.ts | 13 +- .../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 ++++++++++++++++++ .../src/task-store/archive-lifecycle-2.ts | 93 ++++++-- .../core/src/task-store/async-persistence.ts | 26 ++- .../core/src/task-store/lifecycle-outbox.ts | 46 ++++ 14 files changed, 557 insertions(+), 28 deletions(-) create mode 100644 .changeset/fn-8684-task-deleted-outbox-writer.md create mode 100644 packages/core/src/postgres/migrations/0040_fn_8684_task_lifecycle_outbox.sql create mode 100644 packages/core/src/task-store/__tests__/lifecycle-outbox-writer.test.ts create mode 100644 packages/core/src/task-store/lifecycle-outbox.ts diff --git a/.changeset/fn-8684-task-deleted-outbox-writer.md b/.changeset/fn-8684-task-deleted-outbox-writer.md new file mode 100644 index 0000000000..bc099d7eab --- /dev/null +++ b/.changeset/fn-8684-task-deleted-outbox-writer.md @@ -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. diff --git a/docs/architecture.md b/docs/architecture.md index fb674d26d2..c1dec9cf1f 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -2329,3 +2329,9 @@ The auto-recovery dispatcher at `packages/engine/src/auto-recovery.ts` (FN-4533) ### 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`. + +## 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. diff --git a/docs/storage.md b/docs/storage.md index 3edb6be649..a071c91833 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -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. 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. diff --git a/packages/core/src/__tests__/postgres/schema-applier.test.ts b/packages/core/src/__tests__/postgres/schema-applier.test.ts index 074014b09a..b6ec68fe8b 100644 --- a/packages/core/src/__tests__/postgres/schema-applier.test.ts +++ b/packages/core/src/__tests__/postgres/schema-applier.test.ts @@ -81,6 +81,8 @@ import { CHAT_SESSION_TAGS_VERSION, DROP_GLOBAL_CONCURRENCY_VERSION, MISSION_TASK_PREFIX_VERSION, + CREDENTIAL_INSTANCE_SELECTION_VERSION, + TASK_LIFECYCLE_OUTBOX_VERSION, } from "../../postgres/schema-applier.js"; import { ProjectPartitionRekeyError, rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js"; import type { PluginSchemaInitHook } from "../../postgres/plugin-schema-hook.js"; @@ -95,6 +97,11 @@ const PG_AVAILABLE = const pgDescribe = PG_AVAILABLE ? describe : describe.skip; 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", () => { expect(MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION).toBe("0003"); 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 its policy and trigger, including actual second-project read/write isolation. */ +async function assertTaskLifecycleOutboxOwnershipContract(ctx: TestContext): Promise { + 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 { const catalog = (await ctx.db.execute(sql` 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; }); - 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(); // FNXC:PostgresCutover 2026-07-05-15:55: apply the BASELINE only. // 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) // + 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 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. - expect(bySchema.project).toBe(98); + expect(bySchema.project).toBe(100); /* 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 @@ -733,6 +765,40 @@ pgDescribe("schema-applier: VAL-SCHEMA-001 final-schema parity (table counts)", 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: 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, DROP_GLOBAL_CONCURRENCY_VERSION, MISSION_TASK_PREFIX_VERSION, + CREDENTIAL_INSTANCE_SELECTION_VERSION, + TASK_LIFECYCLE_OUTBOX_VERSION, ]); expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false); }); @@ -1743,6 +1811,8 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { CHAT_SESSION_TAGS_VERSION, DROP_GLOBAL_CONCURRENCY_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, DROP_GLOBAL_CONCURRENCY_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, DROP_GLOBAL_CONCURRENCY_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, DROP_GLOBAL_CONCURRENCY_VERSION, MISSION_TASK_PREFIX_VERSION, + CREDENTIAL_INSTANCE_SELECTION_VERSION, + TASK_LIFECYCLE_OUTBOX_VERSION, ]); }); }); diff --git a/packages/core/src/__tests__/task-delete-caller-attribution.test.ts b/packages/core/src/__tests__/task-delete-caller-attribution.test.ts index 03aad04070..02519e28b5 100644 --- a/packages/core/src/__tests__/task-delete-caller-attribution.test.ts +++ b/packages/core/src/__tests__/task-delete-caller-attribution.test.ts @@ -39,13 +39,23 @@ let pgRow: TaskRowShape | null = null; vi.mock("../task-store/async-persistence.js", () => ({ 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", () => ({ findLiveLineageChildren: vi.fn(async () => [] as string[]), projectPartition: vi.fn(() => 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", () => ({ getFeatureByTaskId: vi.fn(async () => null), unlinkFeatureFromTaskId: vi.fn(async () => undefined), diff --git a/packages/core/src/__tests__/task-delete-nonblocking-cleanup.test.ts b/packages/core/src/__tests__/task-delete-nonblocking-cleanup.test.ts index 34c095a10c..2bcebd07e5 100644 --- a/packages/core/src/__tests__/task-delete-nonblocking-cleanup.test.ts +++ b/packages/core/src/__tests__/task-delete-nonblocking-cleanup.test.ts @@ -31,13 +31,23 @@ let lineageChildIds: string[] = []; vi.mock("../task-store/async-persistence.js", () => ({ 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", () => ({ findLiveLineageChildren: vi.fn(async () => lineageChildIds), projectPartition: vi.fn(() => 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", () => ({ getFeatureByTaskId: vi.fn(async () => null), unlinkFeatureFromTaskId: vi.fn(async () => undefined), @@ -94,6 +104,7 @@ function makeDeleteStore(task: Task, children: string[] = []) { auditEvents.push(event); }), makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`), + laneCache: { invalidate: vi.fn() }, withTaskLock: vi.fn(async (_id: string, fn: () => Promise) => fn()), cleanupBranchForTask: vi.fn(async () => [] as string[]), clearNearDuplicateReferencesToFailSoft: vi.fn(async () => undefined), diff --git a/packages/core/src/__tests__/task-delete-notice.test.ts b/packages/core/src/__tests__/task-delete-notice.test.ts index b1a27ceded..3c74aec3eb 100644 --- a/packages/core/src/__tests__/task-delete-notice.test.ts +++ b/packages/core/src/__tests__/task-delete-notice.test.ts @@ -35,7 +35,11 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; vi.mock("../task-store/async-persistence.js", () => ({ 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", () => ({ findLiveLineageChildren: vi.fn(async () => [] as string[]), @@ -120,6 +124,7 @@ function makePgStore(task: Task) { recordRunAuditEventBackend: vi.fn(async () => undefined), makeSyntheticDeleteRunId: vi.fn((id: string) => `synthetic-delete-${id}`), emit: vi.fn(), + laneCache: { invalidate: vi.fn() }, /* `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. */ withTaskLock: vi.fn(async (_id: string, fn: () => Promise) => fn()), diff --git a/packages/core/src/postgres/migrations/0040_fn_8684_task_lifecycle_outbox.sql b/packages/core/src/postgres/migrations/0040_fn_8684_task_lifecycle_outbox.sql new file mode 100644 index 0000000000..468c748ebb --- /dev/null +++ b/packages/core/src/postgres/migrations/0040_fn_8684_task_lifecycle_outbox.sql @@ -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(); diff --git a/packages/core/src/postgres/schema-applier.ts b/packages/core/src/postgres/schema-applier.ts index 49c9e47a78..d79cde0737 100644 --- a/packages/core/src/postgres/schema-applier.ts +++ b/packages/core/src/postgres/schema-applier.ts @@ -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 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. */ export const TASK_DECLARED_SYMBOLS_VERSION = "0028"; 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"; /** 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"; +/** 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. */ 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 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 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 @@ -497,6 +501,7 @@ export async function applySchemaBaseline( const dropGlobalConcurrencyAlreadyApplied = applied.includes(DROP_GLOBAL_CONCURRENCY_VERSION); const missionTaskPrefixAlreadyApplied = applied.includes(MISSION_TASK_PREFIX_VERSION); const credentialInstanceSelectionAlreadyApplied = applied.includes(CREDENTIAL_INSTANCE_SELECTION_VERSION); + const taskLifecycleOutboxAlreadyApplied = applied.includes(TASK_LIFECYCLE_OUTBOX_VERSION); assertBinaryNotOlderThanDatabase(applied); let schemaChanged = false; @@ -1054,6 +1059,14 @@ export async function applySchemaBaseline( 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 }; }); } diff --git a/packages/core/src/postgres/schema/project.ts b/packages/core/src/postgres/schema/project.ts index 79a81115fe..df4e60fee5 100644 --- a/packages/core/src/postgres/schema/project.ts +++ b/packages/core/src/postgres/schema/project.ts @@ -564,6 +564,33 @@ export const symbolLocks = projectSchema.table("symbol_locks", { 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 ──────────────────────────────────────── export const workflowSteps = projectSchema.table("workflow_steps", { id: text("id").primaryKey(), diff --git a/packages/core/src/task-store/__tests__/lifecycle-outbox-writer.test.ts b/packages/core/src/task-store/__tests__/lifecycle-outbox-writer.test.ts new file mode 100644 index 0000000000..fe205fe2db --- /dev/null +++ b/packages/core/src/task-store/__tests__/lifecycle-outbox-writer.test.ts @@ -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; + __afterLifecycleOutboxWriteForTest?: () => void | Promise; +}; + +function deferred() { + let resolve!: (value: T | PromiseLike) => void; + const promise = new Promise((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> | 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); + }); +}); diff --git a/packages/core/src/task-store/archive-lifecycle-2.ts b/packages/core/src/task-store/archive-lifecycle-2.ts index d4d30c6961..5b4deee0bd 100644 --- a/packages/core/src/task-store/archive-lifecycle-2.ts +++ b/packages/core/src/task-store/archive-lifecycle-2.ts @@ -24,7 +24,8 @@ import {normalizeTaskPriority} from "../task-priority.js"; import {generateTaskLineageId} from "../task-lineage.js"; import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.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 { resolveProjectColumnsForRoles } from "../project-lane-vocabulary.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 { +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 { /* 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. @@ -157,7 +167,7 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option // Idempotent: already soft-deleted is a no-op. if (task.deletedAt) { - return task; + return { task, claimed: false }; } // 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 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 }).__beforeDeleteClaimForTest?.(id); // 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. if (lineageChildIds.length > 0) { 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 unlinkMissionFeatureFromTaskId(tx, linkedFeature.id); } - // Soft-delete the task row. - await softDeleteTaskRowInTransaction(tx, id, deletedAt, allowResurrection, layer.projectId); // Record the audit event. await store.recordRunAuditEventBackend(tx, { domain: "database", @@ -222,8 +249,42 @@ export async function deleteTaskBackendImpl(store: TaskStore, id: string, option ...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 }).__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). store.laneCache.invalidate(task.id); 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 }, options?.auditContext, ); - return task; + return deletion; } +export async function deleteTaskBackendImpl(store: TaskStore, id: string, options?: DeleteTaskBackendOptions): Promise { + return (await deleteTaskBackendWithClaimResultImpl(store, id, options)).task; +} + /** PostgreSQL mirror of deleteTaskIfImpl: predicate and deletion share one task lock. */ export async function deleteTaskIfBackendImpl( store: TaskStore, @@ -271,16 +336,14 @@ export async function deleteTaskIfBackendImpl( throw new TaskHasLineageChildrenError(id, lineageChildIds); } 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: - FN-8361 exposes `{ task, deleted }` as the authoritative conditional-delete - result. Return the persisted archived row, not deleteTaskBackendImpl's - pre-delete audit snapshot, so callers can safely inspect an applied result. + FNXC:LifecycleOutbox 2026-08-01-11:12: + `deleted` means this caller won the first-transition claim, not merely that the predicate + observed a live row. The loser carries the transaction-scoped re-read deleted task while + returning false, preventing downstream conditional-delete callers from duplicating work. */ - const deletedRow = await readTaskRowAsync(layer, id, { includeDeleted: true }); - if (!deletedRow) throw new Error(`Task ${id} disappeared after soft-delete`); - return { task: store.rowToTask(store.pgRowToTaskRow(deletedRow)), deleted: true }; + return { task: deletion.task, deleted: deletion.claimed }; }); } diff --git a/packages/core/src/task-store/async-persistence.ts b/packages/core/src/task-store/async-persistence.ts index d4f99c7a1a..82af2529ac 100644 --- a/packages/core/src/task-store/async-persistence.ts +++ b/packages/core/src/task-store/async-persistence.ts @@ -257,6 +257,9 @@ export async function resolveActiveTaskWedgeEpisodeRow( * @param id The task id to soft-delete. * @param deletedAt The deletion timestamp (ISO-8601). * @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( tx: DbTransaction, @@ -264,12 +267,24 @@ export async function softDeleteTaskRowInTransaction( deletedAt: string, allowResurrection = false, projectId?: string, -): Promise { + firstTransitionOnly = false, +): Promise { /* 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. */ - 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) .set({ column: "archived", @@ -277,10 +292,9 @@ export async function softDeleteTaskRowInTransaction( allowResurrection: allowResurrection ? 1 : 0, updatedAt: deletedAt, }) - .where(and( - eq(schema.project.tasks.projectId, projectId?.trim() || "__legacy_unscoped__"), - eq(schema.project.tasks.id, id), - )); + .where(and(...predicates)) + .returning({ id: schema.project.tasks.id }); + return claimed.length === 1; } /** diff --git a/packages/core/src/task-store/lifecycle-outbox.ts b/packages/core/src/task-store/lifecycle-outbox.ts new file mode 100644 index 0000000000..68cabf2637 --- /dev/null +++ b/packages/core/src/task-store/lifecycle-outbox.ts @@ -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 }; +}