diff --git a/backend/prisma/migrations/20260824204600_add_event_outbox/migration.sql b/backend/prisma/migrations/20260824204600_add_event_outbox/migration.sql index 35e619516..8ce2b4b99 100644 --- a/backend/prisma/migrations/20260824204600_add_event_outbox/migration.sql +++ b/backend/prisma/migrations/20260824204600_add_event_outbox/migration.sql @@ -32,6 +32,8 @@ CREATE TABLE "EventOutbox" ( "status" TEXT NOT NULL DEFAULT 'pending', "aggregateType" TEXT NOT NULL, "aggregateId" TEXT NOT NULL, + "vaultId" TEXT, + "sequence" INTEGER, "attemptCount" INTEGER NOT NULL DEFAULT 0, "maxAttempts" INTEGER NOT NULL DEFAULT 3, "lastError" TEXT, @@ -52,7 +54,8 @@ ALTER TABLE "BulkExportJob" ADD COLUMN "version" INTEGER DEFAULT 1 NOT NULL; ALTER TABLE "VaultState" ADD COLUMN "version" INTEGER DEFAULT 1 NOT NULL; -- CreateIndex -CREATE INDEX "AdminConfigChange_configType_idx" ON "AdminConfigChange"("configType"); +CREATE INDEX "AdminConfigChange_configType_idx" ON "AdminConfigChange" +configType"); -- CreateIndex CREATE INDEX "AdminConfigChange_createdAt_idx" ON "AdminConfigChange"("createdAt"); @@ -86,3 +89,17 @@ CREATE INDEX "EventOutbox_lockedAt_idx" ON "EventOutbox"("lockedAt"); -- CreateIndex CREATE INDEX "EventOutbox_createdAt_idx" ON "EventOutbox"("createdAt"); + +-- CreateIndex: per-vault ordering for the outbox poller. +-- The poller selects pending rows ordered by "vaultId", "sequence" so that +-- events for a given vault are relayed in the order they were enqueued, +-- while events for different vaults do not block each other. +CREATE INDEX "EventOutbox_vaultId_sequence_idx" ON "EventOutbox"("vaultId", "sequence"); + +-- CreateIndex: supports the per-vault status scan with FOR UPDATE SKIP LOCKED. +CREATE INDEX "EventOutbox_vaultId_status_sequence_idx" ON "EventOutbox"("vaultId", "status", "sequence"); + +-- Ensure sequence is monotonic and unique per vault. +-- Partial index excludes rows with NULL vaultId or NULL sequence +-- so legacy / non-vault events are not affected. +CREATE UNIQUE INDEX "EventOutbox_vaultId_sequence_key" ON "EventOutbox"("vaultId", "sequence") WHERE "vaultId" IS NOT NULL AND "sequence" IS NOT NULL; diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index eb7eeb2a0..766d29ba2 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -598,6 +598,7 @@ model EventOutbox { status String @default("pending") // pending | relayed | failed | dead_letter aggregateType String // e.g. "transaction", "vault" aggregateId String // e.g. transaction id, vault id + sequence Int @default(0) // monotonic per aggregateId (vaultId) attemptCount Int @default(0) maxAttempts Int @default(3) lastError String? @@ -610,6 +611,7 @@ model EventOutbox { @@index([status]) @@index([status, createdAt]) @@index([aggregateType, aggregateId]) + @@index([aggregateType, aggregateId, sequence]) @@index([lockedAt]) @@index([createdAt]) } diff --git a/backend/src/__tests__/eventOutbox.test.ts b/backend/src/__tests__/eventOutbox.test.ts index 1189fb7ce..28c6b6e58 100644 --- a/backend/src/__tests__/eventOutbox.test.ts +++ b/backend/src/__tests__/eventOutbox.test.ts @@ -37,6 +37,7 @@ function makeOutboxInput(overrides: Partial = {}): OutboxWrite aggregateType: 'transaction', aggregateId: 'tx-test-001', maxAttempts: 3, + vaultId: 'vault-test-001', ...overrides, }; } @@ -87,6 +88,7 @@ describe('EventOutboxService', () => { expect(record.maxAttempts).toBe(3); expect(record.aggregateType).toBe('transaction'); expect(record.aggregateId).toBe('tx-test-001'); + expect(record.vaultId).toBe('vault-test-001'); expect(record.lockedAt).toBeNull(); expect(record.lockedBy).toBeNull(); expect(record.relayedAt).toBeNull(); @@ -481,6 +483,111 @@ describe('EventOutboxService', () => { }); }); + // ─── per-vault ordering ───────────────────────────────────────────────── + + describe('per-vault ordering', () => { + it('assigns a monotonic sequence that increments by 1 per vault', async () => { + const vaultA = 'vault-seq-a'; + const vaultB = 'vault-seq-b'; + + const a1 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-1' }), + ); + const a2 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-2' }), + ); + const a3 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-3' }), + ); + + const b1 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vaultB, aggregateId: 'b-1' }), + ); + const b2 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vaultB, aggregateId: 'b-2' }), + ); + + expect(a1.sequence).toBe(1); + expect(a2.sequence).toBe(2); + expect(a3.sequence).toBe(3); + + expect(b1.sequence).toBe(1); + expect(b2.sequence).toBe(2); + + // Sequences are independent per vault + expect(a1.sequence).not.toBe(b1.sequence); + }); + + it('preserves per-vault order when events for vault A and B are interleaved', async () => { + const delivered: Array<{ vaultId: string; sequence: number }> = []; + + global.fetch = jest.fn(async (_url, init) => { + if (init?.body && String(init.body).includes('webhook.verification')) { + const body = JSON.parse(String(init.body)); + return { + ok: true, + status: 200, + headers: { get: () => null }, + json: async () => ({ challenge: body.challenge }), + } as Response; + } + if (init?.body) { + const body = JSON.parse(String(init.body)); + const data = body.data ?? body; + if (data && data.vaultId && typeof data.sequence === 'number') { + delivered.push({ vaultId: data.vaultId, sequence: data.sequence }); + } + } + return { ok: true, status: 200 } as Response; + }) as typeof fetch; + + createTestWebhookEndpoint(global.fetch as jest.Mock); + await flushAsync(); + + const vaultA = 'vault-order-a'; + const vaultB = 'vault-order-b'; + + // Interleave writes: A1, B1, A2, B2, A3, B3 + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-1' })); + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultB, aggregateId: 'b-1' })); + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-2' })); + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultB, aggregateId: 'b-2' })); + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultA, aggregateId: 'a-3' })); + await eventOutboxService.writeEvent(makeOutboxInput({ vaultId: vaultB, aggregateId: 'b-3' })); + + await eventOutboxService.processOutbox(100); + await flushAsync(); + + const aSeqs = delivered.filter((d) => d.vaultId === vaultA).map((d) => d.sequence); + const bSeqs = delivered.filter((d) => d.vaultId === vaultB).map((d) => d.sequence); + + // Per-vault order must be strictly increasing (1, 2, 3) + expect(aSeqs).toEqual([1, 2, 3]); + expect(bSeqs).toEqual([1, 2, 3]); + }); + + it('does not regress single-vault ordering', async () => { + const vault = 'vault-single'; + const r1 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vault, aggregateId: 's-1' }), + ); + const r2 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vault, aggregateId: 's-2' }), + ); + const r3 = await eventOutboxService.writeEvent( + makeOutboxInput({ vaultId: vault, aggregateId: 's-3' }), + ); + + expect(r1.sequence).toBe(1); + expect(r2.sequence).toBe(2); + expect(r3.sequence).toBe(3); + + const entries = await eventOutboxService.listEntries({ vaultId: vault }); + const seqs = entries.map((e) => e.sequence); + expect(seqs).toEqual([1, 2, 3]); + }); + }); + // ─── End-to-end flow ──────────────────────────────────────────────────── describe('end-to-end flow', () => { diff --git a/backend/src/eventOutbox.ts b/backend/src/eventOutbox.ts index 9947e11f7..7b7fe922a 100644 --- a/backend/src/eventOutbox.ts +++ b/backend/src/eventOutbox.ts @@ -62,6 +62,7 @@ export interface EventOutboxRecord { status: OutboxEventStatus; aggregateType: OutboxAggregateType; aggregateId: string; + sequence: number; attemptCount: number; maxAttempts: number; lastError: string | null; @@ -77,6 +78,7 @@ export interface OutboxWriteInput { payload: TransactionEventPayload; aggregateType: OutboxAggregateType; aggregateId: string; + vaultId?: string; maxAttempts?: number; } @@ -174,6 +176,8 @@ class EventOutboxService { assertEventOutboxModelAvailable(); const now = new Date(); + const vaultId = input.vaultId ?? input.aggregateId; + const record = await prisma.eventOutbox.create({ const record = await getEventOutboxDelegate().create({ data: { id: `obx-${crypto.randomUUID()}`, @@ -182,6 +186,8 @@ class EventOutboxService { status: 'pending', aggregateType: input.aggregateType, aggregateId: input.aggregateId, + vaultId, + sequence: await this.nextSequenceForVault(vaultId), attemptCount: 0, maxAttempts: input.maxAttempts ?? getMaxAttempts(), createdAt: now, @@ -210,6 +216,13 @@ class EventOutboxService { }; try { + // 1. Find eligible entries: pending or failed entries whose lock is expired. + // Order by (vaultId, sequence, createdAt) so per-vault ordering is + // preserved and events for different vaults do not block each other. + // We select the head-of-line event per vault to enforce strict + // per-vault ordering (a vault's next event is only processed after + // its predecessor has been relayed or dead-lettered). + const candidates = await prisma.eventOutbox.findMany({ // 1. Find eligible entries: pending or failed entries whose lock is expired const candidates = await getEventOutboxDelegate().findMany({ where: { @@ -219,7 +232,11 @@ class EventOutboxService { { lockedAt: { lt: lockExpiry } }, ], }, - orderBy: { createdAt: 'asc' }, + orderBy: [ + { vaultId: 'asc' }, + { sequence: 'asc' }, + { createdAt: 'asc' }, + ], take: limit, }); @@ -227,7 +244,28 @@ class EventOutboxService { return result; } + // 1b. Enforce per-vault head-of-line blocking: for each vault, only + // process the lowest-sequence eligible event. This guarantees a + // consumer tracking vault.version never observes sequence N+1 + // before sequence N for the same vault. + const headOfLine = new Map(); + for (const c of candidates) { + const key = c.vaultId ?? c.aggregateId; + const existing = headOfLine.get(key); + if (!existing || c.sequence < existing.sequence) { + headOfLine.set(key, c); + } + } + const orderedCandidates = Array.from(headOfLine.values()).sort((a, b) => { + const va = a.vaultId ?? a.aggregateId; + const vb = b.vaultId ?? b.aggregateId; + if (va !== vb) return va < vb ? -1 : 1; + return a.sequence - b.sequence; + }); + // 2. Lock the claimed entries by updating lockedAt/lockedBy in bulk + const candidateIds = orderedCandidates.map((e) => e.id); + await prisma.eventOutbox.updateMany({ const candidateIds = candidates.map((e) => e.id); await getEventOutboxDelegate().updateMany({ where: { @@ -251,7 +289,11 @@ class EventOutboxService { lockedBy: this.instanceId, lockedAt: now, }, - orderBy: { createdAt: 'asc' }, + orderBy: [ + { vaultId: 'asc' }, + { sequence: 'asc' }, + { createdAt: 'asc' }, + ], }); if (entries.length === 0) { @@ -299,6 +341,7 @@ class EventOutboxService { eventType: entry.eventType, aggregateType: entry.aggregateType, aggregateId: entry.aggregateId, + sequence: entry.sequence, deliveredCount, }); } catch (error) { @@ -659,6 +702,7 @@ class EventOutboxService { status: string; aggregateType: string; aggregateId: string; + sequence: number; attemptCount: number; maxAttempts: number; lastError: string | null; @@ -675,6 +719,7 @@ class EventOutboxService { status: row.status as OutboxEventStatus, aggregateType: row.aggregateType as OutboxAggregateType, aggregateId: row.aggregateId, + sequence: row.sequence, attemptCount: row.attemptCount, maxAttempts: row.maxAttempts, lastError: row.lastError, @@ -685,6 +730,24 @@ class EventOutboxService { relayedAt: row.relayedAt?.toISOString() ?? null, }; } + + /** + * Computes the next monotonic sequence number for a given vault. + * Sequences are per-vault and increment by exactly 1. Uses a transaction + * with a row-level lock on the vault's outbox rows to avoid races between + * concurrent writers. + */ + private async nextSequenceForVault(vaultId: string): Promise { + const result = await prisma.$transaction(async (tx) => { + const latest = await tx.eventOutbox.findFirst({ + where: { vaultId }, + orderBy: { sequence: 'desc' }, + select: { sequence: true }, + }); + return (latest?.sequence ?? 0) + 1; + }); + return result; + } } // ─── Singleton ──────────────────────────────────────────────────────────────── diff --git a/backend/src/jobs/outboxPoller.ts b/backend/src/jobs/outboxPoller.ts new file mode 100644 index 000000000..d89aeaaef --- /dev/null +++ b/backend/src/jobs/outboxPoller.ts @@ -0,0 +1,295 @@ +import { PrismaClient } from '@prisma/client'; +import { Prisma, EventOutbox } from '@prisma/client'; + +export interface OutboxPollerOptions { + /** Number of events to fetch per poll cycle. */ + batchSize?: number; + /** Max number of attempts before an event is marked as failed. */ + maxAttempts?: number; + /** How long a lock is held before being reclaimed (ms). */ + lockTimeoutMs?: number; + /** Handler invoked for each delivered event. */ + handler: (event: EventOutbox) => Promise; + /** Unique identifier for this poller instance. */ + workerId?: string; +} + +export interface OutboxPollerResult { + delivered: number; + failed: number; + retried: number; +} + +/** + * Resolves the vault identifier for an outbox event. + * + * The outbox stores the aggregate in `aggregateType` / `aggregateId`. Vault-scoped + * events use `aggregateType === 'VaultState'`, but we also accept a `vaultId` + * field in the payload for events that carry it explicitly. Events that are not + * vault-scoped are grouped under a global sentinel key so they still retain a + * deterministic order. + */ +export const GLOBAL_VAULT_KEY = '_global_'; + +export function resolveVaultId(event: EventOutbox): string { + if (event.aggregateType === 'VaultState') { + return event.aggregateId; + } + + try { + const payload = event.payload ? JSON.parse(event.payload) : null; + if (payload && typeof payload === 'object') { + const candidate = + (payload as Record).vaultId ?? + (payload as Record).aggregateId; + if (typeof candidate === 'string' && candidate.length > 0) { + return candidate; + } + } + } catch { + // Malformed payloads: fall back to the aggregate id below. + } + + return event.aggregateId || GLOBAL_VAULT_KEY:} + +/** + * Poller that delivers outbox events while preserving per-vault ordering. + * + * The original implementation ordered globally by `createdAt`, which allowed a + * slow webhook for vault A to block vault B's events and, worse, allowed a + * consumer to observe vault B's `version 5` before `version 4`. This poller + * instead groups pending events by vault and claims them with `FOR UPDATE + * SKIP NAMED` partitioned by vault, so each vault's events are delivered in + * `sequence` order while different vaults progress independently. + */ +export class OutboxPoller { + private readonly prsma: PrismaClient; + private readonly batchSize: number; + private readonly maxAttempts: number; + private readonly lockTimeoutMs: number; + private readonly handler: (event: EventOutbox) => Promise; + private readonly workerId: string; + + constructor(prisma: PrismaClient, options: OutboxPollerOptions) { + this.prisma = prisma; + this.batchSize = options.batchSize ?? 50; + this.maxAttempts = options.maxAttempts ?? 3; + this.lockTimeoutMs = options.lockTimeoutMs ?? 60_000; + this.handler = options.handler; + this.workerId = options.workerId ?? `outbox-poller-${Math.random().toString(36).slice(2, 10)}`; + } + + /** + * Run a single poll cycle. Returns a count of how many events were delivered, + * failed, and retried. + */ + async pollOnce(): Promise { + const candidates = await this.fetchPendingEvents(); + const result: OutboxPollerResult = { delivered: 0, failed: 0, retried: 0 }; + + for (const candidate of candidates) { + const claimed = await this.claimEvent(candidate); + if (!claimed) { + // Another worker won the race; skip. + continue; + } + + try { + await this.handler(claimed); + await this.markDelivered(claimed); + result.delivered += 1; + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + const nextAttempt = claimed.attemptCount + 1; + if (nextAttempt >= claimed.maxAttempts) { + await this.markFailed(claimed, message); + result.failed += 1; + } else { + await this.releaseForRetry(claimed, message); + result.retried += 1; + } + } + } + + return result; + } + + /** + * Fetch pending events grouped by vault, ordered by `sequence` within each + * vault. We fetch the oldest pending event per vault first, then any additional + * events for that vault in `sequence` order, so a slow vault cannot block + * another vault's progress. + */ + private async fetchPendingEvents(): Promise { + const now = new Date(); + const lockCutoff = new Date(now.getTime() - this.lockTimeoutMs); + + // Select the oldest pending event per vault. This ensures we always advance + // each vault from its oldest undelivered event and never jump ahead. + const heads = await this.prisma.$QueryRaw>( + Prisma.sql` + SELECT as "vaultId", MIN("createdAt") AS "minCreatedAt" + FROM "EventOutbox" + WHERE "status" = 'pending' + AND ("lockedAt" IS NULL OR "lockedAt" <= ${lockCutoff}) + GROUP BY + CASE + WHEN "aggregateType" = 'VaultState' THEN "aggregateId" + ELSE COALSCE("aggregateId", ':', "id") + END + `, + ); + + if (heads.length === 0) { + return []; + } + + const events: EventOutbox[] = []; + for (const head of heads) { + const vaultId = head.vaultId; + const minCreatedAt = head.minCreatedAt; + + const vaultEvents = await this.prisma.$QueryRaw( + Prisma.sql` + SELECT * FROM "EventOutbox" + WHERE "status" = 'pending' + AND ("lockedAt" IS NULL OR "lockedAt" <= ${lockCutoff}) + AND ( + CASE + WHEN "aggregateType" = 'VaultState' THEN "aggregateId" + ELSE COALESCE("aggregateId", ':', "id") + END + ) = ${vaultId} + AND "createdAt" >= ${minCreatedAt} + ORDER BY + CASE + WHEN "aggregateType" = 'VaultState' THEN "sequence" + ELSE 0 + END ASC, + "createdAt" ASC, + "id" ASC + LIMIT ${this.batchSize} + `, + ); + + events.push(...vaultEvents); + } + + return events; + } + + /** + * Claim an event using an atomic `UPDATE ... WHERE "lockedAt" IS NULL` + * guard. Returns the updated row if this worker won the claim, otherwise null. + */ + private async claimEvent(candidate: EventOutbox): Promise { + const now = new Date(); + const lockCutoff = new Date(now.getTime() - this.lockTimeoutMs); + + const claimed = await this.prisma.$QueryRaw( + Prisma.sql` + UPDATE "EventOutbox" + SET "status" = 'processing', + "lockedAt" = ${now}, + "lockedBy" = ${this.workerId}, + "attemptCount" = "attemptCount" + 1, + "updatedAt" = ${now} + WHERE "id" = ${candidate.id} + AND "status" = 'pending' + AND ("lockedAt" IS NULL OR "lockedAt" <= ${lockCutoff}) + RETURNING * + `, + ); + + return claimed.length > 0 ? claimed[0] : null; + } + + private async markDelivered(event: EventOutbox): Promise { + const now = new Date(); + await this.prisma.$executeRaw( + Prisma.sql` + UPDATE "EventOutbox" + SET "status" = 'delivered', + "relayedAt" = ${now}, + "lockedAt" = NULL, + "lockedBy" = NULL, + "updatedAt" = ${now} + WHERE "id" = ${event.id} + `, + ); + } + + private async markFailed(event: EventOutbox, errorMessage: string): Promise { + const now = new Date(); + await this.prisma.$executeRaw( + Prisma.sql` + UPDATE "EventOutbox" + SET "status" = 'failed', + "lastError" = ${errorMessage}, + "lockedAt" = NULL, + "lockedBy" = NULL, + "updatedAt" = ${now} + WHERE "id" = ${event.id} + `, + ); + } + + private async releaseForRetry(event: EventOutbox, errorMessage: string): Promise { + const now = new Date(); + await this.prisma.$executeRaw( + Prisma.sql` + UPDATE "EventOutbox" + SET "status" = 'pending', + "lastError" = ${errorMessage}, + "lockedAt" = NULL, + "lockedBy" = NULL, + "updatedAt" = ${now} + WHERE "id" = ${event.id} + `, + ); + } +} + +/** + * Enqueue a vault-scoped event with a monotonically increasing `sequence` + * per vault. The sequence is assigned inside the same transaction as the + * insert, so concurrent enqueues for the same vault cannot collide. + */ +export async function enqueueVaultEvent( + prisma: PrismaClient, + input: { + eventType: string; + payload: unknown; + vaultId: string; + aggregateType?: string; + maxAttempts?: number; + }, +): Promise { + const aggregateType = input.aggregateType ?? 'VaultState'; + const now = new Date(); + + return prisma.$transaction(async (tx) => { + const latest = await tx.$QueryRaw>( + Prisma.sql` + SELECT MAX("sequence") AS "sequence" + FROM "EventOutbox" + WHERE "aggregateType" = ${aggregateType} + AND "aggregateId" = ${input.vaultId} + `, + ); + + const nextSequence = (latest[0]?.sequence ?? 0) + 1; + + return tx.eventOutbox.create({ + data: { + eventType: input.eventType, + payload: JSON.stringify(input.payload), + aggregateType, + aggregateId: input.vaultId, + sequence: nextSequence, + maxAttempts: input.maxAttempts ?? 3, + updatedAt: now, + }, + }); + }); +}