40 lines
1.5 KiB
JavaScript
40 lines
1.5 KiB
JavaScript
import { describe, expect, test } from 'bun:test';
|
|
import { createWarmQueue } from './warm-queue.js';
|
|
|
|
describe('stream warm queue', () => {
|
|
test('runs at the concurrency limit and drains queued requests FIFO', async () => {
|
|
const started = [];
|
|
const releases = new Map();
|
|
const q = createWarmQueue({ concurrency: 2, maxQueued: 2, run: (id) => {
|
|
started.push(id);
|
|
return new Promise((resolve) => releases.set(id, resolve));
|
|
} });
|
|
|
|
q.enqueue('a'); q.enqueue('b'); q.enqueue('c');
|
|
expect(q.pending).toBe(3);
|
|
await Promise.resolve();
|
|
expect(started).toEqual(['a', 'b']);
|
|
releases.get('a')();
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
expect(started).toEqual(['a', 'b', 'c']);
|
|
releases.get('b')(); releases.get('c')();
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
expect(q.pending).toBe(0);
|
|
});
|
|
|
|
test('deduplicates in-flight ids and rejects work past the queue bound', async () => {
|
|
const releases = [];
|
|
const q = createWarmQueue({ concurrency: 1, maxQueued: 1, run: () => new Promise((resolve) => releases.push(resolve)) });
|
|
expect(q.enqueue('a')).toBe(true);
|
|
await Promise.resolve();
|
|
expect(q.enqueue('a')).toBe(true);
|
|
expect(q.enqueue('b')).toBe(true);
|
|
expect(q.enqueue('c')).toBe(false);
|
|
releases[0]();
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
releases[1]();
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
expect(q.pending).toBe(0);
|
|
});
|
|
});
|