FN-8685: add durable cross-process task deletion consumers

Deliver durable, replay-safe cross-process task deletion observation.

- Add PostgreSQL lifecycle consumer cursors, leases, acknowledgements, retention, and recovery.
- Start named consumers in dashboard, serve, and engine runtime paths.
- Preserve delete integration metadata while suppressing replayed GitHub and GitLab side effects.
- Cover outbox identity, observed delivery, fencing, and reconciliation behavior.

Files changed:
 ...fn-8685-cross-process-task-deleted-observers.md |   7 +
 .../fn-8685-task-deleted-outbox-consumers.md       |   7 +
 docs/architecture.md                               |   8 +-
 ...tgres-cross-process-task-deleted-observation.md |   8 +-
 docs/storage.md                                    |  10 +-
 packages/cli/src/commands/dashboard.ts             |   9 +-
 packages/cli/src/commands/serve.ts                 |   9 +-
 packages/cli/src/project-context.ts                |   9 +-
 .../task-deleted-outbox-consumer.pg.test.ts        | 157 ++++++++
 ...-deleted-observed-dispatch-side-effects.test.ts |  36 ++
 .../task-lifecycle-consumer-identity.test.ts       |  22 ++
 packages/core/src/index.ts                         |  11 +
 .../0041_fn_8685_task_lifecycle_consumers.sql      |  88 +++++
 packages/core/src/postgres/schema-applier.ts       |  16 +-
 packages/core/src/postgres/schema/project.ts       |  45 +++
 packages/core/src/postgres/startup-factory.ts      |   4 +
 packages/core/src/store.ts                         |  54 ++-
 .../__tests__/lifecycle-outbox-writer.test.ts      |   4 +-
 .../core/src/task-store/archive-lifecycle-2.ts     |   1 +
 packages/core/src/task-store/lifecycle-ops.ts      |  13 +-
 packages/core/src/task-store/lifecycle-outbox.ts   |   2 +
 packages/core/src/task-store/project-store-ops.ts  |   4 +-
 .../src/task-store/task-deleted-outbox-consumer.ts | 333 +++++++++++++++++
 .../task-store/task-lifecycle-consumer-identity.ts |  32 ++
 .../task-store/task-lifecycle-consumer-registry.ts | 396 +++++++++++++++++++++
 .../task-store/task-lifecycle-event-retention.ts   | 104 ++++++
 packages/core/src/task-store/task-mutation-ops.ts  |   1 +
 packages/dashboard/src/github-tracking-state.ts    |  12 +-
 packages/dashboard/src/gitlab-delete-close.ts      |   3 +
 packages/dashboard/src/gitlab-split-close.ts       |   7 +-
 packages/dashboard/src/project-store-resolver.ts   |   9 +-
 packages/engine/src/project-manager.ts             |   4 +-
 packages/engine/src/project-runtime.ts             |   2 +-
 packages/engine/src/runtimes/in-process-runtime.ts |  17 +-
 packages/engine/src/self-healing.ts                |  27 ++
 35 files changed, 1439 insertions(+), 32 deletions(-)

Fusion-Task-Id: FN-8685
Fusion-Task-Lineage: 63eca9ac-d2af-44b0-ba79-388a950148d3
Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-01 06:09:04 -07:00
parent 4b306f10bd
commit 006cc40454
35 changed files with 1439 additions and 32 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": minor
---
summary: Deliver PostgreSQL task deletions to durable cross-process observers.
category: feature
dev: Adds lifecycle consumer identity, cursor, receipt, lease, dead-letter, and retention storage.

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": minor
---
summary: Deliver task deletions safely to configured PostgreSQL runtime consumers.
category: feature
dev: Adds durable per-consumer outbox cursors, receipts, leases, retries, and bounded retention.

View File

@@ -2334,4 +2334,10 @@ A heartbeat `moveTask` failure with the typed `TaskDeletedError` soft-delete mes
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.
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.
FN-8685 adds per-project, per-consumer registration, cursor/lease, receipt, and dead-letter state. An explicit consumer identity is required; stores with no identity are observation-disabled rather than silently sharing a cursor. The supported role identities are `engine`, `dashboard`, `cli`, `child-process-worker:<persisted-node-id>`, and `remote-node:<persisted-node-id>`; instance keys are sourced at wiring sites from named persisted fields, never boot-generated process state, so they remain byte-stable across restart. Observed dispatch evicts the local cache and is explicitly marked `observed`, suppressing writer-owned delete activity/audit effects. Delivery is **at-least-once**: dispatch precedes receipt/cursor acknowledgement, so crash-window duplicates are expected and listeners must be idempotent. The lease has a fencing token, renewal, and token-checked acknowledgement so an expired holder cannot advance a successor's cursor.
Consumers replay within the 30-day retention window. A stale cursor or retained gap uses snapshot-bounded reconciliation: capture the outbox head before reading live tasks, emit observed deletes for missing cached tasks, and CAS-advance only to that captured head so later rows remain available for normal polling. Poison parking atomically inserts the unique `(project_id, consumer_id, event_id)` dead letter, advances the fenced cursor, resets retries, and writes its audit record.
`SelfHealingManager` invokes bounded `pruneTaskLifecycleEvents` no more than once per project every six hours. It uses durable registration liveness rather than cursor existence, retains unacknowledged rows, and uses an age-only 30-day prune when no consumer is live. Run-audit mutation types are `task-deleted-outbox:catch-up`, `task-deleted-outbox:reconciliation-fallback`, `task-deleted-outbox:lease-fenced`, `task-deleted-outbox:dead-letter`, and `task-deleted-outbox:retention-pruned`.

View File

@@ -16,7 +16,7 @@ applies_when: A Fusion deployment has more than one TaskStore process sharing a
Use a **transactional PostgreSQL outbox with per-consumer durable cursors** for cross-process lifecycle observation. The initial scope is `task:deleted`; its event contract must be deliberately extended before other lifecycle events join the stream.
This is a design decision, not an implementation. The unreachable polling replica is removed by FN-8683. Implementation is tracked by the linked follow-up tasks recorded in the FN-8683 task documents.
The unreachable polling replica is removed by FN-8683. FN-8684 lands the transactional writer and FN-8685 lands the consumer-state migration, explicit identity helper, observed dispatcher, fenced polling seam, and bounded self-healing retention invocation. The consumer remains at-least-once: dispatch precedes the durable receipt/cursor acknowledgement, so crash-window duplicates are expected and observed listeners must remain idempotent.
## Problem and current boundary
@@ -123,6 +123,10 @@ This prevents pruning events that an active consumer has not acknowledged and le
5. Add a NOTIFY wake-up only as an optimization after cursor catch-up is proven; it is never the delivery authority.
6. Enable per deployment, observe cursor lag/dead letters, and retain the old full-cache reconciliation fallback during rollout. The removed SQLite polling replica is not a rollback path because it cannot run in PostgreSQL mode.
## Implementation status
FN-8684 landed the transactional writer and FN-8685 landed the consumer-state tables, explicit consumer identity helper, observed dispatch, fenced polling, reconciliation boundary, retry/dead-letter handling, and bounded self-healing retention invocation. Delivery is available only to PostgreSQL `TaskStore` instances constructed with an explicit durable consumer identity; stores without one intentionally remain observation-disabled. Runtime wiring must source any instance key from a named persisted field, never boot-generated process state.
## Consequences
Until the rollout ships, cross-process `task:deleted` observation is a known PostgreSQL gap. In-process deletes and their existing listeners remain supported and unchanged. Any feature that needs cross-process lifecycle correctness must depend on the outbox implementation rather than reintroducing `store.db` polling.
Cross-process `task:deleted` observation is available through the configured outbox consumer. In-process deletes and their existing listeners remain supported and unchanged. Any feature that needs cross-process lifecycle correctness must depend on the outbox implementation rather than reintroducing `store.db` polling.

View File

