mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
perf(desktop): parallelize relay reconnect replay
Finding L3-1: reconnect replay previously awaited every live subscription REQ and paged history replay serially, making reconnect latency grow with active subscription count. Start live REQs together, then replay paged channel catch-up with a bounded cap while preserving ordering inside each subscription's page loop. Co-authored-by: Tyler Longwell <tlongwell@block.xyz> Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
a9a8d0a309
commit
261e047d5a
@@ -172,3 +172,81 @@ test("channel reconnect replay pages the missed window until a short page", asyn
|
||||
]);
|
||||
assert.equal(delivered.length, 1008);
|
||||
});
|
||||
|
||||
test("reconnect replay starts live REQs in parallel and preserves per-sub page order", async () => {
|
||||
const sentPayloads = [];
|
||||
const sendResolvers = [];
|
||||
const historyFiltersByChannel = {
|
||||
"channel-1": [],
|
||||
"channel-2": [],
|
||||
};
|
||||
const pagesByChannel = {
|
||||
"channel-1": [
|
||||
eventRange("c1-full", 1501, 500),
|
||||
eventRange("c1-short", 1490, 2),
|
||||
],
|
||||
"channel-2": [
|
||||
eventRange("c2-full", 1701, 500),
|
||||
eventRange("c2-short", 1690, 2),
|
||||
],
|
||||
};
|
||||
const subscriptions = new Map([
|
||||
[
|
||||
"live-1",
|
||||
{
|
||||
mode: "live",
|
||||
filter: buildChannelFilter("channel-1", 50),
|
||||
onEvent: () => {},
|
||||
lastSeenCreatedAt: 1000,
|
||||
},
|
||||
],
|
||||
[
|
||||
"live-2",
|
||||
{
|
||||
mode: "live",
|
||||
filter: buildChannelFilter("channel-2", 50),
|
||||
onEvent: () => {},
|
||||
lastSeenCreatedAt: 1000,
|
||||
},
|
||||
],
|
||||
]);
|
||||
|
||||
const replayPromise = replayLiveSubscriptions({
|
||||
subscriptions,
|
||||
now: 2000,
|
||||
pageReplayConcurrency: 2,
|
||||
sendRaw: (payload) => {
|
||||
sentPayloads.push(payload);
|
||||
return new Promise((resolve) => {
|
||||
sendResolvers.push(resolve);
|
||||
});
|
||||
},
|
||||
requestHistory: async (filter) => {
|
||||
const channelId = filter["#h"]?.[0];
|
||||
historyFiltersByChannel[channelId].push(filter.until);
|
||||
return pagesByChannel[channelId].shift() ?? [];
|
||||
},
|
||||
});
|
||||
|
||||
await Promise.resolve();
|
||||
|
||||
assert.deepEqual(
|
||||
sentPayloads.map((payload) => payload[1]),
|
||||
["live-1", "live-2"],
|
||||
);
|
||||
assert.equal(sendResolvers.length, 2);
|
||||
assert.deepEqual(historyFiltersByChannel, {
|
||||
"channel-1": [],
|
||||
"channel-2": [],
|
||||
});
|
||||
|
||||
for (const resolve of sendResolvers) {
|
||||
resolve();
|
||||
}
|
||||
await replayPromise;
|
||||
|
||||
assert.deepEqual(historyFiltersByChannel, {
|
||||
"channel-1": [2000, 1501],
|
||||
"channel-2": [2000, 1701],
|
||||
});
|
||||
});
|
||||
|
||||
@@ -7,6 +7,25 @@ import type { RelayEvent } from "@/shared/api/types";
|
||||
|
||||
const RECONNECT_REPLAY_SKEW_SECS = 5;
|
||||
export const RECONNECT_REPLAY_PAGE_LIMIT = 500;
|
||||
export const RECONNECT_REPLAY_PAGE_CONCURRENCY = 4;
|
||||
|
||||
async function runWithConcurrency<T>(
|
||||
items: T[],
|
||||
concurrency: number,
|
||||
worker: (item: T) => Promise<void>,
|
||||
) {
|
||||
const workerCount = Math.min(Math.max(1, concurrency), items.length);
|
||||
let nextIndex = 0;
|
||||
|
||||
await Promise.all(
|
||||
Array.from({ length: workerCount }, async () => {
|
||||
while (nextIndex < items.length) {
|
||||
const item = items[nextIndex++];
|
||||
await worker(item);
|
||||
}
|
||||
}),
|
||||
);
|
||||
}
|
||||
|
||||
export function buildReconnectReplayFilter(
|
||||
filter: RelaySubscriptionFilter,
|
||||
@@ -84,34 +103,60 @@ export async function replayLiveSubscriptions({
|
||||
sendRaw,
|
||||
requestHistory,
|
||||
now = Math.floor(Date.now() / 1_000),
|
||||
pageReplayConcurrency = RECONNECT_REPLAY_PAGE_CONCURRENCY,
|
||||
}: {
|
||||
subscriptions: Map<string, RelaySubscription>;
|
||||
sendRaw: (payload: unknown[]) => Promise<void>;
|
||||
requestHistory: (filter: RelaySubscriptionFilter) => Promise<RelayEvent[]>;
|
||||
now?: number;
|
||||
pageReplayConcurrency?: number;
|
||||
}) {
|
||||
for (const [subId, subscription] of subscriptions) {
|
||||
if (subscription.mode !== "live") continue;
|
||||
const replayRequests = Array.from(subscriptions.entries())
|
||||
.filter(
|
||||
(
|
||||
entry,
|
||||
): entry is [string, Extract<RelaySubscription, { mode: "live" }>] =>
|
||||
entry[1].mode === "live",
|
||||
)
|
||||
.map(([subId, subscription]) => {
|
||||
const replaySince =
|
||||
subscription.lastSeenCreatedAt === undefined
|
||||
? undefined
|
||||
: Math.max(
|
||||
0,
|
||||
subscription.lastSeenCreatedAt - RECONNECT_REPLAY_SKEW_SECS,
|
||||
);
|
||||
const shouldPageReplay =
|
||||
replaySince !== undefined &&
|
||||
shouldPageReconnectReplay(subscription.filter);
|
||||
|
||||
const replaySince =
|
||||
subscription.lastSeenCreatedAt === undefined
|
||||
? undefined
|
||||
: Math.max(
|
||||
0,
|
||||
subscription.lastSeenCreatedAt - RECONNECT_REPLAY_SKEW_SECS,
|
||||
);
|
||||
const shouldPageReplay =
|
||||
replaySince !== undefined &&
|
||||
shouldPageReconnectReplay(subscription.filter);
|
||||
await sendRaw([
|
||||
"REQ",
|
||||
subId,
|
||||
shouldPageReplay
|
||||
? subscription.filter
|
||||
: buildReconnectReplayFilter(subscription.filter, replaySince),
|
||||
]);
|
||||
return { subId, subscription, replaySince, shouldPageReplay };
|
||||
});
|
||||
|
||||
if (shouldPageReplay) {
|
||||
await Promise.all(
|
||||
replayRequests.map(
|
||||
({ subId, subscription, replaySince, shouldPageReplay }) =>
|
||||
sendRaw([
|
||||
"REQ",
|
||||
subId,
|
||||
shouldPageReplay
|
||||
? subscription.filter
|
||||
: buildReconnectReplayFilter(subscription.filter, replaySince),
|
||||
]),
|
||||
),
|
||||
);
|
||||
|
||||
await runWithConcurrency(
|
||||
replayRequests.filter(
|
||||
(
|
||||
request,
|
||||
): request is typeof request & {
|
||||
replaySince: number;
|
||||
shouldPageReplay: true;
|
||||
} => request.shouldPageReplay && request.replaySince !== undefined,
|
||||
),
|
||||
pageReplayConcurrency,
|
||||
async ({ subId, subscription, replaySince }) => {
|
||||
await replayReconnectHistoryPages({
|
||||
subscription,
|
||||
since: replaySince,
|
||||
@@ -119,6 +164,6 @@ export async function replayLiveSubscriptions({
|
||||
isActive: () => subscriptions.get(subId) === subscription,
|
||||
requestHistory,
|
||||
});
|
||||
}
|
||||
}
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user