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
31 changes: 20 additions & 11 deletions src/Lazy/concurrentPool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,14 @@ function concurrentPool<A>(
let prev = Promise.resolve();

async function fetchNextItem() {
if (idleWorkers <= 0) {
if (idleWorkers <= 0 || finished) {
return;
}
// While a consumer is waiting, keep the pool saturated regardless of the
// buffer (order-preserving pools must buffer past a slow head item, as
// p-map does). Only when nobody is waiting is prefetch capped at the
// pool size, so a slow consumer cannot grow the buffer unboundedly.
if (consumer.length === 0 && itemMap.size >= idleWorkers) {
return;
}

Expand Down Expand Up @@ -114,7 +121,7 @@ function concurrentPool<A>(
}
}

async function pullItem(item: Item<A>): Promise<void> {
function deliverFromBuffer() {
while (consumer.length > 0) {
const id = consumer[0][2];
if (!itemMap.has(id)) {
Expand All @@ -132,10 +139,11 @@ function concurrentPool<A>(
resolve(value.success);
}
}
}

async function pullItem(item: Item<A>): Promise<void> {
deliverFromBuffer();

/**
* @TODO ppeeou: If an error occurs, the next iterator should not be executed
*/
if (!item.success.done && consumer.length > 0) {
fetchNextItem();
}
Expand All @@ -152,12 +160,10 @@ function concurrentPool<A>(
}

function slide() {
if (finished && fillId === consumeId) {
while (consumer.length > 0) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const [resolve] = consumer.shift()!;
resolve({ value: undefined, done: true });
}
// After a failure, once nothing is in flight and the buffer is drained,
// remaining consumers can only ever receive `done`.
if (finished && idleWorkers === length && itemMap.size === 0) {
doneQueue();
} else {
processQueue();
}
Expand All @@ -174,6 +180,9 @@ function concurrentPool<A>(
const task: [Resolve<A>, Reject, number] = [resolve, reject, id];

consumer.push(task);
// Serve already-buffered results immediately - delivery must not
// depend on a future fetch settling.
deliverFromBuffer();
slide();
});
},
Expand Down
79 changes: 79 additions & 0 deletions test/Lazy/concurrentPool.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -115,3 +115,82 @@ describe("concurrentPool", function () {
expect(acc).toEqual([1, 2, 3]);
}, 2050);
});

describe("concurrentPool regressions", function () {
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));

it("should deliver buffered items even when the next upstream fetch hangs", async function () {
async function* fastThenHang() {
yield 1;
yield 2;
yield 3;
await new Promise(() => undefined); // upstream hangs on the 4th item
}
const iter = concurrentPool(3, fastThenHang())[Symbol.asyncIterator]();

expect((await iter.next()).value).toBe(1);
// let every pending pull-chain continuation drain
await sleep(50);
// items 2 and 3 are already buffered - they must not wait for a new fetch
const result = await Promise.race([
iter.next(),
sleep(300).then(() => "TIMEOUT" as const),
]);
expect(result).toEqual({ value: 2, done: false });
});

it("should keep prefetch bounded by the pool size with a slow consumer", async function () {
let produced = 0;
async function* fast() {
for (let i = 0; i < 1000; i++) {
produced++;
yield i;
}
}
const iter = concurrentPool(3, fast())[Symbol.asyncIterator]();
let consumed = 0;
for (let i = 0; i < 5; i++) {
await iter.next();
consumed++;
await sleep(15);
}
// in-flight + buffered may run ahead of the consumer by at most the pool size
expect(produced).toBeLessThanOrEqual(consumed + 3);
});

it("should not pull the source again after it has thrown", async function () {
let pullsAfterThrow = 0;
let thrown = false;
const source: AsyncIterable<number> = {
[Symbol.asyncIterator]() {
let n = 0;
return {
async next() {
if (thrown) {
pullsAfterThrow++;
return { value: undefined, done: true } as IteratorResult<number>;
}
if (n >= 1) {
thrown = true;
throw new Error("boom");
}
return { value: n++, done: false };
},
};
},
};
const iter = concurrentPool(2, source)[Symbol.asyncIterator]();
const results = [];
try {
for (let i = 0; i < 5; i++) {
results.push(await iter.next());
}
} catch {
// first rejection is expected
}
await sleep(30);
await iter.next().catch(() => undefined);
await sleep(30);
expect(pullsAfterThrow).toBe(0);
});
});
Loading