Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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");
Expand Down Expand Up @@ -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;
2 changes: 2 additions & 0 deletions backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand All @@ -610,6 +611,7 @@ model EventOutbox {
@@index([status])
@@index([status, createdAt])
@@index([aggregateType, aggregateId])
@@index([aggregateType, aggregateId, sequence])
@@index([lockedAt])
@@index([createdAt])
}
Expand Down
107 changes: 107 additions & 0 deletions backend/src/__tests__/eventOutbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ function makeOutboxInput(overrides: Partial<OutboxWriteInput> = {}): OutboxWrite
aggregateType: 'transaction',
aggregateId: 'tx-test-001',
maxAttempts: 3,
vaultId: 'vault-test-001',
...overrides,
};
}
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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', () => {
Expand Down
67 changes: 65 additions & 2 deletions backend/src/eventOutbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ export interface EventOutboxRecord {
status: OutboxEventStatus;
aggregateType: OutboxAggregateType;
aggregateId: string;
sequence: number;
attemptCount: number;
maxAttempts: number;
lastError: string | null;
Expand All @@ -77,6 +78,7 @@ export interface OutboxWriteInput {
payload: TransactionEventPayload;
aggregateType: OutboxAggregateType;
aggregateId: string;
vaultId?: string;
maxAttempts?: number;
}

Expand Down Expand Up @@ -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()}`,
Expand All @@ -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,
Expand Down Expand Up @@ -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: {
Expand All @@ -219,15 +232,40 @@ class EventOutboxService {
{ lockedAt: { lt: lockExpiry } },
],
},
orderBy: { createdAt: 'asc' },
orderBy: [
{ vaultId: 'asc' },
{ sequence: 'asc' },
{ createdAt: 'asc' },
],
take: limit,
});

if (candidates.length === 0) {
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<string, typeof candidates[number]>();
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: {
Expand All @@ -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) {
Expand Down Expand Up @@ -299,6 +341,7 @@ class EventOutboxService {
eventType: entry.eventType,
aggregateType: entry.aggregateType,
aggregateId: entry.aggregateId,
sequence: entry.sequence,
deliveredCount,
});
} catch (error) {
Expand Down Expand Up @@ -659,6 +702,7 @@ class EventOutboxService {
status: string;
aggregateType: string;
aggregateId: string;
sequence: number;
attemptCount: number;
maxAttempts: number;
lastError: string | null;
Expand All @@ -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,
Expand All @@ -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<number> {
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 ────────────────────────────────────────────────────────────────
Expand Down
Loading