@@ -44,7 +44,7 @@ See the [2026-07-14 PostgreSQL runtime cutover review](./postgres-migration-revi
- Active task readers (`getTask`, `listTasks`, search, dependency scans, scheduler/watcher reads, mission task aggregations) must filter with `deletedAt IS NULL`.
- Archived-task flows (`archiveTask`, archived cleanup/migration) hard-delete from the active `tasks` table after copying to PostgreSQL cold storage. Legacy `archive.db` files are import-only.
- ID reservation is unchanged: soft-deleted IDs remain reserved. `distributed-task-id` and `task-id-integrity` intentionally scan all task rows (including soft-deleted rows), and must not filter on `deletedAt`.
- The legacy SQLite polling replica for cross-process task lifecycle observation no longer exists. In PostgreSQL mode `TaskStore.watch()` only warms its local cache; `task:deleted` is process-local until the transactional outbox defined in [PostgreSQL cross-process `task:deleted` observation](./solutions/architecture/postgres-cross-process-task-deleted-observation.md) is implemented. Observing processes must not recreate writer-owned delete run-audit or mailbox effects.
- The legacy SQLite polling replica for cross-process task lifecycle observation no longer exists. PostgreSQL stores configured with an explicit durable consumer identity use the transactional outbox described in [PostgreSQL cross-process `task:deleted` observation](./solutions/architecture/postgres-cross-process-task-deleted-observation.md); stores without an identity remain observation-disabled. Observers must not recreate writer-owned delete run-audit or mailbox effects.
### Orphaned task-dir reconciliation (FN-6783)
@@ -138,7 +138,9 @@ The PostgreSQL writer locks the active `(project_id, task_id)` parent row before
- Task-linked artifact registration requires an active, non-archived task. Archived tasks are read-only for artifact writes; soft-deleted or missing tasks are rejected.
- Retention follows the existing task lifecycle rather than a separate artifact policy: soft-deleted parent tasks keep artifact rows/files for forensics but normal live-reader APIs hide them; hard deletion from the active `tasks` table cascades artifact metadata through the `taskId` foreign key, and archive cleanup removes the task directory that contains task-scoped artifact binaries. Task-less artifacts live under `<rootDir>/.fusion/artifacts/` and are not tied to task archival cleanup.
- Worktree DB hydration copies task-scoped artifact metadata for the current task/dependency graph alongside task rows and `task_documents`. It intentionally does not copy binary payload files, and it intentionally excludes task-less registry artifacts because dependency hydration is scoped to the active task graph.
- **Cross-instance live refresh boundary (FN-8683).** A project can have more than one `TaskStore` instance open against the same PostgreSQL DB, but TaskStore no longer has a SQLite polling replica to re-emit artifact or task lifecycle events in another process. A receiver refreshes through its normal API reads today; the planned durable outbox in [PostgreSQL cross-process `task:deleted` observation](./solutions/architecture/postgres-cross-process-task-deleted-observation.md) is the only approved path for future cross-process lifecycle delivery, not a revived `checkForChanges()` loop.
- **Cross-instance live refresh boundary (FN-8683/FN-8685).** A project can have more than one `TaskStore` instance open against the same PostgreSQL DB. A configured task-deletion consumer replays the durable outbox per identity; other lifecycle events and stores without an explicit identity remain process-local. The outbox is the only approved cross-process delivery path, not a revived `checkForChanges()` loop.
- **Task-deletion consumer storage.** `project.task_lifecycle_consumer_registrations` is the authoritative liveness set (`registered_at`, `last_seen_at`, `active`). Each `(project_id, consumer_id)` has a cursor/lease row containing its acknowledged sequence, retry backoff, fencing token, and expiry; durable receipts suppress already-committed redelivery. Dead letters are unique on `(project_id, consumer_id, event_id)`, so their insert, fenced cursor advance, retry reset, and `task-deleted-outbox:dead-letter` audit record commit as one transaction.
- **Task-deletion outbox retention.** `pruneTaskLifecycleEvents` is the only pruning seam, invoked by the engine self-healing maintenance sweep at most once per project every six hours with a 5,000-row budget. It deletes only 30-day-old rows acknowledged by every live registered consumer; an active registration without a cursor prevents pruning. Stale/inactive registrations do not pin retention. With zero live identities it age-prunes only rows older than 30 days, preserving within-bound restart catch-up; per-project prune failures are non-fatal and retry on the next sweep.
Agent-facing registration tools are documented in [Artifact registry tools](./agents.md#artifact-registry-tools), and the dashboard browsing surface is documented in [Artifacts View](./dashboard-guide.md#artifacts-view).
@@ -809,4 +811,6 @@ Direct chat tags are stored in the project PostgreSQL schema as `chat_tags` and
`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.
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.
FN-8685 adds `task_lifecycle_consumer_registrations`, `task_lifecycle_consumer_cursors`, `task_lifecycle_consumer_receipts`, and `task_lifecycle_consumer_dead_letters`. Every table is scoped by project and consumer; receipts and dead letters use `(project_id, consumer_id, event_id)` uniqueness. Registrations provide durable liveness for retention. Retention runs from self-healing at most every six hours per project, is bounded to 5,000 rows per sweep, and does not prune rows unacknowledged by live consumers. If no identity is live, it only age-prunes events older than 30 days so a within-bound restart can catch up.

View File

@@ -30,6 +30,7 @@ import {
type WorkflowIrColumn,
type TraitFlags,
createTaskStoreForBackend,
buildConsumerId,
FUSION_RESTART_EXIT_CODE,
FUSION_NON_RETRYABLE_EXIT_CODE,
isPostgresUniqueError,
@@ -951,6 +952,8 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
"backend.factory",
() => createTaskStoreForBackend({
rootDir: cwd,
/* FNXC:CrossProcessDeleteObservation 2026-08-01-12:14: The dashboard store subscribes SSE, so it owns the restart-stable dashboard lifecycle consumer stream. */
consumerId: buildConsumerId("dashboard"),
onMigrationProgress: (event) => migrationHoldingServer?.setMigrationProgress(event),
}),
logPhase,
@@ -1065,7 +1068,11 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
if (!store) throw new Error("cwd TaskStore not yet initialized");
projectStore = store;
} else {
const boot = await createTaskStoreForBackend({ rootDir: projectPath });
const boot = await createTaskStoreForBackend({
rootDir: projectPath,
/* FNXC:CrossProcessDeleteObservation 2026-08-01-12:14: Each dashboard project store uses the dashboard role, isolated by project id in durable consumer state. */
consumerId: buildConsumerId("dashboard"),
});
projectStore = boot.taskStore;
projectStoreShutdowns.set(projectPath, boot.shutdown);
setHostTaskStore(projectPath, projectStore);

View File

@@ -303,7 +303,7 @@ export async function runServe(
* the TaskStore) as the externalTaskStore for the cwd project's engine so
* the connection pool is shared — no second embedded PG instance is started.
*/
const { createTaskStoreForBackend } = await import("@fusion/core");
const { createTaskStoreForBackend, buildConsumerId } = await import("@fusion/core");
/*
* FNXC:PostgresFinalCutover 2026-07-14-17:20:
* Serve must share one successfully booted PostgreSQL layer between CentralCore and the cwd engine. A backend boot error is fatal; constructing a layerless CentralCore would make project discovery appear empty and split control-plane state.
@@ -317,6 +317,13 @@ export async function runServe(
"backend.factory",
() => createTaskStoreForBackend({
rootDir: cwd,
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-13:03:
serve shares this store with the engine, whose runtime starts the durable consumer. Give the
shared store the engine identity so its cursor survives restart and cross-process deletes reach
runtime observers even though this headless path never calls dashboard watch().
*/
consumerId: buildConsumerId("engine"),
onMigrationProgress: (event) => migrationHoldingServer?.setMigrationProgress(event),
}),
logPhase,

View File

@@ -5,7 +5,7 @@
* for operating on tasks across multiple registered projects.
*/
import { createTaskStoreForBackend, type AsyncDataLayer, type RegisteredProject, type TaskStore, CentralCore, GlobalSettingsStore, hasProjectIdentity, isValidSqliteDatabaseFile } from "@fusion/core";
import { buildConsumerId, createTaskStoreForBackend, type AsyncDataLayer, type RegisteredProject, type TaskStore, CentralCore, GlobalSettingsStore, hasProjectIdentity, isValidSqliteDatabaseFile } from "@fusion/core";
import { resolve, dirname, basename } from "node:path";
/** Project context for CLI operations */
@@ -351,7 +351,12 @@ export async function createLocalStore(
// catch-fallbacks (task/pr/backup/memory-backup/branch-group/mcp) boot their
// cwd-rooted store through the same factory instead of constructing a legacy
// SQLite TaskStore directly (its runtime throws in backend mode).
const boot = await createTaskStoreForBackend({ rootDir: projectPath, globalSettingsDir });
const boot = await createTaskStoreForBackend({
rootDir: projectPath,
globalSettingsDir,
/* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: CLI owns its factory-created store, so role-only identity is restart-stable. */
consumerId: buildConsumerId("cli"),
});
/* FNXC:PostgresCliLifecycle 2026-07-14-18:07: CLI project contexts return only TaskStore, so retain the factory owner handle in a WeakMap and release it whenever closeProjectStore closes that context. */
storeOwners.set(boot.taskStore, { backendShutdown: boot.shutdown });
return boot.taskStore;

View File

@@ -0,0 +1,157 @@
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:29:
A reconnect must make its 30-day decision from durable acknowledgement state, not a lease update.
These PostgreSQL checks exercise the real conditional updates so a hand-written query fake cannot
hide a lease write that makes an offline cursor appear current.
*/
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "vitest";
import { and, eq } from "drizzle-orm";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
type SharedPgTaskStoreHarness,
} from "../../__test-utils__/pg-test-harness.js";
import * as schema from "../../postgres/schema/index.js";
import { TaskDeletedOutboxConsumer } from "../../task-store/task-deleted-outbox-consumer.js";
import {
acquireTaskLifecycleLease,
ensureTaskLifecycleConsumerCursor,
advanceTaskLifecycleConsumerCursor,
parkTaskLifecycleConsumerDeadLetter,
readTaskLifecycleConsumerCursor,
readTaskLifecycleEventBounds,
renewTaskLifecycleLease,
} from "../../task-store/task-lifecycle-consumer-registry.js";
const pgTest = pgDescribe;
const TEST_PROJECT_ID = "outbox-consumer-project";
pgTest("task:deleted outbox consumer fences", () => {
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_task_deleted_outbox_consumer",
});
beforeAll(h.beforeAll);
beforeEach(h.beforeEach);
afterEach(h.afterEach);
afterAll(h.afterAll);
it("preserves acknowledgement freshness while acquiring and renewing a lease", async () => {
const layer = { ...h.layer(), projectId: TEST_PROJECT_ID };
const consumerId = "engine";
const acknowledgedAt = "2026-06-01T00:00:00.000Z";
await ensureTaskLifecycleConsumerCursor(layer, consumerId, acknowledgedAt);
const lease = await acquireTaskLifecycleLease(
layer,
consumerId,
"lease-a",
"2026-08-02T00:00:15.000Z",
"2026-08-02T00:00:00.000Z",
);
expect(lease).not.toBeNull();
await expect(renewTaskLifecycleLease(
layer,
consumerId,
lease!,
"2026-08-02T00:00:30.000Z",
"2026-08-02T00:00:15.000Z",
)).resolves.toBe(true);
expect((await readTaskLifecycleConsumerCursor(layer, consumerId))?.updatedAt).toBe(acknowledgedAt);
});
it("reconciles an unacknowledged cursor when retained rows start after sequence zero", async () => {
const layer = { ...h.layer(), projectId: TEST_PROJECT_ID };
await layer.db.insert(schema.project.taskLifecycleEvents).values({
projectId: TEST_PROJECT_ID,
seq: 2n,
eventId: "retained-gap",
eventType: "task:deleted",
taskId: "FN-GAP",
occurredAt: "2026-08-01T00:00:00.000Z",
createdAt: "2026-08-01T00:00:00.000Z",
payload: {},
});
const consumer = new TaskDeletedOutboxConsumer({ asyncLayer: layer } as never);
const needsReconciliation = (consumer as unknown as {
needsReconciliation(lastAckedSeq: bigint, updatedAt: string): Promise<"pruned-gap" | "cursor-older-than-retention-bound" | null>;
}).needsReconciliation.bind(consumer);
await expect(needsReconciliation(0n, "2026-08-01T00:00:00.000Z")).resolves.toBe("pruned-gap");
});
it("uses the durable sequence counter after retention removes all event rows and rejects cursor regression", async () => {
const projectId = "outbox-consumer-pruned-sequence";
const layer = { ...h.layer(), projectId };
const consumerId = "engine-pruned-head";
const now = "2026-08-01T00:00:00.000Z";
await ensureTaskLifecycleConsumerCursor(layer, consumerId, now);
await layer.db.update(schema.project.taskLifecycleConsumerCursors).set({ lastAckedSeq: 100n }).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
));
await layer.db.insert(schema.project.taskLifecycleEventSeq).values({
projectId,
lastSeq: 100n,
});
await expect(readTaskLifecycleEventBounds(layer)).resolves.toMatchObject({
oldestSeq: null,
headSeq: 100n,
});
await expect(advanceTaskLifecycleConsumerCursor(layer, consumerId, 100n, 0n, 0n, now)).resolves.toBe(false);
expect((await readTaskLifecycleConsumerCursor(layer, consumerId))?.lastAckedSeq).toBe(100n);
});
it("reports acknowledgement-age reconciliation separately from a retained-row gap", async () => {
const layer = { ...h.layer(), projectId: "outbox-consumer-acknowledgement-age" };
const consumer = new TaskDeletedOutboxConsumer({ asyncLayer: layer } as never);
const needsReconciliation = (consumer as unknown as {
needsReconciliation(lastAckedSeq: bigint, updatedAt: string): Promise<"pruned-gap" | "cursor-older-than-retention-bound" | null>;
}).needsReconciliation.bind(consumer);
await expect(needsReconciliation(0n, "2026-06-01T00:00:00.000Z")).resolves.toBe("cursor-older-than-retention-bound");
});
it("rolls back a fenced dead-letter park so the poller can emit lease-fenced evidence", async () => {
const layer = { ...h.layer(), projectId: TEST_PROJECT_ID };
const consumerId = "engine-fence";
await ensureTaskLifecycleConsumerCursor(layer, consumerId, "2026-08-01T00:00:00.000Z");
const staleLease = await acquireTaskLifecycleLease(
layer,
consumerId,
"lease-stale",
"2026-08-01T00:00:01.000Z",
"2026-08-01T00:00:00.000Z",
);
const successorLease = await acquireTaskLifecycleLease(
layer,
consumerId,
"lease-successor",
"2026-08-01T00:01:00.000Z",
"2026-08-01T00:00:02.000Z",
);
expect(staleLease).not.toBeNull();
expect(successorLease?.fencingToken).toBeGreaterThan(staleLease!.fencingToken);
await expect(parkTaskLifecycleConsumerDeadLetter(layer, {
consumerId,
eventId: "event-stale",
seq: 1n,
priorSeq: 0n,
attempts: 10,
failureClass: "TypeError",
lease: staleLease!,
now: "2026-08-01T00:00:03.000Z",
})).resolves.toBe(false);
const cursor = await readTaskLifecycleConsumerCursor(layer, consumerId);
expect(cursor?.lastAckedSeq).toBe(0n);
const letters = await layer.db.select().from(schema.project.taskLifecycleConsumerDeadLetters).where(and(
eq(schema.project.taskLifecycleConsumerDeadLetters.projectId, layer.projectId!),
eq(schema.project.taskLifecycleConsumerDeadLetters.consumerId, consumerId),
));
expect(letters).toHaveLength(0);
});
});

View File

@@ -0,0 +1,36 @@
import { describe, expect, it } from "vitest";
import { TaskStore } from "../store.js";
import type { Task } from "../types.js";
const task = { id: "FN-observed", title: "Observed", column: "done" } as Task;
describe("observed task:deleted dispatch", () => {
it("evicts local cache, preserves delete intent, and marks notifications observed", () => {
const store = new TaskStore(process.cwd());
store.taskCache.set(task.id, task);
const received: Array<{
observed?: boolean;
outboxEventId?: string;
githubIssueAction?: "leave";
closureContext?: { kind: "split-into-subtasks"; childTaskIds: string[] };
} | undefined> = [];
store.on("task:deleted", () => { throw new Error("listener failure must not stop fan-out"); });
store.on("task:deleted", (_task, meta) => received.push(meta));
expect(store.emitObservedTaskDeleted(task, "evt_observed", {
githubIssueAction: "leave",
closureContext: { kind: "split-into-subtasks", childTaskIds: ["FN-child"] },
})).toBe(true);
expect(store.taskCache.has(task.id)).toBe(false);
expect(received).toEqual([{
observed: true,
outboxEventId: "evt_observed",
githubIssueAction: "leave",
closureContext: { kind: "split-into-subtasks", childTaskIds: ["FN-child"] },
}]);
// The specified crash-window duplicate remains a harmless cache eviction + bridge notification.
expect(store.emitObservedTaskDeleted(task, "evt_observed")).toBe(true);
expect(received).toHaveLength(2);
});
});

View File

@@ -0,0 +1,22 @@
import { describe, expect, it } from "vitest";
import { buildConsumerId } from "../task-store/task-lifecycle-consumer-identity.js";
describe("buildConsumerId", () => {
it("builds role-only and persisted-instance identities", () => {
expect(buildConsumerId("engine")).toBe("engine");
expect(buildConsumerId("child-process-worker", "node-17")).toBe("child-process-worker:node-17");
});
it("accepts only roles and safe trimmed instance-key format", () => {
expect(() => buildConsumerId("unknown" as "engine")).toThrow("Unknown");
expect(() => buildConsumerId("remote-node", " node")).toThrow("trimmed");
expect(() => buildConsumerId("remote-node", "node:1")).toThrow("safe");
expect(() => buildConsumerId("remote-node", "")).toThrow("safe");
});
it("does not pretend to infer durability from a string", () => {
// Durability is proven at the factory/runtime wiring site, not by pattern guessing here.
expect(buildConsumerId("remote-node", "550e8400-e29b-41d4-a716-446655440000"))
.toBe("remote-node:550e8400-e29b-41d4-a716-446655440000");
});
});

View File

@@ -833,6 +833,17 @@ export {
TransitionRejectionError,
type LegacyAutoMergeStampReconcileResult,
} from "./store.js";
export {
pruneTaskLifecycleEvents,
TASK_LIFECYCLE_RETENTION_DAYS,
TASK_LIFECYCLE_RETENTION_MAX_DELETES,
type TaskLifecycleRetentionResult,
} from "./task-store/task-lifecycle-event-retention.js";
export {
buildConsumerId,
TASK_LIFECYCLE_CONSUMER_ROLES,
type TaskLifecycleConsumerRole,
} from "./task-store/task-lifecycle-consumer-identity.js";
export {
STOPWORDS,
tokenize,

View File

@@ -0,0 +1,88 @@
-- FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
-- FN-8685 adds the durable read half of the PostgreSQL task:deleted outbox. These rows
-- are partitioned by project and consumer identity so each independently observing runtime
-- has its own ordered at-least-once stream, fenced lease, receipts, and poison evidence.
CREATE TABLE IF NOT EXISTS project.task_lifecycle_consumer_registrations (
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
consumer_id text NOT NULL,
registered_at text NOT NULL,
last_seen_at text NOT NULL,
active integer NOT NULL DEFAULT 1,
PRIMARY KEY (project_id, consumer_id)
);
CREATE TABLE IF NOT EXISTS project.task_lifecycle_consumer_cursors (
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
consumer_id text NOT NULL,
last_acked_seq bigint NOT NULL DEFAULT 0,
retry_attempts integer NOT NULL DEFAULT 0,
retry_backoff_until text,
lease_token text,
fencing_token bigint NOT NULL DEFAULT 0,
lease_expires_at text,
updated_at text NOT NULL,
PRIMARY KEY (project_id, consumer_id)
);
CREATE TABLE IF NOT EXISTS project.task_lifecycle_consumer_receipts (
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
consumer_id text NOT NULL,
event_id text NOT NULL,
seq bigint NOT NULL,
processed_at text NOT NULL,
PRIMARY KEY (project_id, consumer_id, event_id)
);
CREATE INDEX IF NOT EXISTS "idxTaskLifecycleConsumerReceiptsSequence" ON project.task_lifecycle_consumer_receipts(project_id, consumer_id, seq);
CREATE TABLE IF NOT EXISTS project.task_lifecycle_consumer_dead_letters (
project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true),
consumer_id text NOT NULL,
event_id text NOT NULL,
seq bigint NOT NULL,
attempts integer NOT NULL,
failure_class text NOT NULL,
parked_at text NOT NULL,
updated_at text NOT NULL,
PRIMARY KEY (project_id, consumer_id, event_id)
);
-- All new state has the same forced RLS and project-id trigger as the outbox writer.
ALTER TABLE project.task_lifecycle_consumer_registrations ENABLE ROW LEVEL SECURITY;
ALTER TABLE project.task_lifecycle_consumer_registrations FORCE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_consumer_registrations;
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_consumer_registrations
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_consumer_registrations;
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_consumer_registrations
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();
ALTER TABLE project.task_lifecycle_consumer_cursors ENABLE ROW LEVEL SECURITY;
ALTER TABLE project.task_lifecycle_consumer_cursors FORCE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_consumer_cursors;
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_consumer_cursors
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_consumer_cursors;
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_consumer_cursors
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();
ALTER TABLE project.task_lifecycle_consumer_receipts ENABLE ROW LEVEL SECURITY;
ALTER TABLE project.task_lifecycle_consumer_receipts FORCE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_consumer_receipts;
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_consumer_receipts
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_consumer_receipts;
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_consumer_receipts
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();
ALTER TABLE project.task_lifecycle_consumer_dead_letters ENABLE ROW LEVEL SECURITY;
ALTER TABLE project.task_lifecycle_consumer_dead_letters FORCE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS fusion_project_isolation ON project.task_lifecycle_consumer_dead_letters;
CREATE POLICY fusion_project_isolation ON project.task_lifecycle_consumer_dead_letters
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_consumer_dead_letters;
CREATE TRIGGER fusion_assign_project_id BEFORE INSERT OR UPDATE OF project_id ON project.task_lifecycle_consumer_dead_letters
FOR EACH ROW EXECUTE FUNCTION project.fusion_assign_project_id();

View File

@@ -54,8 +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.
*/
/* 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:CrossProcessDeleteObservation 2026-08-01-11:39: advance the schema ceiling so durable consumer state exists before observers begin polling FN-8684's outbox. */
export const SCHEMA_BASELINE_VERSION = "0041";
/** 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";
@@ -178,6 +178,8 @@ export const MISSION_TASK_PREFIX_VERSION = "0038";
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";
/** FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: 0041 creates project-scoped consumer registration, cursor, receipt, and dead-letter state. */
export const TASK_LIFECYCLE_CONSUMERS_VERSION = "0041";
/** SECURITY DEFINER helper that only inserts LEGACY_ADOPTION_DRAINED_MARKER. */
export const LEGACY_ADOPTION_DRAINED_MARKER_FUNCTION = "fusion_mark_legacy_adoption_drained";
@@ -392,6 +394,7 @@ const DROP_GLOBAL_CONCURRENCY_MIGRATION_PATH = join(MIGRATIONS_DIR, "0037_drop_g
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");
const TASK_LIFECYCLE_CONSUMERS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0041_fn_8685_task_lifecycle_consumers.sql");
/**
* Ensure the migration bookkeeping table exists. Lives in the public schema so
@@ -502,6 +505,7 @@ export async function applySchemaBaseline(
const missionTaskPrefixAlreadyApplied = applied.includes(MISSION_TASK_PREFIX_VERSION);
const credentialInstanceSelectionAlreadyApplied = applied.includes(CREDENTIAL_INSTANCE_SELECTION_VERSION);
const taskLifecycleOutboxAlreadyApplied = applied.includes(TASK_LIFECYCLE_OUTBOX_VERSION);
const taskLifecycleConsumersAlreadyApplied = applied.includes(TASK_LIFECYCLE_CONSUMERS_VERSION);
assertBinaryNotOlderThanDatabase(applied);
let schemaChanged = false;
@@ -1067,6 +1071,14 @@ export async function applySchemaBaseline(
schemaChanged = true;
}
/* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: consumer migration is separate from the immutable FN-8684 writer migration so upgrades cannot mistake one for the other. */
if (!taskLifecycleConsumersAlreadyApplied) {
const migrationSql = await readFile(TASK_LIFECYCLE_CONSUMERS_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_CONSUMERS_VERSION}) ON CONFLICT (version) DO NOTHING`);
schemaChanged = true;
}
return { applied: schemaChanged, pluginHooksRun: pluginHooks.length };
});
}

