src/transport/throttle.test.ts
v0.2.0 · 4.9 KB
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import { createWeightThrottle, type WeightThrottleOptions } from './throttle.js';
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
});
const fakeNow = () => Date.now();
const TEST_CONFIG: WeightThrottleOptions = {
queriesPerMinute: 600,
queryBurst: 20,
executesPerMinute: 300,
executeBurst: 10,
maxConcurrent: 2,
};
const config = (overrides: Partial<WeightThrottleOptions> = {}): WeightThrottleOptions => ({
...TEST_CONFIG,
...overrides,
});
function deferred<T = void>() {
let resolve!: (v: T) => void;
const promise = new Promise<T>((r) => {
resolve = r;
});
return { promise, resolve };
}
describe('createWeightThrottle', () => {
it('requires all workload limits explicitly', () => {
expect(() => (createWeightThrottle as unknown as () => unknown)()).toThrow(/requires explicit/);
const t = createWeightThrottle(config());
expect(t.stats()).toMatchObject({
queries: { capacity: 20, perMinute: 600 },
executes: { capacity: 10, perMinute: 300 },
inFlight: 0,
});
expect(() => createWeightThrottle(config({ queryBurst: 0 }))).toThrow(RangeError);
expect(() => createWeightThrottle(config({ maxConcurrent: 1.5 }))).toThrow(RangeError);
});
it('query and execute buckets are independent', async () => {
const t = createWeightThrottle(config({
queriesPerMinute: 600,
queryBurst: 20,
executesPerMinute: 600,
executeBurst: 20,
now: fakeNow,
}));
await t.query(20, async () => 1);
expect(t.stats().queries.tokens).toBe(0);
expect(t.stats().executes.tokens).toBe(20);
await t.execute(20, async () => 1);
expect(t.stats().executes.tokens).toBe(0);
});
it('paces at the sustained rate once the burst is spent', async () => {
const t = createWeightThrottle(config({ queriesPerMinute: 600, queryBurst: 20, now: fakeNow })); // 10 weight/s
await t.query(20, async () => 1);
let done = false;
void t
.query(10, async () => 1)
.then(() => {
done = true;
});
await vi.advanceTimersByTimeAsync(900);
expect(done).toBe(false);
await vi.advanceTimersByTimeAsync(150);
expect(done).toBe(true);
expect(t.stats().queries.totalQueued).toBe(1);
});
it('rejects a weight above the bucket capacity immediately instead of hanging (orders for 200 products)', async () => {
const t = createWeightThrottle(config({ queryBurst: 300, now: fakeNow }));
await expect(t.query(400, async () => 1)).rejects.toThrow(/chunk the request/);
await expect(t.query(-1, async () => 1)).rejects.toThrow(RangeError);
await expect(t.query(300, async () => 'ok')).resolves.toBe('ok');
});
it('executes take a concurrency slot before queued queries', async () => {
const t = createWeightThrottle(config({ maxConcurrent: 1, now: fakeNow }));
const gate = deferred();
const order: string[] = [];
const first = t.query(1, async () => {
await gate.promise;
order.push('first');
});
await vi.advanceTimersByTimeAsync(0);
const q = t.query(1, async () => {
order.push('query');
});
const e = t.execute(1, async () => {
order.push('execute');
});
await vi.advanceTimersByTimeAsync(0);
expect(t.stats()).toMatchObject({ inFlight: 1, slotQueued: 2 });
gate.resolve();
await Promise.all([first, q, e]);
expect(order).toEqual(['first', 'execute', 'query']);
expect(t.stats().inFlight).toBe(0);
});
it('abort removes a queued waiter', async () => {
const t = createWeightThrottle(config({ queriesPerMinute: 60, queryBurst: 1, now: fakeNow }));
await t.query(1, async () => 1);
const ctl = new AbortController();
const p = t.query(1, async () => 1, { signal: ctl.signal });
ctl.abort(new Error('stop'));
await expect(p).rejects.toThrow('stop');
expect(t.stats().queries.queued).toBe(0);
await expect(t.query(1, async () => 1, { signal: ctl.signal })).rejects.toThrow('stop');
});
it('counts calls per label and drains them', async () => {
const t = createWeightThrottle(config({ now: fakeNow }));
await t.query(1, async () => 1, { label: 'orders' });
await t.query(1, async () => 1, { label: 'orders' });
await t.execute(1, async () => 1, { label: 'order:close' });
await t.execute(1, async () => 1);
expect(t.drainCounts()).toEqual([
{ label: 'orders', n: 2 },
{ label: 'order:close', n: 1 },
{ label: 'execute', n: 1 },
]);
expect(t.drainCounts()).toEqual([]);
});
it('releases the slot when fn throws', async () => {
const t = createWeightThrottle(config({ maxConcurrent: 1, now: fakeNow }));
await expect(
t.query(1, async () => {
throw new Error('x');
}),
).rejects.toThrow('x');
expect(t.stats().inFlight).toBe(0);
await expect(t.query(1, async () => 2)).resolves.toBe(2);
});
});