Skip to content

Commit db8a9f7

Browse files
authored
Merge pull request #382 from ShantelPeters/feat/ledger-poller-sse
feat: implement LedgerPoller SSE with polling fallback and dispatcher…… integration
2 parents a929c50 + d99962f commit db8a9f7

6 files changed

Lines changed: 359 additions & 104 deletions

File tree

indexer/src/health/health.controller.spec.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ describe('HealthController', () => {
2828
lag_ledgers: 5,
2929
db: 'ok',
3030
redis: 'ok',
31+
redis_latency_ms: 0,
3132
};
3233
healthService.getHealth.mockResolvedValue(okResult);
3334

@@ -41,6 +42,7 @@ describe('HealthController', () => {
4142
lag_ledgers: 250,
4243
db: 'ok',
4344
redis: 'ok',
45+
redis_latency_ms: 0,
4446
};
4547
healthService.getHealth.mockResolvedValue(degradedResult);
4648

@@ -55,6 +57,7 @@ describe('HealthController', () => {
5557
lag_ledgers: 150,
5658
db: 'error',
5759
redis: 'ok',
60+
redis_latency_ms: 0,
5861
};
5962
healthService.getHealth.mockResolvedValue(degradedResult);
6063

@@ -75,6 +78,7 @@ describe('HealthController', () => {
7578
lag_ledgers: 0,
7679
db: 'ok',
7780
redis: 'ok',
81+
redis_latency_ms: 0,
7882
};
7983
healthService.getHealth.mockResolvedValue(okResult);
8084

indexer/src/ingestor/event-handler-registry.service.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,16 @@ export class EventHandlerRegistry implements OnModuleInit {
8080
eventHandlers: {
8181
RaffleCreated: "RaffleCreatedHandler",
8282
TicketPurchased: "TicketPurchasedHandler",
83+
DrawTriggered: "DrawTriggeredHandler",
84+
RandomnessRequested: "RandomnessRequestedHandler",
85+
RandomnessReceived: "RandomnessReceivedHandler",
8386
RaffleFinalized: "RaffleFinalizedHandler",
87+
RaffleCancelled: "RaffleCancelledHandler",
88+
TicketRefunded: "TicketRefundedHandler",
89+
ContractPaused: "ContractPausedHandler",
90+
ContractUnpaused: "ContractUnpausedHandler",
91+
AdminTransferProposed: "AdminTransferProposedHandler",
92+
AdminTransferAccepted: "AdminTransferAcceptedHandler",
8493
},
8594
},
8695
],
Lines changed: 135 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,135 @@
1+
import { Injectable, Logger } from "@nestjs/common";
2+
import { RaffleProcessor } from "../processors/raffle.processor";
3+
import { TicketProcessor } from "../processors/ticket.processor";
4+
import { UserProcessor } from "../processors/user.processor";
5+
import { AdminProcessor } from "../processors/admin.processor";
6+
import { DomainEvent } from "./event.types";
7+
8+
/**
9+
* IngestionDispatcherService
10+
*
11+
* Responsible for orchestrating the processing of parsed DomainEvents.
12+
* It maps event types to the appropriate methods in the various processors.
13+
*/
14+
@Injectable()
15+
export class IngestionDispatcherService {
16+
private readonly logger = new Logger(IngestionDispatcherService.name);
17+
18+
constructor(
19+
private readonly raffleProcessor: RaffleProcessor,
20+
private readonly ticketProcessor: TicketProcessor,
21+
private readonly userProcessor: UserProcessor,
22+
private readonly adminProcessor: AdminProcessor,
23+
) {}
24+
25+
/**
26+
* Dispatches a parsed domain event to the appropriate processor.
27+
*
28+
* @param event The parsed DomainEvent
29+
* @param rawEvent The original raw event from Horizon (containing ledger, txHash, etc.)
30+
*/
31+
async dispatch(event: DomainEvent, rawEvent: any): Promise<void> {
32+
const ledger = Number(rawEvent.ledger);
33+
const txHash = rawEvent.id || rawEvent.paging_token;
34+
35+
this.logger.debug(`Dispatching event: ${event.type} from ledger ${ledger}`);
36+
37+
try {
38+
switch (event.type) {
39+
case "RaffleCreated":
40+
await this.raffleProcessor.handleRaffleCreated(
41+
event.raffle_id,
42+
event.creator,
43+
ledger,
44+
);
45+
break;
46+
47+
case "TicketPurchased":
48+
await this.ticketProcessor.handleTicketPurchased(
49+
event.raffle_id,
50+
event.buyer,
51+
event.ticket_ids,
52+
event.total_paid,
53+
ledger,
54+
txHash,
55+
);
56+
break;
57+
58+
case "RaffleFinalized":
59+
await this.raffleProcessor.handleRaffleFinalized(
60+
event.raffle_id,
61+
event.winner,
62+
event.prize_amount,
63+
);
64+
break;
65+
66+
case "RaffleCancelled":
67+
await this.raffleProcessor.handleRaffleCancelled(
68+
event.raffle_id,
69+
event.reason,
70+
ledger,
71+
txHash,
72+
);
73+
break;
74+
75+
case "TicketRefunded":
76+
await this.ticketProcessor.handleTicketRefunded(
77+
event.raffle_id,
78+
event.ticket_id,
79+
event.recipient,
80+
event.amount,
81+
txHash,
82+
);
83+
break;
84+
85+
case "ContractPaused":
86+
await this.adminProcessor.handleContractPaused(event.admin, ledger, txHash);
87+
break;
88+
89+
case "ContractUnpaused":
90+
await this.adminProcessor.handleContractUnpaused(event.admin, ledger, txHash);
91+
break;
92+
93+
case "AdminTransferProposed":
94+
await this.adminProcessor.handleAdminTransferProposed(
95+
event.current_admin,
96+
event.proposed_admin,
97+
ledger,
98+
txHash,
99+
);
100+
break;
101+
102+
case "AdminTransferAccepted":
103+
await this.adminProcessor.handleAdminTransferAccepted(
104+
event.old_admin,
105+
event.new_admin,
106+
ledger,
107+
txHash,
108+
);
109+
break;
110+
111+
case "DrawTriggered":
112+
this.logger.log(`DrawTriggered for raffle ${event.raffle_id} at ledger ${event.ledger}`);
113+
// Currently, this might just be logged or used for internal state tracking
114+
break;
115+
116+
case "RandomnessRequested":
117+
this.logger.log(`RandomnessRequested for raffle ${event.raffle_id}, request ID ${event.request_id}`);
118+
break;
119+
120+
case "RandomnessReceived":
121+
this.logger.log(`RandomnessReceived for raffle ${event.raffle_id}`);
122+
break;
123+
124+
default:
125+
this.logger.warn(`No processor method found for event type: ${(event as any).type}`);
126+
}
127+
} catch (error: any) {
128+
this.logger.error(
129+
`Failed to dispatch event ${event.type} for tx ${txHash}: ${error.message}`,
130+
error.stack,
131+
);
132+
throw error; // Re-throw to allow LedgerPoller to handle retry/backoff
133+
}
134+
}
135+
}

indexer/src/ingestor/ingestor.module.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,20 +4,24 @@ import { EventParserService } from "./event-parser.service";
44
import { LedgerPollerService } from "./ledger-poller.service";
55
import { EventHandlersModule } from "./event-handlers.module";
66
import { DryRunService } from "./dry-run.service";
7+
import { IngestionDispatcherService } from "./ingestion-dispatcher.service";
8+
import { ProcessorsModule } from "../processors/processors.module";
79

810
@Module({
9-
imports: [EventHandlersModule],
11+
imports: [EventHandlersModule, ProcessorsModule],
1012
providers: [
1113
CursorManagerService,
1214
EventParserService,
1315
LedgerPollerService,
1416
DryRunService,
17+
IngestionDispatcherService,
1518
],
1619
exports: [
1720
CursorManagerService,
1821
EventParserService,
1922
LedgerPollerService,
2023
DryRunService,
24+
IngestionDispatcherService,
2125
EventHandlersModule,
2226
],
2327
})

0 commit comments

Comments
 (0)