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); }); });