diff --git a/desktop/src/shared/api/relayReconnectReplay.test.mjs b/desktop/src/shared/api/relayReconnectReplay.test.mjs index a58aa7f25..63f6402cd 100644 --- a/desktop/src/shared/api/relayReconnectReplay.test.mjs +++ b/desktop/src/shared/api/relayReconnectReplay.test.mjs @@ -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], + }); +}); diff --git a/desktop/src/shared/api/relayReconnectReplay.ts b/desktop/src/shared/api/relayReconnectReplay.ts index d41f2c3ba..65709d351 100644 --- a/desktop/src/shared/api/relayReconnectReplay.ts +++ b/desktop/src/shared/api/relayReconnectReplay.ts @@ -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( + items: T[], + concurrency: number, + worker: (item: T) => Promise, +) { + 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; sendRaw: (payload: unknown[]) => Promise; requestHistory: (filter: RelaySubscriptionFilter) => Promise; 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] => + 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, }); - } - } + }, + ); }