Skip to content
markpaper

src/transport/limiter.test.ts

v0.3.0 · 11.8 KB

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