Skip to content
markpaper

src/transport/throttle.test.ts

v0.2.0 · 4.9 KB

Download file
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);
  });
});
All files