Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions packages/graph-explorer/src/utils/mapWithConcurrency.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
});
});
20 changes: 16 additions & 4 deletions packages/graph-explorer/src/utils/mapWithConcurrency.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -20,12 +26,18 @@ export default async function mapWithConcurrency<Item, Result>(
): Promise<Result[]> {
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;
}
}
}

Expand Down
Loading