Skip to content

Commit 5dd58cb

Browse files
authored
Merge pull request #136 from presidojay1/fix/event-poller-truncation-cursor-115
fix: don't advance indexer cursor past a truncated event batch
2 parents 0011608 + ea4c0cf commit 5dd58cb

2 files changed

Lines changed: 165 additions & 8 deletions

File tree

src/indexer/eventPoller.js

Lines changed: 25 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -59,20 +59,37 @@ class EventPoller {
5959
limit: this.pollLimit,
6060
});
6161

62-
const parsedEvents = (response.events || [])
63-
.map(parseContractEvent)
64-
.filter(Boolean);
62+
const rawEvents = response.events || [];
63+
const parsedEvents = rawEvents.map(parseContractEvent).filter(Boolean);
6564

6665
for (const event of parsedEvents) {
6766
await this.store.saveEvent(event);
6867
}
6968

70-
const latestIndexedLedger = Math.max(
71-
response.latestLedger || previousLedger,
72-
...parsedEvents.map((event) => event.ledger)
73-
);
69+
// response.latestLedger is the chain's current tip, not how far this
70+
// particular call actually got — if the RPC returned a full pollLimit
71+
// batch, more matching events may exist beyond it. Jumping the cursor
72+
// to the tip in that case permanently skips whatever wasn't returned,
73+
// with no retry and no downstream reconciliation that could recover
74+
// it (#115). Only safe to advance to the tip when this batch wasn't
75+
// truncated; otherwise advance only past what was actually processed,
76+
// so the next poll picks up right where this one left off.
77+
const truncated = rawEvents.length >= this.pollLimit;
78+
const eventLedgers = parsedEvents.map((event) => event.ledger);
79+
const latestIndexedLedger = truncated
80+
? Math.max(previousLedger ?? 0, ...eventLedgers)
81+
: Math.max(response.latestLedger || previousLedger || 0, ...eventLedgers);
82+
7483
await this.store.setLastLedger(latestIndexedLedger);
7584

85+
if (truncated) {
86+
this.logger.warn('SmartDrop event poll truncated by pollLimit; more events pending next cycle', {
87+
pollLimit: this.pollLimit,
88+
indexed_events: parsedEvents.length,
89+
resumed_from_ledger: latestIndexedLedger + 1,
90+
});
91+
}
92+
7693
this.latestLedger = response.latestLedger || null;
7794
this.lastRun = new Date().toISOString();
7895
this.lastError = null;
@@ -82,6 +99,7 @@ class EventPoller {
8299
start_ledger: startLedger,
83100
latest_ledger: response.latestLedger,
84101
indexed_events: parsedEvents.length,
102+
truncated,
85103
};
86104
}
87105

test/eventPoller.test.js

Lines changed: 140 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
const { nativeToScVal } = require('stellar-sdk');
44
const { EventPoller } = require('../src/indexer/eventPoller');
55

6-
function contractEvent() {
6+
function contractEvent(overrides = {}) {
77
return {
88
id: 'evt-1',
99
type: 'contract',
@@ -19,9 +19,30 @@ function contractEvent() {
1919
total_amount: 1000n,
2020
expiry_ledger: 500n,
2121
}),
22+
...overrides,
2223
};
2324
}
2425

