From e4c5ce6d0e31146f21497b9ae3d66e5a9f8faa50 Mon Sep 17 00:00:00 2001 From: Kris McGinnes Date: Fri, 25 Sep 2026 16:50:52 -0500 Subject: [PATCH] Stop the request pool pulling new work after a failure --- .../src/utils/mapWithConcurrency.test.ts | 24 +++++++++++++++++++ .../src/utils/mapWithConcurrency.ts | 20 ++++++++++++---- 2 files changed, 40 insertions(+), 4 deletions(-) diff --git a/packages/graph-explorer/src/utils/mapWithConcurrency.test.ts b/packages/graph-explorer/src/utils/mapWithConcurrency.test.ts index 85bdd9742..814d021fc 100644 --- a/packages/graph-explorer/src/utils/mapWithConcurrency.test.ts +++ b/packages/graph-explorer/src/utils/mapWithConcurrency.test.ts @@ -61,4 +61,28 @@ describe("mapWithConcurrency", () => { ), ).rejects.toThrow("boom"); }); + + it("stops pulling work once a callback rejects", async () => { + const started: number[] = []; + const items = [...Array(50).keys()]; + + await expect( + mapWithConcurrency(items, 4, async n => { + started.push(n); + await new Promise(r => setTimeout(r, 5)); + if (n === 0) { + throw new Error("boom"); + } + return n; + }), + ).rejects.toThrow("boom"); + + // Wait out anything still in flight so a lane that kept looping would show up. + await new Promise(r => setTimeout(r, 50)); + + // The three lanes that were already running finish their current item, but + // nothing new starts. Without this, the pool would drain all 50 items in the + // background long after the caller gave up. + expect(started.length).toBeLessThanOrEqual(8); + }); }); diff --git a/packages/graph-explorer/src/utils/mapWithConcurrency.ts b/packages/graph-explorer/src/utils/mapWithConcurrency.ts index 9967cfb3a..288743926 100644 --- a/packages/graph-explorer/src/utils/mapWithConcurrency.ts +++ b/packages/graph-explorer/src/utils/mapWithConcurrency.ts @@ -5,8 +5,14 @@ * the other lanes. Results are returned in the order of `items`. * * A rejected callback propagates: the returned promise rejects once any worker - * throws (already-running callbacks are not cancelled), so callers get - * full-failure semantics. + * throws, so callers get full-failure semantics. The first rejection also stops + * the pool pulling new work, so a caller that gives up on a long queue does not + * leave the remaining lanes issuing requests nobody is waiting for. Callbacks + * already running are not cancelled. + * + * That stop applies to every caller, including the openCypher and SPARQL schema + * and edge connection fetches. They previously drained the whole queue after a + * failure their caller had already given up on, so the change is theirs too. * * @param items The items to process * @param concurrency The maximum number of callbacks running at once @@ -20,12 +26,18 @@ export default async function mapWithConcurrency( ): Promise { const results: Result[] = []; let cursor = 0; + let failed = false; async function worker() { - while (cursor < items.length) { + while (cursor < items.length && !failed) { const index = cursor; cursor += 1; - results[index] = await callback(items[index]); + try { + results[index] = await callback(items[index]); + } catch (error) { + failed = true; + throw error; + } } }