From cc580d1e528fe65bd973aa0343c96a77eb4d2f39 Mon Sep 17 00:00:00 2001 From: Toan Date: Tue, 18 Aug 2026 10:17:45 +0700 Subject: [PATCH] fix(indexer): wrap event handlers in Prisma transaction to prevent poison-pill events (#168) --- .../parsers/event-parser.service.spec.ts | 135 ++++++++- .../indexer/parsers/event-parser.service.ts | 282 +++++++++++------- 2 files changed, 292 insertions(+), 125 deletions(-) diff --git a/server/src/indexer/parsers/event-parser.service.spec.ts b/server/src/indexer/parsers/event-parser.service.spec.ts index 53ac4c5..b5001ae 100644 --- a/server/src/indexer/parsers/event-parser.service.spec.ts +++ b/server/src/indexer/parsers/event-parser.service.spec.ts @@ -3,15 +3,20 @@ import { EventParserService } from './event-parser.service'; import type { RawSorobanEvent } from '../types/soroban-events.types'; type MockPrisma = { + $transaction: jest.Mock; transaction: { findUnique: jest.Mock; create: jest.Mock }; - campaign: { upsert: jest.Mock; update: jest.Mock }; + campaign: { + upsert: jest.Mock; + update: jest.Mock; + updateMany: jest.Mock; + }; user: { upsert: jest.Mock }; investment: { create: jest.Mock }; tranche: { create: jest.Mock }; dispute: { create: jest.Mock; update: jest.Mock; findFirst: jest.Mock }; }; -function makeMockPrisma(): MockPrisma { +function makeModelMocks() { return { transaction: { findUnique: jest.fn().mockResolvedValue(null), @@ -20,6 +25,7 @@ function makeMockPrisma(): MockPrisma { campaign: { upsert: jest.fn().mockResolvedValue({}), update: jest.fn().mockResolvedValue({}), + updateMany: jest.fn().mockResolvedValue({ count: 1 }), }, user: { upsert: jest.fn().mockResolvedValue({}), @@ -38,6 +44,18 @@ function makeMockPrisma(): MockPrisma { }; } +function makeMockPrisma(): MockPrisma { + const models = makeModelMocks(); + return { + $transaction: jest.fn(async (callback: (tx: any) => Promise) => { + // Pass the same model mocks as the transaction client so all calls + // within the transaction are recorded on the same mocks. + await callback(models); + }), + ...models, + }; +} + function rawEvent( id: string, topic: unknown[], @@ -689,9 +707,12 @@ describe('EventParserService', () => { describe('broadcast-after-persist wiring', () => { it('emits a realtime event only after the DB write succeeds', async () => { const emitCampaignEvent = jest.fn(); - const withRealtime = new EventParserService(prisma as any, { - emitCampaignEvent, - } as any); + const withRealtime = new EventParserService( + prisma as any, + { + emitCampaignEvent, + } as any, + ); await withRealtime.processEvent( rawEvent( @@ -713,9 +734,12 @@ describe('EventParserService', () => { it('does not emit for already-persisted (replayed) events', async () => { const emitCampaignEvent = jest.fn(); prisma.transaction.findUnique.mockResolvedValueOnce({ id: 'e-rt2' }); - const withRealtime = new EventParserService(prisma as any, { - emitCampaignEvent, - } as any); + const withRealtime = new EventParserService( + prisma as any, + { + emitCampaignEvent, + } as any, + ); await withRealtime.processEvent( rawEvent( @@ -730,9 +754,12 @@ describe('EventParserService', () => { it('does not emit when the parsed event has no campaignId', async () => { const emitCampaignEvent = jest.fn(); - const withRealtime = new EventParserService(prisma as any, { - emitCampaignEvent, - } as any); + const withRealtime = new EventParserService( + prisma as any, + { + emitCampaignEvent, + } as any, + ); await withRealtime.processEvent( rawEvent( @@ -745,4 +772,90 @@ describe('EventParserService', () => { expect(emitCampaignEvent).not.toHaveBeenCalled(); }); }); + + describe('transaction atomicity (poison-pill prevention)', () => { + it('rolls back all writes when a handler fails mid-way and allows retry', async () => { + // First call: campaign.update fails after investment.create would have succeeded. + // The transaction should roll back, leaving zero rows. + prisma.campaign.update.mockRejectedValueOnce( + new Error('transient DB error'), + ); + + // processEvent catches and logs the error (does not rethrow). + await service.processEvent( + rawEvent( + 'e-tx1', + ['ContribReceived', CAMPAIGN_ID], + [INVESTOR, 1700000000n, 250n], + ), + ); + + // Error should be logged. + expect(errorSpy).toHaveBeenCalled(); + + // KEY ASSERTION: Transaction row (idempotency marker) NOT created on failure, + // so the event can be retried on next poll. + expect(prisma.transaction.create).not.toHaveBeenCalled(); + + // Second call: same event, transient error resolved -> succeeds. + prisma.campaign.update.mockResolvedValueOnce({}); + errorSpy.mockClear(); + + await service.processEvent( + rawEvent( + 'e-tx1', // same event ID - idempotency key + ['ContribReceived', CAMPAIGN_ID], + [INVESTOR, 1700000000n, 250n], + ), + ); + + // Now the full transaction commits: Transaction row created (event marked processed). + expect(prisma.transaction.create).toHaveBeenCalledTimes(1); + expect(prisma.transaction.create).toHaveBeenCalledWith( + expect.objectContaining({ + data: expect.objectContaining({ + id: 'e-tx1', + type: 'campaign.invested', + }), + }), + ); + // No error on retry. + expect(errorSpy).not.toHaveBeenCalled(); + }); + + it('marks event as processed only after full transaction commits', async () => { + // Simulate a handler that fails on the second write. + prisma.campaign.update.mockRejectedValueOnce( + new Error('constraint violation'), + ); + + await service.processEvent( + rawEvent( + 'e-tx2', + ['CampaignFunded', CAMPAIGN_ID], + [1700000000n, 5000n], + ), + ); + + expect(errorSpy).toHaveBeenCalled(); + errorSpy.mockClear(); + + // Event NOT marked processed (no Transaction row) -> can be retried. + expect(prisma.transaction.create).not.toHaveBeenCalled(); + + // Retry succeeds. + prisma.campaign.update.mockResolvedValueOnce({}); + await service.processEvent( + rawEvent( + 'e-tx2', + ['CampaignFunded', CAMPAIGN_ID], + [1700000000n, 5000n], + ), + ); + + // Now marked processed exactly once. + expect(prisma.transaction.create).toHaveBeenCalledTimes(1); + expect(errorSpy).not.toHaveBeenCalled(); + }); + }); }); diff --git a/server/src/indexer/parsers/event-parser.service.ts b/server/src/indexer/parsers/event-parser.service.ts index 3f30f7a..4fea8bc 100644 --- a/server/src/indexer/parsers/event-parser.service.ts +++ b/server/src/indexer/parsers/event-parser.service.ts @@ -834,89 +834,97 @@ export class EventParserService { }); if (existing) return false; - switch (parsed.type) { - case 'campaign.escrow_created': - await this.handleCampaignEscrowCreated(parsed); - break; - case 'campaign.invested': - await this.handleCampaignInvested(parsed); - break; - case 'campaign.funded': - await this.handleCampaignFunded(parsed); - break; - case 'campaign.tranches_configured': - await this.handleTranchesConfigured(parsed); - break; - case 'campaign.tranche_released': - await this.handleTrancheReleased(parsed); - break; - case 'campaign.harvest_reported': - await this.handleHarvestReported(parsed); - break; - case 'campaign.failed': - await this.handleCampaignFailed(parsed); - break; - case 'campaign.return_claimed': - case 'campaign.refund_claimed': - await this.handleClaim(parsed); - break; - case 'campaign.dispute_opened': - await this.handleDisputeOpened(parsed); - break; - case 'campaign.dispute_resolved': - await this.handleDisputeResolved(parsed); - break; - case 'campaign.settled': - await this.handleCampaignSettled(parsed); - break; - case 'campaign.created': - await this.handleCampaignCreated(parsed); - break; - case 'campaign.escrow_linked': - await this.handleCampaignEscrowLinked(parsed); - break; - case 'campaign.status_updated': - await this.handleCampaignStatusUpdated(parsed); - break; - case 'registry.farmer_registered': - await this.handleFarmerRegistered(parsed); - break; - case 'registry.admin_initialized': - case 'registry.admin_updated': - case 'registry.contract_approved': - case 'registry.contract_revoked': - case 'registry.activity_recorded': - // Administrative/audit-only events: no domain row needs to change. - // Still captured below via the generic Transaction audit row. - break; - default: - this.logger.warn( - { type: parsed.type }, - 'No persistence handler for parsed event type', - ); - } + // Wrap handler writes + Transaction row insert in a single transaction + // so a partial failure rolls back cleanly and the event can be retried. + await this.prisma.$transaction(async (tx) => { + switch (parsed.type) { + case 'campaign.escrow_created': + await this.handleCampaignEscrowCreated(tx, parsed); + break; + case 'campaign.invested': + await this.handleCampaignInvested(tx, parsed); + break; + case 'campaign.funded': + await this.handleCampaignFunded(tx, parsed); + break; + case 'campaign.tranches_configured': + await this.handleTranchesConfigured(tx, parsed); + break; + case 'campaign.tranche_released': + await this.handleTrancheReleased(tx, parsed); + break; + case 'campaign.harvest_reported': + await this.handleHarvestReported(tx, parsed); + break; + case 'campaign.failed': + await this.handleCampaignFailed(tx, parsed); + break; + case 'campaign.return_claimed': + case 'campaign.refund_claimed': + await this.handleClaim(tx, parsed); + break; + case 'campaign.dispute_opened': + await this.handleDisputeOpened(tx, parsed); + break; + case 'campaign.dispute_resolved': + await this.handleDisputeResolved(tx, parsed); + break; + case 'campaign.settled': + await this.handleCampaignSettled(tx, parsed); + break; + case 'campaign.created': + await this.handleCampaignCreated(tx, parsed); + break; + case 'campaign.escrow_linked': + await this.handleCampaignEscrowLinked(tx, parsed); + break; + case 'campaign.status_updated': + await this.handleCampaignStatusUpdated(tx, parsed); + break; + case 'registry.farmer_registered': + await this.handleFarmerRegistered(tx, parsed); + break; + case 'registry.admin_initialized': + case 'registry.admin_updated': + case 'registry.contract_approved': + case 'registry.contract_revoked': + case 'registry.activity_recorded': + // Administrative/audit-only events: no domain row needs to change. + // Still captured below via the generic Transaction audit row. + break; + default: + this.logger.warn( + { type: parsed.type }, + 'No persistence handler for parsed event type', + ); + } - await this.prisma.transaction.create({ - data: { - id: parsed.id, - type: parsed.type, - campaignId: parsed.campaignId, - userId: parsed.userAddress, - amount: parsed.amount, - txHash: parsed.txHash, - status: 'Confirmed', - timestamp: BigInt(parsed.ledger), - data: JSON.stringify(parsed, (_key, value) => - typeof value === 'bigint' ? value.toString() : value, - ), - }, + await tx.transaction.create({ + data: { + id: parsed.id, + type: parsed.type, + campaignId: parsed.campaignId, + userId: parsed.userAddress, + amount: parsed.amount, + txHash: parsed.txHash, + status: 'Confirmed', + timestamp: BigInt(parsed.ledger), + data: JSON.stringify(parsed, (_key, value) => + typeof value === 'bigint' ? value.toString() : value, + ), + }, + }); }); return true; } - private async ensureUser(address: string, timestamp: number): Promise { - await this.prisma.user.upsert({ + private async ensureUser( + tx: any, + address: string, + timestamp: number, + ): Promise { + await tx.user.upsert({ where: { address }, update: {}, create: { address, firstSeenAt: BigInt(timestamp) }, @@ -924,12 +932,13 @@ export class EventParserService { } private async handleCampaignEscrowCreated( + tx: any, parsed: ParsedEvent, ): Promise { const data = parsed.data as unknown as CampaignEscrowCreatedData; - await this.ensureUser(data.farmer, data.timestamp); + await this.ensureUser(tx, data.farmer, data.timestamp); - await this.prisma.campaign.upsert({ + await tx.campaign.upsert({ where: { id: data.campaignId }, update: { farmer: data.farmer, targetAmount: BigInt(data.targetAmount) }, create: { @@ -946,15 +955,18 @@ export class EventParserService { ); } - private async handleCampaignInvested(parsed: ParsedEvent): Promise { + private async handleCampaignInvested( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignInvestedData; const campaignId = data.campaignId; - await this.ensureUser(data.investor, data.timestamp); + await this.ensureUser(tx, data.investor, data.timestamp); const amount = BigInt(data.amount); - await this.prisma.investment.create({ + await tx.investment.create({ data: { id: parsed.id, campaignId, @@ -965,7 +977,7 @@ export class EventParserService { }, }); - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: campaignId }, data: { totalFunded: { increment: amount } }, }); @@ -976,18 +988,24 @@ export class EventParserService { ); } - private async handleCampaignFunded(parsed: ParsedEvent): Promise { + private async handleCampaignFunded( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignFundedData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Funded', totalFunded: BigInt(data.totalFunded) }, }); this.logger.log({ campaignId: data.campaignId }, 'Campaign funded'); } - private async handleTranchesConfigured(parsed: ParsedEvent): Promise { + private async handleTranchesConfigured( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as TranchesConfiguredData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { trancheCount: data.trancheCount }, }); @@ -997,10 +1015,13 @@ export class EventParserService { ); } - private async handleTrancheReleased(parsed: ParsedEvent): Promise { + private async handleTrancheReleased( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as TrancheReleasedData; - await this.ensureUser(data.recipient, data.timestamp); - await this.prisma.tranche.create({ + await this.ensureUser(tx, data.recipient, data.timestamp); + await tx.tranche.create({ data: { id: parsed.id, campaignId: data.campaignId, @@ -1010,15 +1031,26 @@ export class EventParserService { txHash: parsed.txHash, }, }); + // The first tranche release transitions the campaign from Funded into + // InProduction (mirroring the escrow contract's own state machine). The + // conditional `where` (id + current status) makes repeated releases on an + // already-InProduction campaign a no-op instead of an unnecessary write. + await tx.campaign.updateMany({ + where: { id: data.campaignId, status: 'Funded' }, + data: { status: 'InProduction' }, + }); this.logger.log( { campaignId: data.campaignId, recipient: data.recipient }, 'Tranche released', ); } - private async handleHarvestReported(parsed: ParsedEvent): Promise { + private async handleHarvestReported( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as HarvestReportedData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Harvested', @@ -1032,16 +1064,19 @@ export class EventParserService { ); } - private async handleCampaignFailed(parsed: ParsedEvent): Promise { + private async handleCampaignFailed( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignFailedData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Failed', refundable: BigInt(data.refundable) }, }); this.logger.log({ campaignId: data.campaignId }, 'Campaign failed'); } - private async handleClaim(parsed: ParsedEvent): Promise { + private async handleClaim(tx: any, parsed: ParsedEvent): Promise { // ReturnClaimed/RefundClaimed report a payout to an investor after // settlement/failure. A claim aggregates an investor's whole position // rather than a single contribution, so it isn't attributed to a @@ -1051,13 +1086,16 @@ export class EventParserService { const data = parsed.data as unknown as | ReturnClaimedData | RefundClaimedData; - await this.ensureUser(data.investor, data.timestamp); + await this.ensureUser(tx, data.investor, data.timestamp); } - private async handleDisputeOpened(parsed: ParsedEvent): Promise { + private async handleDisputeOpened( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as DisputeOpenedData; - await this.ensureUser(data.opener, data.timestamp); - await this.prisma.dispute.create({ + await this.ensureUser(tx, data.opener, data.timestamp); + await tx.dispute.create({ data: { id: parsed.id, campaignId: data.campaignId, @@ -1068,7 +1106,7 @@ export class EventParserService { ledgerSequence: data.ledgerSequence, }, }); - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Disputed' }, }); @@ -1078,16 +1116,19 @@ export class EventParserService { ); } - private async handleDisputeResolved(parsed: ParsedEvent): Promise { + private async handleDisputeResolved( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as DisputeResolvedData; - const openDispute = await this.prisma.dispute.findFirst({ + const openDispute = await tx.dispute.findFirst({ where: { campaignId: data.campaignId, status: 'Open' }, orderBy: { openedAt: 'desc' }, }); if (openDispute) { - await this.prisma.dispute.update({ + await tx.dispute.update({ where: { id: openDispute.id }, data: { status: 'Resolved', @@ -1105,7 +1146,7 @@ export class EventParserService { ); } - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Resolved' }, }); @@ -1115,9 +1156,12 @@ export class EventParserService { ); } - private async handleCampaignSettled(parsed: ParsedEvent): Promise { + private async handleCampaignSettled( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignSettledData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: 'Settled' }, }); @@ -1127,11 +1171,14 @@ export class EventParserService { ); } - private async handleCampaignCreated(parsed: ParsedEvent): Promise { + private async handleCampaignCreated( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignCreatedData; - await this.ensureUser(data.farmer, data.timestamp); + await this.ensureUser(tx, data.farmer, data.timestamp); - await this.prisma.campaign.upsert({ + await tx.campaign.upsert({ where: { id: data.campaignId }, update: data.title !== undefined @@ -1151,9 +1198,12 @@ export class EventParserService { ); } - private async handleCampaignEscrowLinked(parsed: ParsedEvent): Promise { + private async handleCampaignEscrowLinked( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as CampaignEscrowLinkedData; - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { escrowContract: data.escrowContract, farmer: data.farmer }, }); @@ -1164,6 +1214,7 @@ export class EventParserService { } private async handleCampaignStatusUpdated( + tx: any, parsed: ParsedEvent, ): Promise { const data = parsed.data as unknown as CampaignStatusUpdatedData; @@ -1171,7 +1222,7 @@ export class EventParserService { // Activity-log mirror shape - no concrete status value to apply. return; } - await this.prisma.campaign.update({ + await tx.campaign.update({ where: { id: data.campaignId }, data: { status: data.newStatus }, }); @@ -1181,9 +1232,12 @@ export class EventParserService { ); } - private async handleFarmerRegistered(parsed: ParsedEvent): Promise { + private async handleFarmerRegistered( + tx: any, + parsed: ParsedEvent, + ): Promise { const data = parsed.data as unknown as FarmerRegisteredData; - await this.prisma.user.upsert({ + await tx.user.upsert({ where: { address: data.farmer }, update: data.name !== undefined ? { name: data.name } : {}, create: {