View File

@@ -591,6 +591,51 @@ export const taskLifecycleEventSeq = projectSchema.table("task_lifecycle_event_s
lastSeq: bigint("last_seq", { mode: "bigint" }).notNull().default(sql`0`),
}, (t) => [primaryKey({ columns: [t.projectId] })]);
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
The transactional delete outbox needs durable state per independently observing identity.
Registration, cursor/lease, receipt, and dead-letter rows stay project-scoped so one
runtime's acknowledgement never suppresses another runtime's at-least-once delivery.
*/
export const taskLifecycleConsumerRegistrations = projectSchema.table("task_lifecycle_consumer_registrations", {
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
consumerId: text("consumer_id").notNull(),
registeredAt: text("registered_at").notNull(),
lastSeenAt: text("last_seen_at").notNull(),
active: integer("active").notNull().default(1),
}, (t) => [primaryKey({ columns: [t.projectId, t.consumerId] })]);
export const taskLifecycleConsumerCursors = projectSchema.table("task_lifecycle_consumer_cursors", {
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
consumerId: text("consumer_id").notNull(),
lastAckedSeq: bigint("last_acked_seq", { mode: "bigint" }).notNull().default(sql`0`),
retryAttempts: integer("retry_attempts").notNull().default(0),
retryBackoffUntil: text("retry_backoff_until"),
leaseToken: text("lease_token"),
fencingToken: bigint("fencing_token", { mode: "bigint" }).notNull().default(sql`0`),
leaseExpiresAt: text("lease_expires_at"),
updatedAt: text("updated_at").notNull(),
}, (t) => [primaryKey({ columns: [t.projectId, t.consumerId] })]);
export const taskLifecycleConsumerReceipts = projectSchema.table("task_lifecycle_consumer_receipts", {
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
consumerId: text("consumer_id").notNull(),
eventId: text("event_id").notNull(),
seq: bigint("seq", { mode: "bigint" }).notNull(),
processedAt: text("processed_at").notNull(),
}, (t) => [primaryKey({ columns: [t.projectId, t.consumerId, t.eventId] })]);
export const taskLifecycleConsumerDeadLetters = projectSchema.table("task_lifecycle_consumer_dead_letters", {
projectId: text("project_id").notNull().default(sql`current_setting('fusion.project_id', true)`),
consumerId: text("consumer_id").notNull(),
eventId: text("event_id").notNull(),
seq: bigint("seq", { mode: "bigint" }).notNull(),
attempts: integer("attempts").notNull(),
failureClass: text("failure_class").notNull(),
parkedAt: text("parked_at").notNull(),
updatedAt: text("updated_at").notNull(),
}, (t) => [primaryKey({ columns: [t.projectId, t.consumerId, t.eventId] })]);
// ── Workflow step definitions ────────────────────────────────────────
export const workflowSteps = projectSchema.table("workflow_steps", {
id: text("id").primaryKey(),

View File

@@ -765,6 +765,8 @@ export interface CreateTaskStoreForBackendOptions {
* constructor (matching `new TaskStore(rootDir)`).
*/
readonly projectId?: string;
/** Explicit durable lifecycle observer identity; absent deliberately disables observation. */
readonly consumerId?: string;
/*
FNXC:MigrationHoldingPage 2026-07-17-12:20:
During the one-time SQLite→PostgreSQL auto-migration the caller's HTTP server is
@@ -1259,10 +1261,12 @@ export async function createTaskStoreForBackend(
undefined,
options.globalSettingsDir,
asyncLayer,
options.consumerId,
);
} else {
taskStore = new TaskStore(rootDir, options.globalSettingsDir, {
asyncLayer,
...(options.consumerId ? { consumerId: options.consumerId } : {}),
});
await taskStore.init();
}

View File

@@ -119,6 +119,7 @@ import { flushAgentLogBufferImpl, appendAgentLogBatchImpl } from "./task-store/a
import { refineTaskImpl, updateTaskDependenciesImpl } from "./task-store/update-task-deps.js";
import { createWorkflowStepImpl, updateWorkflowStepImpl, updateWorkflowDefinitionImpl, deleteWorkflowDefinitionImpl, setDefaultWorkflowIdImpl, selectTaskWorkflowImpl } from "./task-store/workflow-ops.js";
import { initImpl, setupActivityLogListenersImpl, reconcileOrphanedTaskDirsImpl, watchImpl, migrateAgentLogEntriesImpl, migrateMovedSettingsImpl, recoverStaleTransitionPendingImpl, migrateLegacyWorkflowStepsImpl, emitTaskLifecycleEventSafelyImpl } from "./task-store/lifecycle-ops.js";
import { TaskDeletedOutboxConsumer } from "./task-store/task-deleted-outbox-consumer.js";
import { updateStepImpl, startStepImpl, acquireMergeQueueLeaseImpl, mergeTaskImpl } from "./task-store/merge-queue-ops.js";
import { addCommentImpl, publishArchivedTaskDocumentAdditionImpl, upsertTaskDocumentImpl } from "./task-store/comments-ops.js";
import { deleteTaskImpl, archiveTaskImpl, type DeleteTaskIfResult } from "./task-store/archive-lifecycle.js";
@@ -169,7 +170,13 @@ export interface TaskStoreEvents {
unchanged; absent metadata is unknown, never a legacy-lane claim.
*/
"task:updated": [task: Task, meta?: { lanes?: TaskMoveLanes }];
"task:deleted": [task: Task, meta?: { githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext }];
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
Observed outbox delivery is at-least-once, including a crash-window duplicate. The explicit
marker lets listener paths suppress writer-owned accumulating effects while bridge/cache work
remains idempotent; event identity makes duplicate provenance inspectable without payload prose.
*/
"task:deleted": [task: Task, meta?: { githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext; observed?: boolean; outboxEventId?: string }];
"task:merged": [result: MergeResult];
"settings:updated": [data: { settings: Settings; previous: Settings }];
"workflow:setting-values-updated": [data: {
@@ -373,8 +380,8 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
public static readonly DEFAULT_WORKFLOW_POOL_ID = DEFAULT_WORKFLOW_POOL_ID;
/** FNXC:RuntimeBackendInjection 2026-06-24-14:20: Backend-mode factory. */
static async getOrCreateForProject( projectId?: string, centralCore?: CentralCore, globalSettingsDir?: string, asyncLayer?: AsyncDataLayer, ): Promise<TaskStore> {
return getOrCreateForProjectImpl(this, projectId, centralCore, globalSettingsDir, asyncLayer);
static async getOrCreateForProject( projectId?: string, centralCore?: CentralCore, globalSettingsDir?: string, asyncLayer?: AsyncDataLayer, consumerId?: string, ): Promise<TaskStore> {
return getOrCreateForProjectImpl(this, projectId, centralCore, globalSettingsDir, asyncLayer, consumerId);
}
/** FNXC:PostgresRuntimeStorage 2026-07-14-18:47: Task metadata is authoritative in PostgreSQL; task document/blob files remain on disk. */
@@ -389,6 +396,9 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
* FNXC:PostgresRuntimeStorage 2026-07-14-18:47: Production TaskStores receive an AsyncDataLayer and delegate all persistence to PostgreSQL. A missing layer is a construction error; retained sync members exist only until compatibility tests and types are removed.
*/
public readonly asyncLayer: AsyncDataLayer | null = null;
/** Explicitly absent means cross-process lifecycle observation is disabled. */
public readonly consumerId: string | null;
private taskDeletedOutboxConsumer: TaskDeletedOutboxConsumer | null = null;
private pluginPostgresSchemaExecutor: ((contracts: readonly LoadedPluginSchemaContract[]) => Promise<void>) | null = null;
/*
@@ -517,7 +527,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
}
/** FNXC:RuntimeBackendInjection 2026-06-24-14:05: asyncLayer → backend mode (PostgreSQL, no SQLite); absent → legacy SQLite. */
constructor( public rootDir: string, globalSettingsDir?: string, options?: { asyncLayer?: AsyncDataLayer }, ) {
constructor( public rootDir: string, globalSettingsDir?: string, options?: { asyncLayer?: AsyncDataLayer; consumerId?: string }, ) {
super();
this.setMaxListeners(100);
assertProjectRootDir(rootDir, "TaskStore");
@@ -526,6 +536,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
this.tasksDir = join(this.fusionDir, "tasks");
this.configPath = join(this.fusionDir, "config.json");
this.asyncLayer = options?.asyncLayer ?? null;
this.consumerId = options?.consumerId ?? null;
const resolvedGlobalSettingsDir = globalSettingsDir
?? (process.env.VITEST === "true" ? join(rootDir, ".fusion-global-settings") : undefined);
this.globalSettingsDir = resolvedGlobalSettingsDir;
@@ -536,7 +547,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
Safe lifecycle emission invokes listeners directly to isolate listener failures, so it bypasses
EventEmitter.emit. Decorate this path too; otherwise the hot update surfaces silently miss lanes.
*/
public emitTaskLifecycleEventSafely( event: "task:created" | "task:updated", args: TaskStoreEvents["task:created"] | TaskStoreEvents["task:updated"], ): boolean {
public emitTaskLifecycleEventSafely( event: "task:created" | "task:updated" | "task:deleted", args: TaskStoreEvents["task:created"] | TaskStoreEvents["task:updated"] | TaskStoreEvents["task:deleted"], ): boolean {
if (event === "task:updated" && args.length === 1) {
const task = args[0] as Task;
const lanes = this.laneCache.get(task.id);
@@ -545,6 +556,39 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
return emitTaskLifecycleEventSafelyImpl(this, event, args);
}
/** Dispatch a committed outbox row without replaying delete-writer side effects. */
public emitObservedTaskDeleted(
task: Task,
outboxEventId: string,
metadata: Pick<NonNullable<TaskStoreEvents["task:deleted"][1]>, "githubIssueAction" | "closureContext"> = {},
): boolean {
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-13:03:
Observed delivery must retain the delete's integration intent while marking it replay-safe.
Bridges can report the original split handoff/action, but listener-owned GitHub/GitLab mutation
paths must not repeat a committed writer's effects during at-least-once crash-window delivery.
*/
this.taskCache.delete(task.id);
return emitTaskLifecycleEventSafelyImpl(this, "task:deleted", [task, {
...metadata,
observed: true,
outboxEventId,
}]);
}
/** Start durable task:deleted observation only for explicitly named backend consumers. */
async startTaskDeletedOutboxConsumer(): Promise<void> {
if (!this.asyncLayer || !this.consumerId || this.taskDeletedOutboxConsumer) return;
this.taskDeletedOutboxConsumer = new TaskDeletedOutboxConsumer(this);
await this.taskDeletedOutboxConsumer.start();
}
async stopTaskDeletedOutboxConsumer(): Promise<void> {
const consumer = this.taskDeletedOutboxConsumer;
this.taskDeletedOutboxConsumer = null;
await consumer?.stop();
}
/**
* FNXC:RuntimeBackendInjection 2026-06-24-14:10: In backend mode this getter must never be reached (all access via async layer).
* Reaching it is a programming error — throws rather than constructing SQLite.

View File

@@ -64,7 +64,7 @@ function attachMailbox(store: TaskStore, notices: string[]): void {
const payload: TaskDeletedLifecyclePayload = {
taskId: "FN-1", previousColumn: "todo", previousStatus: null,
deletedAt: "2026-08-01T10:33:00.000Z", allowResurrection: false,
githubIssueAction: null, deletedBy: null,
githubIssueAction: null, closureContext: null, deletedBy: null,
};
describe("task lifecycle outbox identity", () => {
@@ -73,7 +73,7 @@ describe("task lifecycle outbox identity", () => {
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",
"allowResurrection", "closureContext", "deletedAt", "deletedBy", "githubIssueAction",
"previousColumn", "previousStatus", "taskId",
]);
});

View File

@@ -266,6 +266,7 @@ async function deleteTaskBackendWithClaimResultImpl(store: TaskStore, id: string
deletedAt,
allowResurrection,
githubIssueAction: options?.githubIssueAction ?? null,
closureContext: options?.closureContext ?? null,
deletedBy: options?.auditContext?.agentId ?? null,
},
});

View File

@@ -362,7 +362,14 @@ export function setupActivityLogListenersImpl(store: TaskStore): void {
});
// Task deleted
store.on("task:deleted", (task) => {
store.on("task:deleted", (task, meta) => {
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
Observed outbox replay is intentionally at-least-once. Activity is writer-owned and
accumulating, so observed events must not create a duplicate activity/audit side effect.
Cache eviction and bridge fan-out remain safe for duplicate observed notifications.
*/
if (meta?.observed) return;
store.recordActivityFromListener(
{
type: "task:deleted",
@@ -604,6 +611,8 @@ export async function watchImpl(store: TaskStore): Promise<void> {
for (const task of tasks) {
store.taskCache.set(task.id, { ...task });
}
// Cache must be populated before observed deletes can safely evict local task state.
await store.startTaskDeletedOutboxConsumer();
}
export async function migrateAgentLogEntriesImpl(store: TaskStore): Promise<void> {
@@ -1127,7 +1136,7 @@ export async function migrateLegacyWorkflowStepsImpl(store: TaskStore): Promise<
*/
export function emitTaskLifecycleEventSafelyImpl(
store: TaskStore,
event: "task:created" | "task:updated",
event: "task:created" | "task:updated" | "task:deleted",
args: Parameters<TaskStore["emitTaskLifecycleEventSafely"]>[1],
): boolean {
const listeners = store.listeners(event) as Array<(...listenerArgs: typeof args) => unknown>;

View File

@@ -1,6 +1,7 @@
import { createHash } from "node:crypto";
import { sql } from "drizzle-orm";
import type { DbTransaction } from "../postgres/data-layer.js";
import type { TaskDeleteClosureContext } from "../types.js";
export type TaskDeletedLifecyclePayload = {
taskId: string;
@@ -9,6 +10,7 @@ export type TaskDeletedLifecyclePayload = {
deletedAt: string;
allowResurrection: boolean;
githubIssueAction: string | null;
closureContext: TaskDeleteClosureContext | null;
deletedBy: string | null;
};

View File

@@ -45,7 +45,7 @@ import {recordRunAuditEvent as recordRunAuditEventAsync} from "../postgres/data-
import {listGoalCitations as listGoalCitationsAsync} from "../task-store/async-events.js";
import type {RunAuditEventRow} from "../task-store/row-types.js";
export async function getOrCreateForProjectImpl(store: typeof TaskStore, projectId?: string, centralCore?: CentralCore, globalSettingsDir?: string, asyncLayer?: AsyncDataLayer,): Promise<TaskStore> {
export async function getOrCreateForProjectImpl(store: typeof TaskStore, projectId?: string, centralCore?: CentralCore, globalSettingsDir?: string, asyncLayer?: AsyncDataLayer, consumerId?: string,): Promise<TaskStore> {
if (!asyncLayer) {
throw new Error("TaskStore.getOrCreateForProject requires a project-bound PostgreSQL AsyncDataLayer");
}
@@ -79,7 +79,7 @@ export async function getOrCreateForProjectImpl(store: typeof TaskStore, project
const store = new TaskStore(
context.workingDirectory,
resolvedGlobalSettingsDir,
{ asyncLayer },
{ asyncLayer, ...(consumerId ? { consumerId } : {}) },
);
await store.init();
return store;

View File

@@ -0,0 +1,333 @@
import { randomUUID } from "node:crypto";
import { and, eq, isNull } from "drizzle-orm";
import type { TaskStore } from "../store.js";
import { createLogger } from "../logger.js";
import { recordRunAuditEvent } from "../postgres/data-layer.js";
import * as schema from "../postgres/schema/index.js";
import {
acknowledgeTaskLifecycleEvent,
acquireTaskLifecycleLease,
advanceTaskLifecycleConsumerCursor,
hasTaskLifecycleConsumerReceipt,
listTaskLifecycleEvents,
registerTaskLifecycleConsumer,
releaseTaskLifecycleLease,
readTaskLifecycleConsumerCursor,
readTaskLifecycleEventBounds,
renewTaskLifecycleLease,
setTaskLifecycleConsumerActive,
setTaskLifecycleConsumerRetry,
parkTaskLifecycleConsumerDeadLetter,
type TaskLifecycleLease,
} from "./task-lifecycle-consumer-registry.js";
export const TASK_DELETED_OUTBOX_POLL_MS = 5_000;
export const TASK_DELETED_OUTBOX_LEASE_MS = 15_000;
export const TASK_DELETED_OUTBOX_BATCH_SIZE = 100;
export const TASK_DELETED_OUTBOX_RETENTION_DAYS = 30;
const outboxConsumerLog = createLogger("task-deleted-outbox-consumer");
type OutboxEventForValidation = {
eventId: string;
eventType: string;
taskId: string;
occurredAt: string;
payload: unknown;
};
type ReconciliationReason = "cursor-older-than-retention-bound" | "pruned-gap";
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-12:14:
* Reject malformed durable rows so poison handling, rather than acknowledgement, owns them.
*/
function assertTaskDeletedOutboxEvent(event: OutboxEventForValidation): void {
const payload = event.payload;
if (event.eventType !== "task:deleted" || !event.eventId || !event.taskId || !event.occurredAt
|| !payload || typeof payload !== "object" || Array.isArray(payload)) {
throw new TypeError("Malformed task:deleted lifecycle outbox event");
}
const deleted = payload as Record<string, unknown>;
if (deleted.taskId !== event.taskId || typeof deleted.previousColumn !== "string"
|| (deleted.previousStatus !== null && typeof deleted.previousStatus !== "string")
|| typeof deleted.deletedAt !== "string" || typeof deleted.allowResurrection !== "boolean"
|| (deleted.githubIssueAction !== null && typeof deleted.githubIssueAction !== "string")
|| (deleted.closureContext !== null && (!deleted.closureContext || typeof deleted.closureContext !== "object"
|| Array.isArray(deleted.closureContext)))
|| (deleted.deletedBy !== null && typeof deleted.deletedBy !== "string")) {
throw new TypeError("Malformed task:deleted lifecycle outbox payload");
}
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* The PostgreSQL outbox is authoritative for cross-process task deletion. Delivery dispatches
* before the durable receipt/cursor acknowledgement, intentionally yielding at-least-once
* observed notifications in the crash window; observed dispatch has no writer-owned effects.
*/
export class TaskDeletedOutboxConsumer {
private pollTimer: ReturnType<typeof setInterval> | null = null;
private renewalTimer: ReturnType<typeof setInterval> | null = null;
private running = false;
private lease: TaskLifecycleLease | null = null;
constructor(private readonly store: TaskStore) {}
async start(): Promise<void> {
if (this.running || !this.store.asyncLayer || !this.store.consumerId) return;
this.running = true;
await this.pollSafely();
this.pollTimer = setInterval(() => void this.pollSafely(), TASK_DELETED_OUTBOX_POLL_MS);
}
async stop(): Promise<void> {
this.running = false;
if (this.pollTimer) clearInterval(this.pollTimer);
if (this.renewalTimer) clearInterval(this.renewalTimer);
this.pollTimer = null;
this.renewalTimer = null;
const layer = this.store.asyncLayer;
const consumerId = this.store.consumerId;
if (layer && consumerId && this.lease) {
try {
await releaseTaskLifecycleLease(layer, consumerId, this.lease);
} catch (error) {
outboxConsumerLog.warn("Could not release task:deleted outbox lease during shutdown", error);
}
}
this.lease = null;
if (layer && consumerId) {
try {
await setTaskLifecycleConsumerActive(layer, consumerId, false);
} catch (error) {
outboxConsumerLog.warn("Could not deactivate task:deleted outbox consumer during shutdown", error);
}
}
}
private async pollSafely(): Promise<void> {
try {
await this.poll();
} catch (error) {
outboxConsumerLog.warn("Task:deleted outbox poll failed; delivery will retry", error);
}
}
async poll(): Promise<void> {
const layer = this.store.asyncLayer;
const consumerId = this.store.consumerId;
if (!this.running || !layer || !consumerId) return;
await registerTaskLifecycleConsumer(layer, consumerId);
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:14:
stop() can race the asynchronous registration write. Re-marking this identity inactive after
that race prevents a cleanly stopped consumer from pinning retention as live.
*/
if (!this.running) {
await setTaskLifecycleConsumerActive(layer, consumerId, false);
return;
}
const now = new Date();
const acquired = await acquireTaskLifecycleLease(
layer,
consumerId,
randomUUID(),
new Date(now.getTime() + TASK_DELETED_OUTBOX_LEASE_MS).toISOString(),
now.toISOString(),
);
if (!acquired) return;
if (!this.running) {
await releaseTaskLifecycleLease(layer, consumerId, acquired);
return;
}
this.lease = acquired;
this.startRenewal(acquired);
try {
const cursor = await readTaskLifecycleConsumerCursor(layer, consumerId);
if (!cursor) return;
if (cursor.retryBackoffUntil && Date.parse(cursor.retryBackoffUntil) > Date.now()) return;
const reconciliationReason = await this.needsReconciliation(cursor.lastAckedSeq, cursor.updatedAt);
if (reconciliationReason) {
const reconciled = await this.reconcile(cursor.lastAckedSeq, acquired, reconciliationReason);
if (!reconciled) {
await this.recordLeaseFenced(acquired, 0);
return;
}
}
const currentCursor = await readTaskLifecycleConsumerCursor(layer, consumerId);
if (!currentCursor) return;
const events = await listTaskLifecycleEvents(layer, currentCursor.lastAckedSeq, TASK_DELETED_OUTBOX_BATCH_SIZE);
let priorSeq = currentCursor.lastAckedSeq;
let dispatchedCount = 0;
for (const event of events) {
if (!this.running || this.lease?.fencingToken !== acquired.fencingToken) break;
try {
assertTaskDeletedOutboxEvent(event);
if (await hasTaskLifecycleConsumerReceipt(layer, consumerId, event.eventId)) {
priorSeq = event.seq;
continue;
}
const task = await this.readDeletedTask(event.taskId);
if (!task) {
// Cache absence is not an idempotency gate: commit a receipt for every valid row.
const acknowledged = await acknowledgeTaskLifecycleEvent(layer, {
consumerId, eventId: event.eventId, seq: event.seq, priorSeq, fencingToken: acquired.fencingToken,
});
if (!acknowledged) {
await this.recordLeaseFenced(acquired, 1);
break;
}
priorSeq = event.seq;
continue;
}
const payload = event.payload as {
githubIssueAction: import("../types.js").GithubIssueAction | null;
closureContext: import("../types.js").TaskDeleteClosureContext | null;
};
this.store.emitObservedTaskDeleted(task, event.eventId, {
githubIssueAction: payload.githubIssueAction ?? "auto",
...(payload.closureContext ? { closureContext: payload.closureContext } : {}),
});
dispatchedCount++;
const acknowledged = await acknowledgeTaskLifecycleEvent(layer, {
consumerId, eventId: event.eventId, seq: event.seq, priorSeq, fencingToken: acquired.fencingToken,
});
if (!acknowledged) {
await this.recordLeaseFenced(acquired, 1);
break;
}
priorSeq = event.seq;
} catch (error) {
const attempts = currentCursor.retryAttempts + 1;
const failureClass = error instanceof Error ? error.name : "unknown";
if (attempts >= 10) {
const parked = await parkTaskLifecycleConsumerDeadLetter(layer, {
consumerId, eventId: event.eventId, seq: event.seq, priorSeq, attempts, failureClass, lease: acquired,
});
if (!parked) await this.recordLeaseFenced(acquired, 1);
else priorSeq = event.seq;
break;
}
const delayMs = [1_000, 5_000, 30_000, 300_000, 900_000][Math.min(attempts - 1, 4)]!;
const retried = await setTaskLifecycleConsumerRetry(layer, consumerId, acquired, attempts,
new Date(Date.now() + delayMs).toISOString(),
);
if (!retried) await this.recordLeaseFenced(acquired, 1);
break;
}
}
if (this.running && events.length > 0) {
await recordRunAuditEvent(layer, {
agentId: "system",
runId: `task-deleted-outbox:${consumerId}`,
domain: "task-lifecycle",
mutationType: "task-deleted-outbox:catch-up",
target: consumerId,
metadata: { projectId: layer.projectId, consumerId, fromSeq: currentCursor.lastAckedSeq.toString(), toSeq: priorSeq.toString(), dispatchedCount },
});
}
if (this.running) await setTaskLifecycleConsumerActive(layer, consumerId, true);
} finally {
if (this.renewalTimer) clearInterval(this.renewalTimer);
this.renewalTimer = null;
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:06:
Every completed batch releases its own fenced lease instead of waiting for TTL expiry. This
keeps normal polling responsive while the token predicate protects a successor's reclaim.
*/
if (this.lease?.token === acquired.token) {
await releaseTaskLifecycleLease(layer, consumerId, acquired);
this.lease = null;
}
}
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* Capture the outbox head before reading task state. Advancing only to that fenced snapshot
* preserves rows inserted during reconciliation for the following ordinary poll.
*/
private async needsReconciliation(lastAckedSeq: bigint, updatedAt: string): Promise<ReconciliationReason | null> {
const bounds = await readTaskLifecycleEventBounds(this.store.asyncLayer!);
if (bounds.oldestSeq !== null && lastAckedSeq + 1n < bounds.oldestSeq) return "pruned-gap";
if (Date.parse(updatedAt) < Date.now() - TASK_DELETED_OUTBOX_RETENTION_DAYS * 86_400_000) {
return "cursor-older-than-retention-bound";
}
return null;
}
private async reconcile(
priorSeq: bigint,
lease: TaskLifecycleLease,
reason: ReconciliationReason,
): Promise<boolean> {
const layer = this.store.asyncLayer!;
const consumerId = this.store.consumerId!;
const bounds = await readTaskLifecycleEventBounds(layer);
const headSeq = bounds.headSeq;
const liveRows = await layer.db.select({ id: schema.project.tasks.id })
.from(schema.project.tasks)
.where(and(eq(schema.project.tasks.projectId, layer.projectId!), isNull(schema.project.tasks.deletedAt)));
const liveIds = new Set(liveRows.map((row) => row.id));
let dispatchedCount = 0;
for (const task of this.store.taskCache.values()) {
if (!liveIds.has(task.id)) {
this.store.emitObservedTaskDeleted(task, `reconciliation:${task.id}:${headSeq}`);
dispatchedCount++;
}
}
const advanced = await advanceTaskLifecycleConsumerCursor(layer, consumerId, priorSeq, headSeq, lease.fencingToken);
if (!advanced) return false;
await recordRunAuditEvent(layer, {
agentId: "system", runId: `task-deleted-outbox:${consumerId}`, domain: "task-lifecycle",
mutationType: "task-deleted-outbox:reconciliation-fallback", target: consumerId,
metadata: { projectId: layer.projectId, consumerId, reason, reconciliationHeadSeq: headSeq.toString(), dispatchedCount, scannedCount: liveRows.length },
});
return true;
}
private async recordLeaseFenced(lease: TaskLifecycleLease, abortedCount: number): Promise<void> {
const layer = this.store.asyncLayer;
const consumerId = this.store.consumerId;
if (!layer || !consumerId) return;
const cursor = await readTaskLifecycleConsumerCursor(layer, consumerId);
await recordRunAuditEvent(layer, {
agentId: "system", runId: `task-deleted-outbox:${consumerId}`, domain: "task-lifecycle",
mutationType: "task-deleted-outbox:lease-fenced", target: consumerId,
metadata: {
projectId: layer.projectId, consumerId, staleToken: lease.fencingToken.toString(),
currentToken: (cursor?.fencingToken ?? lease.fencingToken).toString(), abortedCount,
},
});
}
private startRenewal(lease: TaskLifecycleLease): void {
if (this.renewalTimer) clearInterval(this.renewalTimer);
this.renewalTimer = setInterval(() => {
const layer = this.store.asyncLayer;
const consumerId = this.store.consumerId;
if (!layer || !consumerId || this.lease?.fencingToken !== lease.fencingToken) return;
const now = new Date();
void renewTaskLifecycleLease(layer, consumerId, lease,
new Date(now.getTime() + TASK_DELETED_OUTBOX_LEASE_MS).toISOString(), now.toISOString(),
).then((renewed) => {
if (!renewed && this.lease?.fencingToken === lease.fencingToken) this.lease = null;
}).catch((error) => {
outboxConsumerLog.warn("Could not renew task:deleted outbox lease", error);
});
}, Math.floor(TASK_DELETED_OUTBOX_LEASE_MS / 3));
}
private async readDeletedTask(taskId: string) {
const cached = this.store.taskCache.get(taskId);
if (cached) return cached;
const layer = this.store.asyncLayer!;
if (!layer.projectId) return null;
const [row] = await layer.db.select().from(schema.project.tasks).where(and(
eq(schema.project.tasks.projectId, layer.projectId),
eq(schema.project.tasks.id, taskId),
)).limit(1);
return row ? this.store.rowToTask(this.store.pgRowToTaskRow(row as Record<string, unknown>)) : null;
}
}

View File

@@ -0,0 +1,32 @@
export const TASK_LIFECYCLE_CONSUMER_ROLES = [
"engine",
"dashboard",
"cli",
"child-process-worker",
"remote-node",
] as const;
export type TaskLifecycleConsumerRole = (typeof TASK_LIFECYCLE_CONSUMER_ROLES)[number];
const INSTANCE_KEY_PATTERN = /^[A-Za-z0-9._-]{1,128}$/;
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
Consumer identity validation covers role and string format only. A helper cannot distinguish a
persisted node ID from a UUID minted at boot, so restart durability is a wiring-site provenance
obligation that must be proven where the instance key originates.
*/
/** Build a stable durable-consumer identity. */
export function buildConsumerId(
role: TaskLifecycleConsumerRole,
instanceKey?: string,
): string {
if (!TASK_LIFECYCLE_CONSUMER_ROLES.includes(role)) {
throw new Error(`Unknown task lifecycle consumer role: ${String(role)}`);
}
if (instanceKey === undefined) return role;
if (instanceKey.trim() !== instanceKey || !INSTANCE_KEY_PATTERN.test(instanceKey)) {
throw new Error("Task lifecycle consumer instance key must be a trimmed 1-128 character safe identifier");
}
return `${role}:${instanceKey}`;
}

View File

@@ -0,0 +1,396 @@
import { and, asc, eq, gt, lte, or, sql } from "drizzle-orm";
import * as schema from "../postgres/schema/index.js";
import { recordRunAuditEventWithinTransaction, type AsyncDataLayer, type DbTransaction } from "../postgres/data-layer.js";
/** Durable state of one independently observing task-lifecycle consumer. */
export interface TaskLifecycleConsumerCursor {
readonly lastAckedSeq: bigint;
readonly retryAttempts: number;
readonly retryBackoffUntil: string | null;
readonly leaseToken: string | null;
readonly fencingToken: bigint;
readonly leaseExpiresAt: string | null;
readonly updatedAt: string;
}
export interface TaskLifecycleLease {
readonly token: string;
readonly fencingToken: bigint;
readonly expiresAt: string;
}
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:27:
A fenced dead-letter park must roll back its row and audit write together. This private signal
crosses only the transaction boundary, converting a stale holder into the consumer's required
lease-fenced audit path rather than an unhandled polling failure.
*/
class TaskLifecycleConsumerFenceRejectedError extends Error {}
function projectIdFor(layer: AsyncDataLayer): string {
if (!layer.projectId) {
throw new Error("Task lifecycle consumer state requires asyncLayer.projectId");
}
return layer.projectId;
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* Registration, rather than cursor existence, declares a durable consumer live. A registered
* identity with no cursor deliberately blocks retention because it has not acknowledged any event.
*/
export async function registerTaskLifecycleConsumer(
layer: AsyncDataLayer,
consumerId: string,
now = new Date().toISOString(),
): Promise<void> {
const projectId = projectIdFor(layer);
await layer.db.insert(schema.project.taskLifecycleConsumerRegistrations).values({
projectId,
consumerId,
registeredAt: now,
lastSeenAt: now,
active: 1,
}).onConflictDoUpdate({
target: [
schema.project.taskLifecycleConsumerRegistrations.projectId,
schema.project.taskLifecycleConsumerRegistrations.consumerId,
],
set: { lastSeenAt: now, active: 1 },
});
}
export async function setTaskLifecycleConsumerActive(
layer: AsyncDataLayer,
consumerId: string,
active: boolean,
now = new Date().toISOString(),
): Promise<void> {
const projectId = projectIdFor(layer);
await layer.db.update(schema.project.taskLifecycleConsumerRegistrations)
.set({ active: active ? 1 : 0, lastSeenAt: now })
.where(and(
eq(schema.project.taskLifecycleConsumerRegistrations.projectId, projectId),
eq(schema.project.taskLifecycleConsumerRegistrations.consumerId, consumerId),
));
}
export async function readTaskLifecycleConsumerCursor(
layer: AsyncDataLayer,
consumerId: string,
): Promise<TaskLifecycleConsumerCursor | null> {
const projectId = projectIdFor(layer);
const [row] = await layer.db.select().from(schema.project.taskLifecycleConsumerCursors).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
)).limit(1);
return row ?? null;
}
/** Create an empty cursor only when a consumer starts observing. */
export async function ensureTaskLifecycleConsumerCursor(
layer: AsyncDataLayer,
consumerId: string,
now = new Date().toISOString(),
): Promise<TaskLifecycleConsumerCursor> {
const projectId = projectIdFor(layer);
await layer.db.insert(schema.project.taskLifecycleConsumerCursors).values({
projectId,
consumerId,
updatedAt: now,
}).onConflictDoNothing();
const cursor = await readTaskLifecycleConsumerCursor(layer, consumerId);
if (!cursor) throw new Error("Task lifecycle consumer cursor was not persisted");
return cursor;
}
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:40:
Retention can remove every retained event while the project-scoped sequence counter remains ahead.
Read the counter for the reconciliation boundary so an empty retained window never regresses a
consumer cursor and replays already-accounted-for lifecycle history.
*/
export async function readTaskLifecycleEventBounds(layer: AsyncDataLayer): Promise<{ oldestSeq: bigint | null; oldestOccurredAt: string | null; headSeq: bigint }> {
const projectId = projectIdFor(layer);
const [oldest] = await layer.db.select({ seq: schema.project.taskLifecycleEvents.seq, occurredAt: schema.project.taskLifecycleEvents.occurredAt })
.from(schema.project.taskLifecycleEvents).where(eq(schema.project.taskLifecycleEvents.projectId, projectId))
.orderBy(asc(schema.project.taskLifecycleEvents.seq)).limit(1);
const [sequence] = await layer.db.select({ lastSeq: schema.project.taskLifecycleEventSeq.lastSeq })
.from(schema.project.taskLifecycleEventSeq)
.where(eq(schema.project.taskLifecycleEventSeq.projectId, projectId))
.limit(1);
return { oldestSeq: oldest?.seq ?? null, oldestOccurredAt: oldest?.occurredAt ?? null, headSeq: sequence?.lastSeq ?? 0n };
}
export async function listTaskLifecycleEvents(
layer: AsyncDataLayer,
afterSeq: bigint,
limit: number,
) {
const projectId = projectIdFor(layer);
return layer.db.select().from(schema.project.taskLifecycleEvents).where(and(
eq(schema.project.taskLifecycleEvents.projectId, projectId),
gt(schema.project.taskLifecycleEvents.seq, afterSeq),
)).orderBy(asc(schema.project.taskLifecycleEvents.seq)).limit(limit);
}
export async function hasTaskLifecycleConsumerReceipt(
layer: AsyncDataLayer,
consumerId: string,
eventId: string,
): Promise<boolean> {
const projectId = projectIdFor(layer);
const rows = await layer.db.select({ eventId: schema.project.taskLifecycleConsumerReceipts.eventId })
.from(schema.project.taskLifecycleConsumerReceipts)
.where(and(
eq(schema.project.taskLifecycleConsumerReceipts.projectId, projectId),
eq(schema.project.taskLifecycleConsumerReceipts.consumerId, consumerId),
eq(schema.project.taskLifecycleConsumerReceipts.eventId, eventId),
)).limit(1);
return rows.length > 0;
}
/**
* Atomically records the durable receipt and advances a cursor only for the lease fencing token
* that dispatched it. A false return is a stale-holder signal; callers must stop their batch.
*/
export async function acknowledgeTaskLifecycleEvent(
layer: AsyncDataLayer,
input: { consumerId: string; eventId: string; seq: bigint; priorSeq: bigint; fencingToken: bigint; now?: string },
): Promise<boolean> {
const projectId = projectIdFor(layer);
const now = input.now ?? new Date().toISOString();
return layer.transactionImmediate(async (tx) => {
const advanced = await advanceCursorWithFence(tx, projectId, input.consumerId, input.priorSeq, input.seq, input.fencingToken, now);
if (!advanced) return false;
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:06:
A successful acknowledgement starts the next event with a clean retry budget. Persist this
reset in the same receipt/cursor transaction so one transient failure cannot poison a later row.
*/
await tx.update(schema.project.taskLifecycleConsumerCursors).set({
retryAttempts: 0,
retryBackoffUntil: null,
updatedAt: now,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, input.consumerId),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, input.fencingToken),
eq(schema.project.taskLifecycleConsumerCursors.lastAckedSeq, input.seq),
));
await tx.insert(schema.project.taskLifecycleConsumerReceipts).values({
projectId,
consumerId: input.consumerId,
eventId: input.eventId,
seq: input.seq,
processedAt: now,
}).onConflictDoNothing();
return true;
});
}
/** Shared CAS primitive for normal acknowledgement, reconciliation, and poison parking. */
export async function advanceCursorWithFence(
tx: DbTransaction,
projectId: string,
consumerId: string,
priorSeq: bigint,
nextSeq: bigint,
fencingToken: bigint,
now: string,
): Promise<boolean> {
if (nextSeq < priorSeq) return false;
const result = await tx.update(schema.project.taskLifecycleConsumerCursors).set({
lastAckedSeq: nextSeq,
updatedAt: now,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
eq(schema.project.taskLifecycleConsumerCursors.lastAckedSeq, priorSeq),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, fencingToken),
)).returning({ consumerId: schema.project.taskLifecycleConsumerCursors.consumerId });
return result.length === 1;
}
/** Advance a captured reconciliation boundary without reading a newer outbox head. */
export async function advanceTaskLifecycleConsumerCursor(
layer: AsyncDataLayer,
consumerId: string,
priorSeq: bigint,
nextSeq: bigint,
fencingToken: bigint,
now = new Date().toISOString(),
): Promise<boolean> {
const projectId = projectIdFor(layer);
return layer.transactionImmediate((tx) =>
advanceCursorWithFence(tx, projectId, consumerId, priorSeq, nextSeq, fencingToken, now));
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* Lease claims are conditional on expiry and increment the fencing token. TTL is only a liveness
* aid: every acknowledgement still checks this token so an expired holder cannot advance state.
*/
export async function acquireTaskLifecycleLease(
layer: AsyncDataLayer,
consumerId: string,
token: string,
expiresAt: string,
now = new Date().toISOString(),
): Promise<TaskLifecycleLease | null> {
const projectId = projectIdFor(layer);
await ensureTaskLifecycleConsumerCursor(layer, consumerId, now);
const claimed = await layer.db.update(schema.project.taskLifecycleConsumerCursors).set({
leaseToken: token,
leaseExpiresAt: expiresAt,
fencingToken: sql`${schema.project.taskLifecycleConsumerCursors.fencingToken} + 1`,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
or(
lte(schema.project.taskLifecycleConsumerCursors.leaseExpiresAt, now),
sql`${schema.project.taskLifecycleConsumerCursors.leaseExpiresAt} IS NULL`,
),
)).returning({ fencingToken: schema.project.taskLifecycleConsumerCursors.fencingToken });
const row = claimed[0];
return row ? { token, fencingToken: row.fencingToken, expiresAt } : null;
}
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:27:
Lease bookkeeping must not overwrite the cursor's acknowledgement timestamp. Reconnect logic uses
that timestamp to select the 30-day reconciliation fallback, so only acknowledgement/reconciliation
advances update it; lease claims, renewals, releases, and retry scheduling do not.
*/
/** Renewal preserves the fencing token; only an expired-lease reclaim mints a higher one. */
export async function renewTaskLifecycleLease(
layer: AsyncDataLayer,
consumerId: string,
lease: TaskLifecycleLease,
expiresAt: string,
_now = new Date().toISOString(),
): Promise<boolean> {
const projectId = projectIdFor(layer);
const renewed = await layer.db.update(schema.project.taskLifecycleConsumerCursors).set({
leaseExpiresAt: expiresAt,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
eq(schema.project.taskLifecycleConsumerCursors.leaseToken, lease.token),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, lease.fencingToken),
)).returning({ consumerId: schema.project.taskLifecycleConsumerCursors.consumerId });
return renewed.length === 1;
}
/** Clean shutdown releases only its own fenced lease and cannot clear a successor's lease. */
export async function setTaskLifecycleConsumerRetry(
layer: AsyncDataLayer,
consumerId: string,
lease: TaskLifecycleLease,
attempts: number,
backoffUntil: string | null,
_now = new Date().toISOString(),
): Promise<boolean> {
const projectId = projectIdFor(layer);
const updated = await layer.db.update(schema.project.taskLifecycleConsumerCursors).set({
retryAttempts: attempts,
retryBackoffUntil: backoffUntil,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
eq(schema.project.taskLifecycleConsumerCursors.leaseToken, lease.token),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, lease.fencingToken),
)).returning({ consumerId: schema.project.taskLifecycleConsumerCursors.consumerId });
return updated.length === 1;
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* Poison parking is one transaction: an idempotent dead-letter insert, fenced cursor advance,
* retry reset, and audit row commit or roll back together. A failed fenced advance throws only
* inside the transaction so the catch below can report a stale lease without committing an orphan
* dead-letter or audit row.
*/
export async function parkTaskLifecycleConsumerDeadLetter(
layer: AsyncDataLayer,
input: {
consumerId: string;
eventId: string;
seq: bigint;
priorSeq: bigint;
attempts: number;
failureClass: string;
lease: TaskLifecycleLease;
now?: string;
},
): Promise<boolean> {
const projectId = projectIdFor(layer);
const now = input.now ?? new Date().toISOString();
try {
return await layer.transactionImmediate(async (tx) => {
const inserted = await tx.insert(schema.project.taskLifecycleConsumerDeadLetters).values({
projectId,
consumerId: input.consumerId,
eventId: input.eventId,
seq: input.seq,
attempts: input.attempts,
failureClass: input.failureClass,
parkedAt: now,
updatedAt: now,
}).onConflictDoNothing().returning({ eventId: schema.project.taskLifecycleConsumerDeadLetters.eventId });
// A prior committed park is complete; do not manufacture another audit row on a replay.
if (inserted.length === 0) return true;
const advanced = await advanceCursorWithFence(
tx, projectId, input.consumerId, input.priorSeq, input.seq, input.lease.fencingToken, now,
);
if (!advanced) throw new TaskLifecycleConsumerFenceRejectedError();
await tx.update(schema.project.taskLifecycleConsumerCursors).set({
retryAttempts: 0,
retryBackoffUntil: null,
updatedAt: now,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, input.consumerId),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, input.lease.fencingToken),
));
await recordRunAuditEventWithinTransaction(tx, {
taskId: undefined,
agentId: "system",
runId: `task-deleted-outbox:${input.consumerId}`,
domain: "task-lifecycle",
mutationType: "task-deleted-outbox:dead-letter",
target: input.eventId,
metadata: {
projectId,
consumerId: input.consumerId,
eventId: input.eventId,
seq: input.seq.toString(),
attempts: input.attempts,
failureClass: input.failureClass,
},
});
return true;
});
} catch (error) {
if (error instanceof TaskLifecycleConsumerFenceRejectedError) return false;
throw error;
}
}
export async function releaseTaskLifecycleLease(
layer: AsyncDataLayer,
consumerId: string,
lease: TaskLifecycleLease,
_now = new Date().toISOString(),
): Promise<void> {
const projectId = projectIdFor(layer);
await layer.db.update(schema.project.taskLifecycleConsumerCursors).set({
leaseToken: null,
leaseExpiresAt: null,
}).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, consumerId),
eq(schema.project.taskLifecycleConsumerCursors.leaseToken, lease.token),
eq(schema.project.taskLifecycleConsumerCursors.fencingToken, lease.fencingToken),
));
}

View File

@@ -0,0 +1,104 @@
import { and, asc, eq, lt, lte } from "drizzle-orm";
import * as schema from "../postgres/schema/index.js";
import { recordRunAuditEvent, type AsyncDataLayer } from "../postgres/data-layer.js";
export const TASK_LIFECYCLE_RETENTION_DAYS = 30;
export const TASK_LIFECYCLE_RETENTION_MAX_DELETES = 5_000;
export interface TaskLifecycleRetentionResult {
readonly prunedCount: number;
readonly oldestRetainedSeq: bigint | null;
readonly minAckedSeq: bigint | null;
readonly liveConsumerCount: number;
readonly staleConsumerCount: number;
readonly budgetExhausted: boolean;
}
/**
* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
* This is the sole outbox pruning seam. Live registrations, not cursor rows, protect restartable
* consumers; with no live identity, age-only pruning preserves the 30-day replay contract.
*/
export async function pruneTaskLifecycleEvents(
layer: AsyncDataLayer,
projectId: string,
options: { now?: Date; retentionDays?: number; livenessDays?: number; maxDeletes?: number } = {},
): Promise<TaskLifecycleRetentionResult> {
const now = options.now ?? new Date();
const retentionDays = options.retentionDays ?? TASK_LIFECYCLE_RETENTION_DAYS;
const livenessDays = options.livenessDays ?? retentionDays;
const maxDeletes = options.maxDeletes ?? TASK_LIFECYCLE_RETENTION_MAX_DELETES;
const cutoff = new Date(now.getTime() - retentionDays * 86_400_000).toISOString();
const liveCutoff = new Date(now.getTime() - livenessDays * 86_400_000).toISOString();
const registrations = await layer.db.select().from(schema.project.taskLifecycleConsumerRegistrations)
.where(eq(schema.project.taskLifecycleConsumerRegistrations.projectId, projectId));
const live = registrations.filter((row) => row.active === 1 && row.lastSeenAt >= liveCutoff);
const staleConsumerCount = registrations.length - live.length;
let minAckedSeq: bigint | null = null;
if (live.length > 0) {
const cursorRows = await Promise.all(live.map(async (registration) => {
const [cursor] = await layer.db.select().from(schema.project.taskLifecycleConsumerCursors).where(and(
eq(schema.project.taskLifecycleConsumerCursors.projectId, projectId),
eq(schema.project.taskLifecycleConsumerCursors.consumerId, registration.consumerId),
)).limit(1);
return cursor;
}));
// A registered identity may have started but not acked yet: never prune its history.
if (cursorRows.some((cursor) => !cursor)) {
return { prunedCount: 0, oldestRetainedSeq: null, minAckedSeq: null, liveConsumerCount: live.length, staleConsumerCount, budgetExhausted: false };
}
minAckedSeq = cursorRows.reduce<bigint>((minimum, cursor) =>
cursor!.lastAckedSeq < minimum ? cursor!.lastAckedSeq : minimum,
cursorRows[0]!.lastAckedSeq);
}
const candidates = await layer.db.select({ seq: schema.project.taskLifecycleEvents.seq })
.from(schema.project.taskLifecycleEvents)
.where(live.length === 0
? and(eq(schema.project.taskLifecycleEvents.projectId, projectId), lt(schema.project.taskLifecycleEvents.occurredAt, cutoff))
: and(
eq(schema.project.taskLifecycleEvents.projectId, projectId),
lt(schema.project.taskLifecycleEvents.occurredAt, cutoff),
lte(schema.project.taskLifecycleEvents.seq, minAckedSeq!),
))
.orderBy(asc(schema.project.taskLifecycleEvents.seq))
.limit(maxDeletes);
if (candidates.length > 0) {
await layer.db.delete(schema.project.taskLifecycleEvents).where(and(
eq(schema.project.taskLifecycleEvents.projectId, projectId),
// Bounded candidate selection avoids an unbounded conditional DELETE under concurrent writers.
lte(schema.project.taskLifecycleEvents.seq, candidates[candidates.length - 1]!.seq),
lt(schema.project.taskLifecycleEvents.occurredAt, cutoff),
...(minAckedSeq === null ? [] : [lte(schema.project.taskLifecycleEvents.seq, minAckedSeq)]),
));
}
const [oldest] = await layer.db.select({ seq: schema.project.taskLifecycleEvents.seq })
.from(schema.project.taskLifecycleEvents).where(eq(schema.project.taskLifecycleEvents.projectId, projectId))
.orderBy(asc(schema.project.taskLifecycleEvents.seq)).limit(1);
const result = {
prunedCount: candidates.length,
oldestRetainedSeq: oldest?.seq ?? null,
minAckedSeq,
liveConsumerCount: live.length,
staleConsumerCount,
budgetExhausted: candidates.length === maxDeletes,
};
await recordRunAuditEvent(layer, {
agentId: "system",
runId: "task-deleted-outbox:retention",
domain: "task-lifecycle",
mutationType: "task-deleted-outbox:retention-pruned",
target: projectId,
metadata: {
projectId,
prunedCount: result.prunedCount,
oldestRetainedSeq: result.oldestRetainedSeq?.toString() ?? null,
minAckedSeq: result.minAckedSeq?.toString() ?? null,
liveConsumerCount: result.liveConsumerCount,
staleConsumerCount: result.staleConsumerCount,
budgetExhausted: result.budgetExhausted,
},
});
return result;
}

View File

@@ -1199,6 +1199,7 @@ export async function listApprovedCliAutonomyAdaptersImpl(store: TaskStore): Pro
export async function closeImpl(store: TaskStore): Promise<void> {
store.closing = true;
await store.stopTaskDeletedOutboxConsumer();
if (store.deferredTaskCreatedWork.size > 0) {
await Promise.allSettled([...store.deferredTaskCreatedWork]);
}

View File

@@ -139,7 +139,11 @@ type GitHubIssueActionEvent = {
error?: string;
};
type TaskDeletedMeta = { githubIssueAction?: GithubIssueAction; closureContext?: TaskDeleteClosureContext };
type TaskDeletedMeta = {
githubIssueAction?: GithubIssueAction;
closureContext?: TaskDeleteClosureContext;
observed?: boolean;
};
type SplitCommentLog = (message: string, details: string) => Promise<void>;
@@ -522,6 +526,12 @@ export class GitHubTrackingStateService {
}
private async handleTaskDeleted(store: TaskStore, task: Task, meta?: TaskDeletedMeta): Promise<void> {
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-13:03:
Durable observers replay committed deletes at least once. GitHub close/delete/comment actions are
writer-owned side effects, so observed notifications are bridge-only and must never invoke them.
*/
if (meta?.observed) return;
if (task.githubTracking?.enabled !== true) {
await this.handleSourceIssueDelete(store, task, meta);
return;

View File

@@ -5,6 +5,7 @@ import { updateGitLabTargetState } from "./gitlab-tracking-state.js";
type TaskDeletedMeta = {
githubIssueAction?: GithubIssueAction;
closureContext?: TaskDeleteClosureContext | { kind?: string };
observed?: boolean;
};
export type GitLabDeleteAction =
@@ -59,6 +60,8 @@ export class GitLabDeleteCloseService {
}
private async handleTaskDeleted(store: TaskStore, task: Task, meta?: TaskDeletedMeta): Promise<void> {
/* FNXC:CrossProcessDeleteObservation 2026-08-01-13:03: Observed outbox replays never repeat GitLab close/delete-side effects. */
if (meta?.observed) return;
let stage = "resolve";
try {
const target = resolveGitLabTarget(task, { fallbackToSourceOnInvalidTracking: true });

View File

@@ -2,7 +2,10 @@ import type { Task, TaskDeleteClosureContext, TaskStore } from "@fusion/core";
import { resolveGitLabClient, resolveGitLabTarget, safeLogGitLabEntry, type GitLabLifecycleTarget } from "./gitlab-lifecycle.js";
import { retryTransient, updateGitLabTargetState } from "./gitlab-tracking-state.js";
type TaskDeletedMeta = { closureContext?: TaskDeleteClosureContext | { kind?: string; childTaskIds?: unknown } };
type TaskDeletedMeta = {
closureContext?: TaskDeleteClosureContext | { kind?: string; childTaskIds?: unknown };
observed?: boolean;
};
export type GitLabSplitNoteResult =
| { status: "posted"; target: GitLabLifecycleTarget }
@@ -84,6 +87,8 @@ export class GitLabSplitCloseService {
}
private async handleTaskDeleted(store: TaskStore, task: Task, meta?: TaskDeletedMeta): Promise<void> {
/* FNXC:CrossProcessDeleteObservation 2026-08-01-13:03: Split notes are writer-owned and must not repeat on at-least-once observed delivery. */
if (meta?.observed) return;
try {
const result = await postGitLabSplitNoteBeforeClose(store, task, meta);
if (result.status === "posted") {

View File

@@ -14,7 +14,7 @@
* const store = await getOrCreateProjectStore(projectId);
*/
import { countRunningAgentTasks, enrichRunningAgentTaskShape, resolveWorkflowIrForTask, type TaskStore } from "@fusion/core";
import { buildConsumerId, countRunningAgentTasks, enrichRunningAgentTaskShape, resolveWorkflowIrForTask, type TaskStore } from "@fusion/core";
/**
* Internal cache: projectId → TaskStore instance.
@@ -110,7 +110,12 @@ export async function getOrCreateProjectStore(projectId: string): Promise<TaskSt
// default when DATABASE_URL is unset, external PG when DATABASE_URL is
// set. The factory applies the schema baseline and integrates the
// dual-read harness when FUSION_DUAL_READ=1.
const backendBoot = await createTaskStoreForBackend({ projectId });
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:14:
This dashboard-owned, SSE-subscribed store needs its own project-scoped stream. The role-only
identity is restart-stable and must not share the engine consumer's cursor or fenced lease.
*/
const backendBoot = await createTaskStoreForBackend({ projectId, consumerId: buildConsumerId("dashboard") });
// FNXC:PostgresFinalCutover 2026-07-14-17:20: Project-scoped dashboard
// stores are always factory-backed PostgreSQL stores.
const store = backendBoot.taskStore;

View File

@@ -18,7 +18,7 @@ export interface ProjectManagerEvents {
/** Emitted when a task is created in any project */
"task:created": [data: { projectId: string; projectName: string; task: Task }];
/** Emitted when a task is deleted in any project */
"task:deleted": [data: { projectId: string; projectName: string; task: Task; meta?: { githubIssueAction?: import("@fusion/core").GithubIssueAction } }];
"task:deleted": [data: { projectId: string; projectName: string; task: Task; meta?: { githubIssueAction?: import("@fusion/core").GithubIssueAction; observed?: boolean; outboxEventId?: string } }];
/** Emitted when a task is moved in any project */
"task:moved": [
data: {
@@ -381,7 +381,7 @@ export class ProjectManager extends EventEmitter<ProjectManagerEvents> {
});
// Forward task:deleted
runtime.on("task:deleted", (task: Task, meta?: { githubIssueAction?: import("@fusion/core").GithubIssueAction }) => {
runtime.on("task:deleted", (task: Task, meta?: { githubIssueAction?: import("@fusion/core").GithubIssueAction; observed?: boolean; outboxEventId?: string }) => {
this.emit("task:deleted", { projectId, projectName, task, meta });
});

View File

@@ -94,7 +94,7 @@ export interface ProjectRuntimeEvents {
/** Emitted when a task is updated */
"task:updated": [task: Task];
/** Emitted when a task is deleted */
"task:deleted": [task: Task, meta?: { githubIssueAction?: GithubIssueAction }];
"task:deleted": [task: Task, meta?: { githubIssueAction?: GithubIssueAction; observed?: boolean; outboxEventId?: string }];
/** Emitted when a cross-node assignment event is observed */
"task:assigned": [data: { taskId: string; agentId: string; assignedAt: string; source?: string }];
/** Emitted when an error occurs in the runtime */

View File

@@ -818,6 +818,7 @@ export class InProcessRuntime
// InProcessRuntime.start(). When the factory returns a backend result,
// the engine owns the result's shutdown() for process teardown.
createTaskStoreForBackend,
buildConsumerId,
createProjectScopedPluginMcpProvider,
registerTaskDeleteNoticeMailbox,
} = await import("@fusion/core");
@@ -828,6 +829,8 @@ export class InProcessRuntime
const backendBoot = await createTaskStoreForBackend({
rootDir: this.config.workingDirectory,
projectId: this.config.projectId,
/* FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: the engine owns this store, so its role-only identity is restart-stable and never derives from boot state. */
consumerId: buildConsumerId("engine"),
onMigrationProgress: this.config.onMigrationProgress,
});
// FNXC:PostgresFinalCutover 2026-07-14-17:20: Engine runtimes must fail
@@ -1712,6 +1715,12 @@ export class InProcessRuntime
// 8. Set up event forwarding from TaskStore
this.setupEventForwarding();
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-13:03:
Engine-owned stores do not call watch(), so runtime startup owns durable delete observation.
Start only after its bridge is attached so the initial poll cannot lose a cross-process delete.
*/
await this.taskStore.startTaskDeletedOutboxConsumer();
const startupSettings = await this.taskStore.getSettings();
if (startupSettings.globalPause || startupSettings.enginePaused) {
@@ -2645,8 +2654,12 @@ export class InProcessRuntime
this.emit("task:updated", task);
});
// Forward task:deleted events
this.taskStore.on("task:deleted", (task: Task, meta?: { githubIssueAction?: GithubIssueAction }) => {
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-12:02:
Preserve observed outbox provenance through the runtime bridge. Downstream project and IPC
bridges must distinguish cross-process deletion delivery without re-running writer effects.
*/
this.taskStore.on("task:deleted", (task: Task, meta?: { githubIssueAction?: GithubIssueAction; observed?: boolean; outboxEventId?: string }) => {
this.recordActivity();
this.emit("task:deleted", task, meta);
});

View File

@@ -35,6 +35,7 @@ import { type TaskMoveLanes, resolveColumnFlags, IN_REVIEW_STALL_DEADLOCK_LOG_PR
TERMINAL_ROLES,
resolveProjectColumnsForRoles,
REVIEW_ROLES,
pruneTaskLifecycleEvents,
} from "@fusion/core";
import { finalizePlanningSegment } from "@fusion/core";
import type { MeshLeaseManager } from "./mesh-lease-manager.js";
@@ -915,6 +916,7 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
private symbolLockNoActionAudited = false;
private preservedQueuedOverlapLogged = new Map<string, string>();
private maintenanceTickCounter = 0;
private readonly taskLifecycleRetentionLastPrunedAt = new Map<string, number>();
private readonly processBootStartedAt = Date.now();
private lastDbCorruptionNotifiedAt: number | null = null;
@@ -2159,6 +2161,27 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
}, intervalMs);
}
/*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
Periodic self-healing owns outbox retention. Each already-open PostgreSQL project is gated to
one bounded prune every six hours; failures stay diagnostic-only so retention never blocks
task execution or another maintenance step.
*/
private async pruneTaskLifecycleEventsForMaintenance(): Promise<void> {
const layer = this.store.getAsyncLayer();
const projectId = layer?.projectId;
if (!layer || !projectId) return;
const now = Date.now();
const lastPrunedAt = this.taskLifecycleRetentionLastPrunedAt.get(projectId) ?? 0;
if (now - lastPrunedAt < 6 * 60 * 60 * 1000) return;
try {
await pruneTaskLifecycleEvents(layer, projectId);
this.taskLifecycleRetentionLastPrunedAt.set(projectId, now);
} catch (error) {
log.warn(`Task lifecycle retention failed for project ${projectId}: ${error instanceof Error ? error.message : String(error)}`);
}
}
private isPastInterruptedMergeGrace(task: Task, timeoutMs: number): boolean {
const updatedAt = task.updatedAt ? Date.parse(task.updatedAt) : 0;
if (!Number.isFinite(updatedAt) || updatedAt <= 0) return false;
@@ -2527,6 +2550,10 @@ export class SelfHealingManager extends SelfHealingGitEvidence {
// Batch 1 — housekeeping (safe under pause: filesystem/db cleanup only)
const batch1Fns: Array<{ name: string; fn: () => Promise<unknown> }> = [
{ name: "prune-worktrees", fn: () => this.pruneWorktrees() },
{
name: "prune-task-lifecycle-events",
fn: async () => this.pruneTaskLifecycleEventsForMaintenance(),
},
{ name: "cleanup-orphans", fn: () => this.cleanupOrphans() },
{
name: "cleanup-stale-temp-merge-worktrees",