SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
3 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%
1.7 KB · 31 lines typescript
Raw Blame History
1/** Per-host token bucket + concurrency gate. In-process; the worker runs one process per node so this is sufficient. */2interface Bucket { tokens: number; last: number; rpm: number; inflight: number; concurrency: number; waiters: Array<() => void> }3const buckets = new Map<string, Bucket>();45export function configureHost(host: string, rpm: number, concurrency: number): void {6  const b = buckets.get(host) ?? { tokens: rpm, last: Date.now(), rpm, inflight: 0, concurrency, waiters: [] };7  b.rpm = rpm; b.concurrency = concurrency;8  buckets.set(host, b);9}1011export async function acquire(host: string, minDelayMs = 0): Promise<() => void> {12  const b = buckets.get(host) ?? (buckets.set(host, { tokens: 30, last: Date.now(), rpm: 30, inflight: 0, concurrency: 2, waiters: [] }), buckets.get(host)!);13  for (;;) {14    const now = Date.now();15    b.tokens = Math.min(b.rpm, b.tokens + ((now - b.last) / 60_000) * b.rpm);16    b.last = now;17    if (b.tokens >= 1 && b.inflight < b.concurrency) { b.tokens -= 1; b.inflight += 1; break; }18    const wait = b.tokens < 1 ? Math.ceil(((1 - b.tokens) / b.rpm) * 60_000) : 50;19    await new Promise((r) => setTimeout(r, Math.max(25, Math.min(wait, 5000), minDelayMs)));20  }21  if (minDelayMs > 0) await new Promise((r) => setTimeout(r, minDelayMs));22  let released = false;23  return () => { if (!released) { released = true; b.inflight -= 1; } };24}2526export function hostStats(): Record<string, { rpm: number; inflight: number; tokens: number }> {27  const out: Record<string, { rpm: number; inflight: number; tokens: number }> = {};28  for (const [h, b] of buckets) out[h] = { rpm: b.rpm, inflight: b.inflight, tokens: Math.round(b.tokens) };29  return out;30}31