diff --git a/backend/src/indexer/projections.ts b/backend/src/indexer/projections.ts index faf6d6e0..d2ad9b3c 100644 --- a/backend/src/indexer/projections.ts +++ b/backend/src/indexer/projections.ts @@ -1,4 +1,5 @@ import type { Prisma } from '@prisma/client'; +import type { PrismaClientKnownRequestError } from '@prisma/client/runtime/library'; import { prisma } from '../db/prisma.js'; import { logger } from '../common/utils/logger.js'; import type { DecodedEvent } from './sorobanClient.js'; @@ -73,15 +74,31 @@ async function persistEventLog(event: DecodedEvent): Promise { }); if (existing) return false; - await prisma.eventLog.create({ - data: { - topic: event.topic, - ledger: event.ledger, - txHash: event.txHash, - data: (event.value ?? {}) as Prisma.InputJsonValue, - }, - }); - return true; + try { + await prisma.eventLog.create({ + data: { + topic: event.topic, + ledger: event.ledger, + txHash: event.txHash, + data: (event.value ?? {}) as Prisma.InputJsonValue, + }, + }); + return true; + } catch (err) { + // Two workers may process the same ledger range concurrently. The unique + // (txHash, topic, ledger) constraint (schema) is the source of truth: a + // P2002 here means another worker already persisted this exact event, so + // this is a replay, not an error — treat it as already-seen and continue. + const error = err as PrismaClientKnownRequestError; + if (error?.code === 'P2002') { + logger.debug( + { txHash: event.txHash, topic: event.topic, ledger: event.ledger }, + 'Event insert raced with a duplicate; treating as replay', + ); + return false; + } + throw err; + } } /** Upsert the Tip row. txHash is unique, so replays are no-ops. */ diff --git a/backend/tests/idempotency.test.ts b/backend/tests/idempotency.test.ts new file mode 100644 index 00000000..c4485216 --- /dev/null +++ b/backend/tests/idempotency.test.ts @@ -0,0 +1,114 @@ +import { beforeEach, describe, expect, it } from 'vitest'; +import { prisma } from './helpers/db.js'; +import { resetDb } from './helpers/db.js'; +import { projectEvent } from '../src/indexer/projections.js'; +import { fullFixtureEventPage, ADDR_A } from '../src/indexer/fixtures/events.js'; + +beforeEach(() => resetDb()); + +/** Aggregate counts for every projection-relevant table, used to compare state. */ +async function snapshotState(): Promise> { + const [tips, refunds, users, goals, subscriptions, credits, creditHistory, eventLogs] = + await Promise.all([ + prisma.tip.count(), + prisma.refund.count(), + prisma.user.count(), + prisma.goal.count(), + prisma.subscription.count(), + prisma.creditScore.count(), + prisma.creditScoreHistory.count(), + prisma.eventLog.count(), + ]); + return { tips, refunds, users, goals, subscriptions, credits, creditHistory, eventLogs }; +} + +describe('indexer projection idempotency', () => { + it('processes every fixture twice and leaves identical state', async () => { + const events = fullFixtureEventPage.events; + expect(events.length).toBeGreaterThan(0); + + for (const event of events) { + await projectEvent(event); + } + const afterFirst = await snapshotState(); + + for (const event of events) { + await projectEvent(event); + } + const afterSecond = await snapshotState(); + + expect(afterSecond).toEqual(afterFirst); + }); + + it('every event type projects to exactly one row after double processing', async () => { + const events = fullFixtureEventPage.events; + + for (let pass = 0; pass < 2; pass++) { + for (const event of events) { + await projectEvent(event); + } + } + + // Each event type must appear exactly once in the event log. + for (const event of events) { + const logs = await prisma.eventLog.findMany({ + where: { txHash: event.txHash, topic: event.topic }, + }); + expect(logs).toHaveLength(1); + } + + // The single goal / subscription / credit rows for the fixtures stay singular. + const goalCount = await prisma.goal.count(); + expect(goalCount).toBe(await prisma.goal.findMany().then((r) => r.length)); + expect(await prisma.goal.count()).toBeGreaterThan(0); + }); + + it('out-of-order delivery (older ledger replayed after newer) produces no duplicate rows', async () => { + // Deliver in reversed order (out-of-order) to prove order independence. + const events = [...fullFixtureEventPage.events].reverse(); + + for (let pass = 0; pass < 2; pass++) { + for (const event of events) { + await projectEvent(event); + } + } + + for (const event of events) { + const logs = await prisma.eventLog.findMany({ + where: { txHash: event.txHash, topic: event.topic }, + }); + expect(logs).toHaveLength(1); + } + }); + + it('is robust under concurrent duplicate replay', async () => { + const events = fullFixtureEventPage.events; + await Promise.all(events.flatMap((event) => [projectEvent(event), projectEvent(event)])); + + for (const event of events) { + const logs = await prisma.eventLog.findMany({ + where: { txHash: event.txHash, topic: event.topic }, + }); + expect(logs).toHaveLength(1); + } + }); + + it('re-running projections over tip events is not a blind increment', async () => { + const tipEvent = fullFixtureEventPage.events.find((e) => e.topic === 'tip_sent'); + const profileEvent = fullFixtureEventPage.events.find((e) => e.topic === 'profile_register'); + expect(tipEvent).toBeDefined(); + expect(profileEvent).toBeDefined(); + + // Ensure the profile user row exists first (matching real ordering). + await projectEvent(profileEvent!); + + const before = await prisma.tip.count(); + // Deliver the same tip ledger tuple twice plus a concurrent one. + await projectEvent(tipEvent!); + await Promise.all([projectEvent(tipEvent!), projectEvent(tipEvent!)]); + const after = await prisma.tip.count(); + + expect(after).toBe(before + 1); + expect(ADDR_A).toBeDefined(); + }); +}); \ No newline at end of file