src/transport/limiter.test.ts
v0.3.0 · 11.8 KB
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import {
MIN_SAFE_BURST_CAPACITY,
createWeightLimiter,
getSharedWeightLimiter,
resetSharedWeightLimiters,
type Priority,
type WeightLimiterOptions,
} from './limiter.js';
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
resetSharedWeightLimiters();
});
const fakeNow = () => Date.now();
const TEST_LIMITER_OPTIONS: WeightLimiterOptions = {
weightPerMinute: 600,
burstCapacity: 80,
maxConcurrent: 3,
};
function makeLimiter(options: Partial<WeightLimiterOptions> = {}) {
return createWeightLimiter({ ...TEST_LIMITER_OPTIONS, ...options });
}
function deferred<T = void>() {
let resolve!: (v: T) => void;
const promise = new Promise<T>((r) => {
resolve = r;
});
return { promise, resolve };
}
describe('createWeightLimiter config', () => {
it('requires every workload limit explicitly', () => {
expect(() => (createWeightLimiter as unknown as () => unknown)()).toThrow(/requires explicit/);
const l = makeLimiter();
expect(l.weightPerMinute).toBe(600);
expect(l.burstCapacity).toBe(80);
expect(l.maxConcurrent).toBe(3);
expect(MIN_SAFE_BURST_CAPACITY).toBe(60);
});
it('validates options', () => {
expect(() => makeLimiter({ weightPerMinute: 0 })).toThrow(RangeError);
expect(() => makeLimiter({ burstCapacity: -1 })).toThrow(RangeError);
expect(() => makeLimiter({ maxConcurrent: 1.5 })).toThrow(RangeError);
expect(() => makeLimiter({ weightPerMinute: Number.NaN })).toThrow(RangeError);
});
});
describe('token bucket', () => {
it('admits a burst up to capacity instantly, then paces at the sustained rate', async () => {
// 600/min = 10 weight per second
const l = makeLimiter({ weightPerMinute: 600, burstCapacity: 80, now: fakeNow });
for (let i = 0; i < 4; i++) await l.acquire(20);
expect(l.stats().tokens).toBe(0);
let admitted = false;
void l.acquire(20).then(() => {
admitted = true;
});
await vi.advanceTimersByTimeAsync(1900);
expect(admitted).toBe(false);
await vi.advanceTimersByTimeAsync(150);
expect(admitted).toBe(true);
expect(l.stats().totalQueued).toBe(1);
});
it('rejects a weight above capacity immediately instead of hanging forever', async () => {
const l = makeLimiter({ burstCapacity: 80, now: fakeNow });
await expect(l.acquire(81)).rejects.toThrow(/chunk the request/);
await expect(l.acquire(-1)).rejects.toThrow(RangeError);
await expect(l.acquire(80)).resolves.toBeUndefined();
});
it('never accumulates more than capacity during idle time', async () => {
const l = makeLimiter({ weightPerMinute: 6000, burstCapacity: 50, now: fakeNow });
await vi.advanceTimersByTimeAsync(3_600_000);
expect(l.stats().tokens).toBe(50);
});
it('a clock going backwards does not mint tokens', async () => {
let t = 1_000_000;
const l = makeLimiter({ weightPerMinute: 60, burstCapacity: 10, now: () => t });
await l.acquire(10);
t -= 500_000;
expect(l.stats().tokens).toBe(0);
t += 1000; // 1 s later at 1 weight/s
expect(l.stats().tokens).toBe(1);
});
it('serves strictly by tier and FIFO within a tier (u1, h1, h2, n1, n2)', async () => {
// 1500/min = 25/s: one token every 40 ms
const l = makeLimiter({ weightPerMinute: 1500, burstCapacity: 1, now: fakeNow });
await l.acquire(1, { priority: 'normal' }); // eat the starting token
const order: string[] = [];
const push = (name: string, priority: Priority) => l.acquire(1, { priority }).then(() => order.push(name));
const all = Promise.all([push('n1', 'normal'), push('h1', 'high'), push('u1', 'urgent'), push('h2', 'high'), push('n2', 'normal')]);
await vi.advanceTimersByTimeAsync(400);
await all;
expect(order).toEqual(['u1', 'h1', 'h2', 'n1', 'n2']);
});
it('an urgent CLOSE queued after 4 high orders goes first', async () => {
const l = makeLimiter({ weightPerMinute: 1500, burstCapacity: 1, now: fakeNow });
await l.acquire(1);
const order: string[] = [];
const highs = [1, 2, 3, 4].map((i) => l.acquire(1, { priority: 'high' }).then(() => order.push(`h${i}`)));
await vi.advanceTimersByTimeAsync(5);
const close = l.acquire(1, { priority: 'urgent' }).then(() => order.push('close'));
await vi.advanceTimersByTimeAsync(400);
await Promise.all([...highs, close]);
expect(order[0]).toBe('close');
});
it('a heavy head-of-line waiter is not overtaken by lighter ones of the same tier', async () => {
const l = makeLimiter({ weightPerMinute: 600, burstCapacity: 20, now: fakeNow });
await l.acquire(15); // 5 tokens left
const order: string[] = [];
const heavy = l.acquire(20).then(() => order.push('heavy'));
const light = l.acquire(2).then(() => order.push('light'));
await vi.advanceTimersByTimeAsync(3000);
await Promise.all([heavy, light]);
expect(order).toEqual(['heavy', 'light']);
});
it('abort removes a waiter and unblocks the queue behind it', async () => {
const l = makeLimiter({ weightPerMinute: 60, burstCapacity: 20, now: fakeNow });
await l.acquire(20);
const ctl = new AbortController();
const heavy = l.acquire(20, { signal: ctl.signal });
let lightDone = false;
const light = l.acquire(1).then(() => {
lightDone = true;
});
await vi.advanceTimersByTimeAsync(1100);
expect(lightDone).toBe(false); // blocked behind heavy
ctl.abort(new Error('cancelled'));
await expect(heavy).rejects.toThrow('cancelled');
await vi.advanceTimersByTimeAsync(10);
await light;
expect(lightDone).toBe(true);
expect(l.stats().normalQueued).toBe(0);
});
it('rejects at once for an already aborted signal', async () => {
const l = makeLimiter({ now: fakeNow });
const ctl = new AbortController();
ctl.abort(new Error('nope'));
await expect(l.acquire(1, { signal: ctl.signal })).rejects.toThrow('nope');
});
it('charge() can overdraw the bucket and delays the next request', async () => {
const l = makeLimiter({ weightPerMinute: 600, burstCapacity: 80, now: fakeNow });
await l.acquire(20);
l.charge(100); // 60 - 100 = -40
expect(l.stats().tokens).toBe(-40);
let done = false;
void l.acquire(20).then(() => {
done = true;
});
await vi.advanceTimersByTimeAsync(5900);
expect(done).toBe(false);
await vi.advanceTimersByTimeAsync(200);
expect(done).toBe(true);
expect(l.stats().totalWeight).toBe(140);
});
it('charge() while a waiter sleeps reschedules its wake-up', async () => {
const l = makeLimiter({ weightPerMinute: 600, burstCapacity: 10, now: fakeNow });
await l.acquire(10);
let done = false;
void l.acquire(10).then(() => {
done = true;
});
l.charge(10); // now needs 2 s instead of 1 s
await vi.advanceTimersByTimeAsync(1500);
expect(done).toBe(false);
await vi.advanceTimersByTimeAsync(600);
expect(done).toBe(true);
l.charge(0);
l.charge(Number.NaN);
});
});
describe('concurrency semaphore', () => {
it('never runs more than maxConcurrent at once', async () => {
vi.useRealTimers();
const l = makeLimiter({ maxConcurrent: 3, burstCapacity: 1000, weightPerMinute: 60_000 });
let running = 0;
let peak = 0;
const tasks = Array.from({ length: 12 }, () =>
l.schedule(1, async () => {
running++;
peak = Math.max(peak, running);
await new Promise((r) => setTimeout(r, 2));
running--;
}),
);
await Promise.all(tasks);
expect(peak).toBe(3);
expect(l.stats().inFlight).toBe(0);
});
it('hands a released slot to the highest-priority waiter', async () => {
const l = makeLimiter({ maxConcurrent: 1, now: fakeNow });
const gate = deferred();
const order: string[] = [];
const first = l.withSlot(() => gate.promise);
const n = l.withSlot(async () => void order.push('normal'), { priority: 'normal' });
const h = l.withSlot(async () => void order.push('high'), { priority: 'high' });
const u = l.withSlot(async () => void order.push('urgent'), { priority: 'urgent' });
expect(l.stats().slotQueued).toBe(3);
gate.resolve();
await Promise.all([first, n, h, u]);
expect(order).toEqual(['urgent', 'high', 'normal']);
});
it('keeps at most maxConcurrent across a release boundary with a late arrival', async () => {
const l = makeLimiter({ maxConcurrent: 1, now: fakeNow });
const gate = deferred();
let running = 0;
let peak = 0;
const track = async () => {
running++;
peak = Math.max(peak, running);
await Promise.resolve();
running--;
};
const first = l.withSlot(async () => {
running++;
peak = Math.max(peak, running);
await gate.promise;
running--;
});
const queued = l.withSlot(track);
// Arrives right after the release; the slot was handed to the queued waiter atomically,
// so the late caller must queue instead of running concurrently.
const fresh = first.then(() => l.withSlot(track));
gate.resolve();
await Promise.all([first, queued, fresh]);
expect(peak).toBe(1);
});
it('releases the slot when the task throws, and a slot waiter can be aborted', async () => {
const l = makeLimiter({ maxConcurrent: 1, now: fakeNow });
const gate = deferred();
const first = l.withSlot(async () => {
await gate.promise;
throw new Error('task failed');
});
const ctl = new AbortController();
const aborted = l.withSlot(async () => 'never', { signal: ctl.signal });
ctl.abort(new Error('left the queue'));
await expect(aborted).rejects.toThrow('left the queue');
gate.resolve();
await expect(first).rejects.toThrow('task failed');
await expect(l.withSlot(async () => 'next')).resolves.toBe('next');
expect(l.stats().inFlight).toBe(0);
});
});
describe('shared limiter registry', () => {
it('returns one limiter per egress key', () => {
const a = getSharedWeightLimiter('default', TEST_LIMITER_OPTIONS);
const b = getSharedWeightLimiter('default');
const c = getSharedWeightLimiter('lane-b', TEST_LIMITER_OPTIONS);
expect(a).toBe(b);
expect(c).not.toBe(a);
});
it('refuses a silent split with different explicit options', () => {
getSharedWeightLimiter('x', TEST_LIMITER_OPTIONS);
expect(getSharedWeightLimiter('x', TEST_LIMITER_OPTIONS)).toBe(getSharedWeightLimiter('x'));
expect(() => getSharedWeightLimiter('x', { ...TEST_LIMITER_OPTIONS, burstCapacity: 240 })).toThrow(RangeError);
expect(() => getSharedWeightLimiter('missing')).toThrow(/requires explicit options/);
});
it('is stored on globalThis so duplicated package copies share it', () => {
const a = getSharedWeightLimiter('g', TEST_LIMITER_OPTIONS);
const map = (globalThis as unknown as Record<symbol, Map<string, { limiter: unknown }>>)[
Symbol.for('@markpaper/hl-kit/weight-limiters')
];
expect(map?.get('g')?.limiter).toBe(a);
resetSharedWeightLimiters('g');
expect(getSharedWeightLimiter('g', TEST_LIMITER_OPTIONS)).not.toBe(a);
});
});
describe('limiter boundaries', () => {
it('admits a weight exactly equal to capacity and a zero weight without waiting', async () => {
const l = makeLimiter({ burstCapacity: 80, now: () => 0 });
await l.acquire(80);
expect(l.stats().tokens).toBe(0);
await l.acquire(0);
await expect(l.acquire(80.0001)).rejects.toThrow(/chunk the request/);
await expect(l.acquire(-1)).rejects.toThrow(RangeError);
});
it('a configured bucket can hold the heaviest known /info request (userRole = 60)', async () => {
const l = makeLimiter({ now: () => 0 });
await expect(l.acquire(60)).resolves.toBeUndefined();
});
it('charge() ignores zero, negative and non-finite weights', () => {
const l = makeLimiter({ now: () => 0 });
l.charge(0);
l.charge(-5);
l.charge(Number.NaN);
expect(l.stats().tokens).toBe(80);
expect(l.stats().totalWeight).toBe(0);
});
});