Skip to content

Commit f71d689

Browse files
authored
Improve support for sparse data (#196)
2 parents 57ee8e4 + d802fb1 commit f71d689

5 files changed

Lines changed: 46 additions & 2 deletions

File tree

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
{
2+
"type": "prerelease",
3+
"comment": "rpc: periodically send empty blocks",
4+
"packageName": "@apibara/evm-rpc",
5+
"email": "francesco@ceccon.me",
6+
"dependentChangeType": "patch"
7+
}
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
{
2+
"type": "prerelease",
3+
"comment": "rpc: periodically send empty blocks",
4+
"packageName": "@apibara/protocol",
5+
"email": "francesco@ceccon.me",
6+
"dependentChangeType": "patch"
7+
}

packages/evm-rpc/src/stream-config.ts

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,7 @@ export class EvmRpcStream extends RpcStreamConfig<Filter, Block> {
121121
async fetchBlockRange({
122122
startBlock,
123123
finalizedBlock,
124+
force,
124125
filter,
125126
}: FetchBlockRangeArgs<Filter>): Promise<FetchBlockRangeResult<Block>> {
126127
const { start: fromBlock, end: toBlock } = this.blockRangeOracle.clampRange(
@@ -151,6 +152,12 @@ export class EvmRpcStream extends RpcStreamConfig<Filter, Block> {
151152
}),
152153
);
153154
}
155+
} else if (force && blockNumbers.length === 0) {
156+
blockNumberResponses.push(
157+
this.fetchBlockHeaderByNumberWithRetry({
158+
blockNumber: toBlock,
159+
}),
160+
);
154161
} else {
155162
for (const blockNumber of blockNumbers) {
156163
blockNumberResponses.push(
@@ -180,6 +187,7 @@ export class EvmRpcStream extends RpcStreamConfig<Filter, Block> {
180187
async fetchBlockByNumber({
181188
blockNumber,
182189
expectedParentBlockHash,
190+
isAtHead,
183191
filter,
184192
}: FetchBlockByNumberArgs<Filter>): Promise<FetchBlockByNumberResult<Block>> {
185193
// Fetch block header and check it matches the expected parent block hash.
@@ -226,8 +234,12 @@ export class EvmRpcStream extends RpcStreamConfig<Filter, Block> {
226234

227235
let block = null;
228236

229-
// TODO: handle header on new block.
230-
if (filter.header === "always" || logs.length > 0) {
237+
const shouldSendBlock =
238+
filter.header === "always" ||
239+
logs.length > 0 ||
240+
(filter.header === "on_data_or_on_new_block" && isAtHead);
241+
242+
if (shouldSendBlock) {
231243
block = {
232244
header,
233245
logs,

packages/protocol/src/rpc/config.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ import type { Bytes, Cursor } from "../common";
33
export type FetchBlockRangeArgs<TFilter> = {
44
startBlock: bigint;
55
finalizedBlock: bigint;
6+
force: boolean;
67
filter: TFilter;
78
};
89

@@ -27,6 +28,7 @@ export type BlockInfo = {
2728
export type FetchBlockByNumberArgs<TFilter> = {
2829
blockNumber: bigint;
2930
expectedParentBlockHash: Bytes;
31+
isAtHead: boolean;
3032
filter: TFilter;
3133
};
3234

packages/protocol/src/rpc/data-stream.ts

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@ type State<TFilter, TBlock> = {
1616
lastFinalizedRefresh: number;
1717
// When the last heartbeat was sent.
1818
lastHeartbeat: number;
19+
// When the last backfill message was sent.
20+
lastBackfillMessage: number;
1921
// Track the chain's state.
2022
chainTracker: ChainTracker;
2123
// Heartbeat interval in milliseconds.
@@ -95,6 +97,7 @@ export class RpcDataStream<TFilter, TBlock> {
9597
cursor,
9698
lastHeartbeat: Date.now(),
9799
lastFinalizedRefresh: Date.now(),
100+
lastBackfillMessage: Date.now(),
98101
chainTracker,
99102
config: this.config,
100103
heartbeatIntervalMs: this.heartbeatIntervalMs,
@@ -165,9 +168,14 @@ async function* backfillFinalizedBlocks<TFilter, TBlock>(
165168
const { cursor, chainTracker, config, filter } = state;
166169
const finalized = chainTracker.finalized();
167170

171+
// While backfilling we want to regularly send some blocks (even if empty) so
172+
// that the client can store the cursor.
173+
const force = shouldForceBackfill(state);
174+
168175
const filterData = await config.fetchBlockRange({
169176
startBlock: cursor.orderKey + 1n,
170177
finalizedBlock: finalized.orderKey,
178+
force,
171179
filter,
172180
});
173181

@@ -179,6 +187,7 @@ async function* backfillFinalizedBlocks<TFilter, TBlock>(
179187

180188
for (const data of filterData.data) {
181189
state.lastHeartbeat = Date.now();
190+
state.lastBackfillMessage = Date.now();
182191
yield {
183192
_tag: "data",
184193
data: {
@@ -213,6 +222,7 @@ async function* produceNextBlock<TFilter, TBlock>(
213222

214223
const result = await state.config.fetchBlockByNumber({
215224
blockNumber: state.cursor.orderKey + 1n,
225+
isAtHead: isAtHead(state),
216226
expectedParentBlockHash: currentBlockHash,
217227
filter: state.filter,
218228
});
@@ -325,6 +335,12 @@ function shouldSendHeartbeat(state: State<unknown, unknown>): boolean {
325335
return now - lastHeartbeat >= heartbeatIntervalMs;
326336
}
327337

338+
function shouldForceBackfill(state: State<unknown, unknown>): boolean {
339+
const { lastBackfillMessage, heartbeatIntervalMs } = state;
340+
const now = Date.now();
341+
return now - lastBackfillMessage >= heartbeatIntervalMs;
342+
}
343+
328344
function shouldContinue(state: State<unknown, unknown>): boolean {
329345
const { endingCursor } = state.options || {};
330346
if (endingCursor === undefined) return true;

0 commit comments

Comments
 (0)