Skip to content

Commit 23b0c9a

Browse files
authored
feat(store-sync): make pending logs sync more resilient (#3743)
1 parent 0f5c75b commit 23b0c9a

6 files changed

Lines changed: 367 additions & 75 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@latticexyz/store-sync": patch
3+
---
4+
5+
The sync stack now handles downtime in the pending logs API and reconnects once it's available again.

packages/entrykit/playground/common.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,13 @@
11
import { Hex } from "viem";
22
import { anvil } from "viem/chains";
3-
import { garnet, redstone } from "@latticexyz/common/chains";
3+
import { garnet, pyrope, redstone } from "@latticexyz/common/chains";
44

55
const testWorlds = {
66
// TODO: get this from somewhere else, like playground deploy output
77
[anvil.id]: "0x60e7e3caed67b9d2cca14519b6cd7700a7d4ee66",
88
[redstone.id]: "0xf75b1b7bdb6932e487c4aa8d210f4a682abeacf0",
99
[garnet.id]: "0x1d1BBb7e359a7425428adBe799d84601B221b4E8",
10+
[pyrope.id]: "0xe8d5C603b4A501B6B5EAB46a14dA05f38e44F64f",
1011
} as Partial<Record<string, Hex>>;
1112

1213
const searchParams = new URLSearchParams(window.location.search);

packages/entrykit/playground/wagmiConfig.ts

Lines changed: 47 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,30 @@
11
import { Chain, http } from "viem";
2-
import { anvil, mainnet } from "viem/chains";
2+
import { anvil, redstone, mainnet } from "viem/chains";
33
import { createWagmiConfig } from "../src/createWagmiConfig";
44
import { chainId } from "./common";
5-
import { garnet } from "@latticexyz/common/chains";
5+
import { garnet, pyrope } from "@latticexyz/common/chains";
66
import { wiresaw } from "@latticexyz/common/internal";
77

8+
const redstoneWithPaymaster = {
9+
...redstone,
10+
rpcUrls: {
11+
...redstone.rpcUrls,
12+
wiresaw: {
13+
http: ["https://wiresaw.redstonechain.com"],
14+
webSocket: ["wss://wiresaw.redstonechain.com"],
15+
},
16+
bundler: {
17+
http: ["https://rpc.redstonechain.com"],
18+
webSocket: ["wss://rpc.redstonechain.com"],
19+
},
20+
},
21+
contracts: {
22+
quarryPaymaster: {
23+
address: "0x2d70F1eFFbFD865764CAF19BE2A01a72F3CE774f",
24+
},
25+
},
26+
};
27+
828
const garnetWithPaymaster = {
929
...garnet,
1030
rpcUrls: {
@@ -47,12 +67,29 @@ const anvilWithPaymaster = {
4767
},
4868
};
4969

50-
const chains = [mainnet, garnetWithPaymaster, anvilWithPaymaster] as const satisfies Chain[];
70+
const pyropeWithPaymaster = {
71+
...pyrope,
72+
contracts: {
73+
quarryPaymaster: {
74+
address: "0xD40C9cFc97b855B2183D7e1c8925edF77C309b85",
75+
},
76+
},
77+
};
78+
79+
const chains = [
80+
mainnet,
81+
garnetWithPaymaster,
82+
anvilWithPaymaster,
83+
redstoneWithPaymaster,
84+
pyropeWithPaymaster,
85+
] as const satisfies Chain[];
5186

5287
const transports = {
5388
[mainnet.id]: http(),
54-
[anvil.id]: http(),
55-
[garnet.id]: wiresaw(),
89+
[anvilWithPaymaster.id]: http(),
90+
[garnetWithPaymaster.id]: wiresaw(),
91+
[redstoneWithPaymaster.id]: wiresaw(),
92+
[pyropeWithPaymaster.id]: http(),
5693
} as const;
5794

5895
export const wagmiConfig = createWagmiConfig({
@@ -62,7 +99,10 @@ export const wagmiConfig = createWagmiConfig({
6299
chains,
63100
transports,
64101
pollingInterval: {
65-
[anvil.id]: 500,
66-
[garnet.id]: 2000,
102+
[mainnet.id]: 2000,
103+
[anvilWithPaymaster.id]: 500,
104+
[garnetWithPaymaster.id]: 2000,
105+
[redstoneWithPaymaster.id]: 2000,
106+
[pyropeWithPaymaster.id]: 2000,
67107
},
68108
});
Lines changed: 232 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,232 @@
1+
import {
2+
catchError,
3+
combineLatest,
4+
concatMap,
5+
from,
6+
map,
7+
mergeMap,
8+
Observable,
9+
of,
10+
Subject,
11+
tap,
12+
throwError,
13+
switchMap,
14+
merge,
15+
filter,
16+
startWith,
17+
} from "rxjs";
18+
import { StorageAdapterBlock, StoreEventsLog, SyncFilter } from "./common";
19+
import { watchLogs } from "./watchLogs";
20+
import { Hex } from "viem";
21+
import { fromEventSource } from "./fromEventSource";
22+
import { isLogsApiResponse } from "./indexer-client/isLogsApiResponse";
23+
import { toStorageAdapterBlock } from "./indexer-client/toStorageAdapterBlock";
24+
import { fetchAndStoreLogs } from "./fetchAndStoreLogs";
25+
import { storeEventsAbi } from "@latticexyz/store";
26+
import { bigIntMax, isDefined } from "@latticexyz/common/utils";
27+
import { getRpcClient, GetRpcClientOptions } from "@latticexyz/block-logs-stream";
28+
import { debug } from "./debug";
29+
30+
type PreconfirmedBlockStreamOptions = GetRpcClientOptions & {
31+
fromBlock: bigint;
32+
preconfirmedLogsUrl: string;
33+
indexerUrl?: string;
34+
chainId: number;
35+
address?: Hex;
36+
filters: SyncFilter[];
37+
latestBlockNumber$: Observable<bigint>;
38+
maxBlockRange?: bigint;
39+
};
40+
41+
export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptions): Observable<StorageAdapterBlock> {
42+
const recreatePreconfirmedStream$ = new Subject<void>();
43+
const recreateLatestStream$ = new Subject<void>();
44+
45+
let restartBlockNumber = opts.fromBlock;
46+
let initialCatchUpBlockNumber: bigint | undefined = undefined;
47+
getRpcClient(opts)
48+
.request({ method: "eth_blockNumber" })
49+
.then((blockNumber) => {
50+
console.log("initial catch up block number", BigInt(blockNumber));
51+
initialCatchUpBlockNumber = BigInt(blockNumber);
52+
});
53+
54+
const latestBlock$ = recreateLatestStream$.pipe(
55+
startWith(undefined),
56+
switchMap(() =>
57+
createLatestBlockStream({ ...opts, fromBlock: restartBlockNumber }).pipe(
58+
catchError((e) => {
59+
debug("Error in latest block stream, recreating", e);
60+
recreateLatestStream$.next();
61+
return throwError(() => e);
62+
}),
63+
),
64+
),
65+
);
66+
67+
let processedBlockLogs: { [blockNumber: string]: { [logIndex: number]: boolean } } = {};
68+
let preconfirmedLogsState: "initializing" | "initialized" | "waiting" = "waiting";
69+
const preconfirmedLogs$ = recreatePreconfirmedStream$.pipe(
70+
tap(() => {
71+
debug("initializing preconfirmed logs stream");
72+
preconfirmedLogsState = "initializing";
73+
processedBlockLogs = {};
74+
}),
75+
switchMap(() =>
76+
watchLogs({
77+
...opts,
78+
url: opts.preconfirmedLogsUrl,
79+
fromBlock: restartBlockNumber,
80+
}).logs$.pipe(
81+
catchError((e) => {
82+
debug("Error in preconfirmed logs stream, recreating", e);
83+
recreatePreconfirmedStream$.next();
84+
return throwError(() => e);
85+
}),
86+
),
87+
),
88+
tap((block) => {
89+
debug("preconfirmed block", block.blockNumber, "with", block.logs.length, "logs");
90+
preconfirmedLogsState = "initialized";
91+
restartBlockNumber = block.blockNumber;
92+
const seenLogs = (processedBlockLogs[String(block.blockNumber)] ??= {});
93+
block.logs.forEach((log) => {
94+
seenLogs[log.logIndex!] = true;
95+
});
96+
debug("got preconfirmed block", block.blockNumber, "with", block.logs.length, "logs");
97+
}),
98+
);
99+
100+
const missingLogs$ = latestBlock$.pipe(
101+
map((block) => {
102+
const missingBlock = processedBlockLogs[String(block.blockNumber)] == null;
103+
const seenLogs = processedBlockLogs[String(block.blockNumber)] ?? {};
104+
const missingLogs = block.logs.filter((log) => !seenLogs[log.logIndex!]);
105+
delete processedBlockLogs[String(block.blockNumber)];
106+
restartBlockNumber = block.blockNumber + 1n;
107+
108+
debug(
109+
"got latest block",
110+
block.blockNumber,
111+
"with",
112+
block.logs.length,
113+
"logs (",
114+
missingBlock ? "missing block," : "block seen,",
115+
`${missingLogs.length} new logs`,
116+
")",
117+
);
118+
119+
if (preconfirmedLogsState === "waiting") {
120+
// Once the initial catch up block is reached, we can start the preconfirmed logs stream
121+
if (
122+
initialCatchUpBlockNumber != null && // initial catch up block fetched
123+
block.blockNumber >= initialCatchUpBlockNumber // initial catch up block reached
124+
) {
125+
debug("initial catch up block number", initialCatchUpBlockNumber, "reached, creating preconfirmed stream");
126+
recreatePreconfirmedStream$.next();
127+
}
128+
// While the preconfirmed logs stream is waiting, pass the block through
129+
return block;
130+
}
131+
132+
// While the preconfirmed logs stream is initializing, don't recreate it and pass the block through
133+
if (preconfirmedLogsState === "initializing") {
134+
debug("preconfirmed logs stream is initializing, not recreating");
135+
return block;
136+
}
137+
138+
// If the preconfirmed logs stream is initialized but there are missing logs, recreate it and pass the block through.
139+
// Pass all logs from this block, not just the missing ones, to make sure they appear in the right order.
140+
if (preconfirmedLogsState === "initialized" && (missingLogs.length > 0 || missingBlock)) {
141+
debug("missing logs found in latest block", block.blockNumber, "recreating preconfirmed stream", {
142+
missingLogs: missingLogs.length,
143+
missingBlock,
144+
});
145+
recreatePreconfirmedStream$.next();
146+
return block;
147+
}
148+
149+
debug("no missing logs found in latest block", block.blockNumber, "not recreating preconfirmed stream");
150+
return;
151+
}),
152+
filter(isDefined),
153+
);
154+
155+
return merge(preconfirmedLogs$, missingLogs$);
156+
}
157+
158+
// TODO: refactor to reduce duplication with indexer/rpc stream in `createStoreSync.ts`
159+
function createLatestBlockStream({
160+
fromBlock,
161+
indexerUrl,
162+
chainId,
163+
address,
164+
filters,
165+
latestBlockNumber$,
166+
maxBlockRange,
167+
...opts
168+
}: PreconfirmedBlockStreamOptions): Observable<StorageAdapterBlock> {
169+
const indexerBlocks$ = indexerUrl
170+
? of(indexerUrl).pipe(
171+
mergeMap((indexerUrl) => {
172+
const url = new URL(
173+
`api/logs-live?${new URLSearchParams({
174+
input: JSON.stringify({ chainId, address, filters }),
175+
block_num: fromBlock.toString(),
176+
include_tx_hash: "true",
177+
})}`,
178+
indexerUrl,
179+
);
180+
return fromEventSource<string>(url);
181+
}),
182+
map((messageEvent) => {
183+
const data = JSON.parse(messageEvent.data);
184+
if (!isLogsApiResponse(data)) {
185+
throw new Error("Received unexpected from indexer:" + messageEvent.data);
186+
}
187+
return toStorageAdapterBlock(data);
188+
}),
189+
)
190+
: throwError(() => new Error("No indexer URL provided"));
191+
192+
let lastBlockNumberProcessed = 0n;
193+
const ethRpcBlocks$ = combineLatest([of(fromBlock), latestBlockNumber$]).pipe(
194+
map(([startBlock, endBlock]) => ({ startBlock, endBlock })),
195+
concatMap((range) => {
196+
const storedBlocks = fetchAndStoreLogs({
197+
...opts,
198+
address,
199+
events: storeEventsAbi,
200+
maxBlockRange,
201+
fromBlock: lastBlockNumberProcessed
202+
? bigIntMax(range.startBlock, lastBlockNumberProcessed + 1n)
203+
: range.startBlock,
204+
toBlock: range.endBlock,
205+
logFilter: filters.length
206+
? (log: StoreEventsLog): boolean =>
207+
filters.some(
208+
(filter) =>
209+
filter.tableId === log.args.tableId &&
210+
(filter.key0 == null || filter.key0 === log.args.keyTuple[0]) &&
211+
(filter.key1 == null || filter.key1 === log.args.keyTuple[1]),
212+
)
213+
: undefined,
214+
storageAdapter: () => Promise.resolve(),
215+
});
216+
return from(storedBlocks);
217+
}),
218+
tap((block) => {
219+
lastBlockNumberProcessed = block.blockNumber;
220+
}),
221+
);
222+
223+
const latestBlock$ = indexerBlocks$.pipe(
224+
catchError((error) => {
225+
debug("failed to stream logs from indexer:", error.message);
226+
debug("falling back to streaming logs from ETH RPC");
227+
return ethRpcBlocks$;
228+
}),
229+
);
230+
231+
return latestBlock$;
232+
}

0 commit comments

Comments
 (0)