import test from 'node:test'; import assert from 'node:assert/strict'; import { OverloadError, createInferenceQueue } from '../src/inference-queue.js'; function deferred() { let resolve; let reject; const promise = new Promise((res, rej) => { resolve = res; reject = rej; }); return { promise, resolve, reject }; } test('a queued request past its deadline is expired deterministically and its place is reusable', async () => { let clock = 0; const queue = createInferenceQueue({ concurrency: 1, maxQueueDepth: 1, requestTimeoutMs: 100, now: () => clock }); const wedgeGate = deferred(); const active = queue.run(signal => { return new Promise((resolve, reject) => { signal.addEventListener('abort', () => reject(new Error('upstream aborted'))); wedgeGate.promise.then(resolve, reject); }); }); active.catch(() => {}); await Promise.resolve(); let started = false; const queued = queue.run(() => { started = true; return 'never'; }); await Promise.resolve(); clock += 150; // deadline (submitted at t=0, limit 100) has elapsed const probe = queue.run(() => 'ok'); // queue activity triggers the deadline sweep await assert.rejects(() => queued, error => { assert.ok(error instanceof OverloadError); assert.match(error.message, /timed out|busy/i); return true; }); assert.equal(started, false, 'the expired request must never reach the provider'); // Capacity is intact: cancelling the wedged holder lets the surviving request through. queue.cancel(active, 'cleanup'); assert.equal(await Promise.race([probe, new Promise((_, reject) => setTimeout(() => reject(new Error('probe starved')), 200))]), 'ok'); }); test('client disconnect cancels a queued request before it reaches the provider', async () => { let clock = 0; const queue = createInferenceQueue({ concurrency: 1, maxQueueDepth: 2, requestTimeoutMs: 60_000, now: () => clock }); const wedgeGate = deferred(); const active = queue.run(() => wedgeGate.promise); active.catch(() => {}); await Promise.resolve(); let sawStart = false; const queued = queue.run(() => { sawStart = true; return 'too late'; }); await Promise.resolve(); queue.cancel(queued, 'client disconnected'); await assert.rejects(() => queued, error => error.message === 'client disconnected'); assert.equal(sawStart, false, 'cancelled request never invoked its task'); const admitted = queue.run(() => 'admitted-after-cancel'); wedgeGate.resolve('wedge-done'); assert.equal(await admitted, 'admitted-after-cancel'); assert.equal(await active, 'wedge-done'); }); test('cancelling an already-running task aborts its signal so buffers can be released', async () => { let clock = 0; const queue = createInferenceQueue({ concurrency: 1, maxQueueDepth: 1, requestTimeoutMs: 60_000, now: () => clock }); const observed = []; const running = queue.run(async signal => { observed.push(signal); signal.addEventListener('abort', () => observed.push(`aborted:${signal.reason?.message ?? signal.reason}`)); await new Promise((resolve, reject) => { signal.addEventListener('abort', () => reject(new Error('task stopped'))); setTimeout(resolve, 5_000); }); }); // Let the wrapper admit the request and actually invoke the task body. await Promise.resolve(); await Promise.resolve(); assert.equal(observed.length, 1, 'task is genuinely running before cancellation'); queue.cancel(running, 'stop'); await assert.rejects(() => running, /task stopped/); assert.ok(observed[0] instanceof AbortSignal, 'task received a real AbortSignal'); assert.match(String(observed[1]), /stop/, 'running task observed the abort with its reason'); // The slot was released by the cancelled task, so the next request starts immediately. const next = queue.run(() => 'next-runs'); assert.equal(await next, 'next-runs'); });