Skip to content

Commit a3918e0

Browse files
authored
fix(store-sync): setup message listener before setting up subscription (#3765)
1 parent bb5f221 commit a3918e0

4 files changed

Lines changed: 40 additions & 13 deletions

File tree

.changeset/sharp-shrimps-dream.md

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+
Fixed a race condition in the preconfirmed logs stream by setting up the message listener before setting up the subscription.

packages/store-sync/package.json

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@
7070
"clean:js": "shx rm -rf dist",
7171
"dev": "tsup --watch",
7272
"lint": "eslint .",
73-
"playground": "DEBUG=mud:* tsx playground/index.ts",
73+
"playground": "DEBUG=mud:* tsx playground/index.ts | tee playground.log",
7474
"test": "vitest",
7575
"test:ci": "vitest --run"
7676
},

packages/store-sync/src/createPreconfirmedBlockStream.ts

Lines changed: 26 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -73,11 +73,13 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
7373

7474
let preconfirmedTransactionLogs: { [txHash: string]: Partial<StoreEventsLog>[] | undefined } = {};
7575
let preconfirmedLogsState: "initializing" | "initialized" | "waiting" = "waiting";
76+
let firstPreconfirmedBlockNumber: bigint | undefined = undefined;
7677
let attempt = 0;
7778
const preconfirmedBlockLogs$ = recreatePreconfirmedStream$.pipe(
7879
tap(() => {
7980
if (attempt !== 0) debug(`waiting ${attempt * 500}ms before initializing preconfirmed logs stream`);
8081
preconfirmedLogsState = "initializing";
82+
firstPreconfirmedBlockNumber = undefined;
8183
preconfirmedTransactionLogs = {};
8284
}),
8385
switchMap(() => timer(attempt * 500)),
@@ -114,6 +116,10 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
114116
return !isProcessedBlock;
115117
}),
116118
tap((block) => {
119+
if (preconfirmedLogsState !== "initialized") {
120+
firstPreconfirmedBlockNumber = block.blockNumber;
121+
debug("first preconfirmed block number", firstPreconfirmedBlockNumber);
122+
}
117123
debug("preconfirmed block", block.blockNumber, "with", block.logs.length, "logs");
118124
preconfirmedLogsState = "initialized";
119125
attempt = 0;
@@ -125,6 +131,16 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
125131
}
126132
(preconfirmedTransactionLogs[txHash] ??= []).push(log);
127133
});
134+
debug(
135+
"preconfirmed logs state",
136+
Object.fromEntries(
137+
Object.entries(preconfirmedTransactionLogs).map(([txHash, logs]) => [
138+
txHash,
139+
{ numLogs: logs?.length, blockNumber: logs?.[0]?.blockNumber },
140+
]),
141+
),
142+
"transactions",
143+
);
128144
}),
129145
);
130146

@@ -133,7 +149,12 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
133149
processedLatestBlockNumber = block.blockNumber;
134150

135151
const mismatchingTransactions: string[] = [];
136-
if (preconfirmedLogsState === "initialized") {
152+
const confirmPreconfirmedLogs =
153+
preconfirmedLogsState === "initialized" &&
154+
firstPreconfirmedBlockNumber &&
155+
block.blockNumber >= firstPreconfirmedBlockNumber;
156+
157+
if (confirmPreconfirmedLogs) {
137158
const logsByTransaction = groupBy(
138159
block.logs.filter((log) => log.transactionHash) as StoreEventsLog[],
139160
(log) => log.transactionHash,
@@ -148,7 +169,7 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
148169
JSON.stringify(
149170
{
150171
txHash,
151-
numPreconfirmedLogs: preconfirmedLogs?.length,
172+
numPreconfirmedLogs: preconfirmedLogs?.length ?? "none",
152173
numLatestLogs: latestLogs.length,
153174
missingLogs: latestLogs.filter(
154175
(log) => !preconfirmedLogs?.find((preconfirmedLog) => log.logIndex === preconfirmedLog.logIndex),
@@ -190,14 +211,14 @@ export function createPreconfirmedBlockStream(opts: PreconfirmedBlockStreamOptio
190211
}
191212

192213
// While the preconfirmed logs stream is initializing, don't recreate it and pass the block through
193-
if (preconfirmedLogsState === "initializing") {
194-
debug("preconfirmed logs stream is initializing, not recreating");
214+
if (!confirmPreconfirmedLogs) {
215+
debug("block is before first preconfirmed block, not recreating preconfirmed stream");
195216
return block;
196217
}
197218

198219
// If the preconfirmed logs stream is initialized but there are mismatching logs, recreate it and pass the block through.
199220
// Pass all logs from this block, not just the mismatching ones, to make sure they appear in the right order.
200-
if (preconfirmedLogsState === "initialized" && mismatchingTransactions.length > 0) {
221+
if (mismatchingTransactions.length > 0) {
201222
debug("mismatching transactions found in latest block", block.blockNumber, "recreating preconfirmed stream");
202223
recreatePreconfirmedStream$.next();
203224
return block;

packages/store-sync/src/watchLogs.ts

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ export function watchLogs({ url, address, fromBlock }: WatchLogsInput): WatchLog
4343
// Buffer the live logs received until the gap from `fromBlock` to `currentBlock` is closed
4444
let caughtUp = false;
4545
const logBuffer: StoreEventsLog[] = [];
46+
let subscriptionId: Hex | undefined = undefined;
4647

4748
ws = new WebSocket(url);
4849

@@ -78,13 +79,6 @@ export function watchLogs({ url, address, fromBlock }: WatchLogsInput): WatchLog
7879
debug(`ws${wsId} upgrade`);
7980
});
8081

81-
const subscriptionId = await request<Hex>({
82-
ws,
83-
method: "wiresaw_watchLogs",
84-
params: [{ address, topics }],
85-
wsId,
86-
});
87-
8882
ws.on("message", (message) => {
8983
const data = JSON.parse(message.toString());
9084

@@ -112,6 +106,13 @@ export function watchLogs({ url, address, fromBlock }: WatchLogsInput): WatchLog
112106
debug(`ws${wsId} message`);
113107
});
114108

109+
subscriptionId = await request<Hex>({
110+
ws,
111+
method: "wiresaw_watchLogs",
112+
params: [{ address, topics }],
113+
wsId,
114+
});
115+
115116
// Catch up to the pending logs
116117
fetchInitialLogs({ ws, address, fromBlock, topics, wsId })
117118
.then((initialLogs) => {

0 commit comments

Comments
 (0)