/** Per-host token bucket + concurrency gate. In-process; the worker runs one process per node so this is sufficient. */ interface Bucket { tokens: number; last: number; rpm: number; inflight: number; concurrency: number; waiters: Array<() => void> } const buckets = new Map(); export function configureHost(host: string, rpm: number, concurrency: number): void { const b = buckets.get(host) ?? { tokens: rpm, last: Date.now(), rpm, inflight: 0, concurrency, waiters: [] }; b.rpm = rpm; b.concurrency = concurrency; buckets.set(host, b); } export async function acquire(host: string, minDelayMs = 0): Promise<() => void> { const b = buckets.get(host) ?? (buckets.set(host, { tokens: 30, last: Date.now(), rpm: 30, inflight: 0, concurrency: 2, waiters: [] }), buckets.get(host)!); for (;;) { const now = Date.now(); b.tokens = Math.min(b.rpm, b.tokens + ((now - b.last) / 60_000) * b.rpm); b.last = now; if (b.tokens >= 1 && b.inflight < b.concurrency) { b.tokens -= 1; b.inflight += 1; break; } const wait = b.tokens < 1 ? Math.ceil(((1 - b.tokens) / b.rpm) * 60_000) : 50; await new Promise((r) => setTimeout(r, Math.max(25, Math.min(wait, 5000), minDelayMs))); } if (minDelayMs > 0) await new Promise((r) => setTimeout(r, minDelayMs)); let released = false; return () => { if (!released) { released = true; b.inflight -= 1; } }; } export function hostStats(): Record { const out: Record = {}; for (const [h, b] of buckets) out[h] = { rpm: b.rpm, inflight: b.inflight, tokens: Math.round(b.tokens) }; return out; }