diff --git a/desktop/src/shared/api/concurrency.ts b/desktop/src/shared/api/concurrency.ts new file mode 100644 index 000000000..925af8921 --- /dev/null +++ b/desktop/src/shared/api/concurrency.ts @@ -0,0 +1,20 @@ +export async function collectWithConcurrency( + items: T[], + concurrency: number, + worker: (item: T) => Promise, +): Promise { + const workerCount = Math.min(Math.max(1, concurrency), items.length); + const results = new Array(items.length); + let nextIndex = 0; + + await Promise.all( + Array.from({ length: workerCount }, async () => { + while (nextIndex < items.length) { + const currentIndex = nextIndex++; + results[currentIndex] = await worker(items[currentIndex]); + } + }), + ); + + return results; +} diff --git a/desktop/src/shared/api/relayClientSession.ts b/desktop/src/shared/api/relayClientSession.ts index c2462d08c..63b71497d 100644 --- a/desktop/src/shared/api/relayClientSession.ts +++ b/desktop/src/shared/api/relayClientSession.ts @@ -29,6 +29,7 @@ import { buildChannelMentionFilter, buildGlobalStreamFilter, } from "@/shared/api/relayChannelFilters"; +import { collectWithConcurrency } from "@/shared/api/concurrency"; import { replayLiveSubscriptions } from "@/shared/api/relayReconnectReplay"; import { RelayConnectionStateEmitter } from "@/shared/api/relayConnectionStateEmitter"; import { @@ -41,7 +42,8 @@ import { buildThreadReferenceTags } from "@/features/messages/lib/threading"; const RECONNECT_BASE_DELAY_MS = 1_000, RECONNECT_MAX_DELAY_MS = 30_000, - EVENT_BATCH_MS = 16; + EVENT_BATCH_MS = 16, + AUX_BACKFILL_CONCURRENCY = 4; /** * Passive liveness check. The relay sends heartbeat pings every 30s; if no @@ -217,10 +219,11 @@ export class RelayClient { chunks.push(eventIds.slice(i, i + AUX_BACKFILL_CHUNK_SIZE)); } - const batches: RelayEvent[][] = []; - for (const ids of chunks) { - batches.push(await this.requestHistory(buildFilter(channelId, ids))); - } + const batches = await collectWithConcurrency( + chunks, + AUX_BACKFILL_CONCURRENCY, + (ids) => this.requestHistory(buildFilter(channelId, ids)), + ); return batches.flat(); }