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:
gsxdsm
2026-08-01 04:23:51 -07:00
parent a7591853eb
commit 8c9346ee94
14 changed files with 557 additions and 28 deletions

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

View File

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

View File

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

View File

@@ -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<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> {
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,
]);
});
});

View File

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

View File

@@ -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<unknown>) => fn()),
cleanupBranchForTask: vi.fn(async () => [] as string[]),
clearNearDuplicateReferencesToFailSoft: vi.fn(async () => undefined),

View File

@@ -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<unknown>) => fn()),

View File

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

View File

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

View File

@@ -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(),

View File

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

View File

@@ -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<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:
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<void> }).__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<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).
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<Task> {
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 };
});
}

View File

@@ -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<void> {
firstTransitionOnly = false,
): Promise<boolean> {
/*
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;
}
/**

View 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 };
}