26+
// N distinct events at consecutive ledgers starting at `startLedger`, used
27+
// to simulate a burst/backlog large enough to fill a batch of size `n`.
28+
function contractEvents(n, startLedger) {
29+
return Array.from({ length: n }, (_, i) => {
30+
const ledger = startLedger + i;
31+
return contractEvent({
32+
id: `evt-${ledger}`,
33+
ledger,
34+
pagingToken: `${ledger}-1`,
35+
value: nativeToScVal({
36+
airdrop_id: `drop-${ledger}`,
37+
creator: 'GCREATOR',
38+
token: 'USDC',
39+
total_amount: 1000n,
40+
expiry_ledger: 500n,
41+
}),
42+
});
43+
});
44+
}
45+
2546
describe('EventPoller', () => {
2647
test('polls Soroban RPC, stores parsed events, and advances last ledger', async () => {
2748
const server = {
@@ -95,4 +116,122 @@ describe('EventPoller', () => {
95116

96117
await expect(poller.pollOnce()).resolves.toMatchObject({ skipped: true });
97118
});
119+
120+
describe('truncated batch (#115)', () => {
121+
test('advances last_ledger only to the last processed event, not to the chain tip', async () => {
122+
const pollLimit = 5;
123+
// Simulates a real burst/backlog: exactly pollLimit events returned,
124+
// last event's ledger (24) is far behind the chain tip (500).
125+
const events = contractEvents(pollLimit, 20);
126+
const server = {
127+
getEvents: jest.fn(async () => ({ latestLedger: 500, events })),
128+
};
129+
const store = {
130+
getLastLedger: jest.fn(async () => null),
131+
saveEvent: jest.fn(async () => {}),
132+
setLastLedger: jest.fn(async () => {}),
133+
};
134+
const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() };
135+
136+
const poller = new EventPoller({
137+
enabled: true,
138+
contractId: 'CCONTRACT',
139+
startLedger: 10,
140+
pollLimit,
141+
server,
142+
store,
143+
logger,
144+
});
145+
146+
const result = await poller.pollOnce();
147+
148+
// Not 500 (response.latestLedger) — that would permanently skip
149+
// whatever exists between ledger 24 and the tip.
150+
expect(store.setLastLedger).toHaveBeenCalledWith(24);
151+
expect(result).toMatchObject({ truncated: true, indexed_events: pollLimit });
152+
expect(logger.warn).toHaveBeenCalledWith(
153+
'SmartDrop event poll truncated by pollLimit; more events pending next cycle',
154+
expect.objectContaining({ pollLimit, indexed_events: pollLimit, resumed_from_ledger: 25 }),
155+
);
156+
});
157+
158+
test('the next poll resumes from the last processed event, picking up previously-skippable events', async () => {
159+
const pollLimit = 5;
160+
const firstBatch = contractEvents(pollLimit, 20); // ledgers 20-24
161+
// Events that would have been silently skipped pre-fix: they sit
162+
// between the last processed ledger (24) and the previous poll's
163+
// chain-tip snapshot (500).
164+
const skippableRangeEvents = contractEvents(2, 100); // ledgers 100-101
165+
166+
const server = { getEvents: jest.fn() };
167+
server.getEvents
168+
.mockImplementationOnce(async () => ({ latestLedger: 500, events: firstBatch }))
169+
.mockImplementationOnce(async () => ({ latestLedger: 500, events: skippableRangeEvents }));
170+
171+
// Stateful store, so the second pollOnce() actually reads back what
172+
// the first one wrote — proves resumption across ticks, not just
173+
// within one call.
174+
let lastLedger = null;
175+
const store = {
176+
getLastLedger: jest.fn(async () => lastLedger),
177+
saveEvent: jest.fn(async () => {}),
178+
setLastLedger: jest.fn(async (ledger) => {
179+
lastLedger = ledger;
180+
}),
181+
};
182+
183+
const poller = new EventPoller({
184+
enabled: true,
185+
contractId: 'CCONTRACT',
186+
startLedger: 10,
187+
pollLimit,
188+
server,
189+
store,
190+
logger: { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() },
191+
});
192+
193+
const first = await poller.pollOnce();
194+
expect(first.truncated).toBe(true);
195+
expect(lastLedger).toBe(24);
196+
197+
const second = await poller.pollOnce();
198+
199+
// Resumed from 25 (last processed + 1), not 501 (tip + 1) — the
200+
// pre-fix bug would have started here at 501, skipping ledgers
201+
// 25-499 (including skippableRangeEvents) forever.
202+
expect(server.getEvents.mock.calls[1][0].startLedger).toBe(25);
203+
expect(second.indexed_events).toBe(2);
204+
expect(store.saveEvent).toHaveBeenCalledTimes(pollLimit + 2);
205+
});
206+
207+
test('a batch smaller than pollLimit still advances to the chain tip and logs no warning', async () => {
208+
const pollLimit = 100;
209+
const events = contractEvents(3, 20);
210+
const server = {
211+
getEvents: jest.fn(async () => ({ latestLedger: 500, events })),
212+
};
213+
const store = {
214+
getLastLedger: jest.fn(async () => null),
215+
saveEvent: jest.fn(async () => {}),
216+
setLastLedger: jest.fn(async () => {}),
217+
};
218+
const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() };
219+
220+
const poller = new EventPoller({
221+
enabled: true,
222+
contractId: 'CCONTRACT',
223+
startLedger: 10,
224+
pollLimit,
225+
server,
226+
store,
227+
logger,
228+
});
229+
230+
const result = await poller.pollOnce();
231+
232+
expect(store.setLastLedger).toHaveBeenCalledWith(500);
233+
expect(result.truncated).toBe(false);
234+
expect(logger.warn).not.toHaveBeenCalled();
235+
});
236+
});
98237
});

0 commit comments

Comments
 (0)