/** * Shared premium-credit budgets in Redis, so that daily caps hold across worker processes and restarts. * * Keys (UTC day `YYYY-MM-DD`, values are floats written with INCRBYFLOAT, 3-day TTL): * dci:budget:: credits spent today per provider (scrapfly | firecrawl) * dci:budget:connector:: credits spent today per connector (enforces fetch.maxCreditsPerDay) * * The API (`/api/admin/ops`) can read them with plain GET / MGET; `budgetKey()` / `connectorBudgetKey()` build the * names. `RedisBudgetStore` implements the connectors package's pluggable `BudgetStore`: `spend()` is fire-and-forget * (never blocks a fetch, never throws), `used()` returns the last value read by the periodic refresh (every * DCI_BUDGET_REFRESH_MS) — the in-process counter of the fetchers guarantees the value never lags below what this * process spent itself. */ import type { Redis } from "ioredis"; import { budgetDay, setBudgetStore, type BudgetStore, type PremiumProvider } from "@dci/connectors"; export const BUDGET_KEY_PREFIX = "dci:budget"; export const BUDGET_TTL_SECONDS = 3 * 86_400; export const PROVIDERS: readonly PremiumProvider[] = ["scrapfly", "firecrawl"]; export function budgetKey(provider: PremiumProvider, day = budgetDay()): string { return `${BUDGET_KEY_PREFIX}:${provider}:${day}`; } export function connectorBudgetKey(connectorId: string, day = budgetDay()): string { return `${BUDGET_KEY_PREFIX}:connector:${connectorId}:${day}`; } async function incr(redis: Redis, key: string, credits: number): Promise { const res = (await redis.multi().incrbyfloat(key, credits).expire(key, BUDGET_TTL_SECONDS).exec()) ?? []; const first = res[0]; if (!first) throw new Error("empty MULTI reply"); if (first[0]) throw first[0]; return Number(first[1]) || 0; } export class RedisBudgetStore implements BudgetStore { private readonly cache = new Map(); private readonly connectorCache = new Map(); private timer: NodeJS.Timeout | null = null; private lastError: string | null = null; constructor(private readonly redis: Redis, private readonly refreshMs = 15_000, private readonly log: (msg: string) => void = (m) => console.error(`[budget] ${m}`)) {} used(provider: PremiumProvider, day: string): number | null { return this.cache.get(budgetKey(provider, day)) ?? null; } spend(provider: PremiumProvider, day: string, credits: number): void { if (!(credits > 0)) return; const key = budgetKey(provider, day); this.cache.set(key, (this.cache.get(key) ?? 0) + credits); // optimistic, corrected by the next refresh incr(this.redis, key, credits).then((v) => this.cache.set(key, v), (e: Error) => this.fail(`spend ${key}: ${e.message}`)); } /** Re-read today's provider counters (also called by start()). Never throws. */ async refresh(): Promise { const day = budgetDay(); const keys = PROVIDERS.map((p) => budgetKey(p, day)); try { const vals = await this.redis.mget(...keys); keys.forEach((k, i) => this.cache.set(k, Number(vals[i] ?? 0) || 0)); if (this.lastError) { this.log("redis budget counters reachable again"); this.lastError = null; } } catch (e) { this.fail(`refresh: ${(e as Error).message}`); } } /** Credits a connector spent today (cached ≤ 10 s; 0 when Redis is unreachable — the global caps still apply). */ async connectorUsed(connectorId: string, day = budgetDay()): Promise { const key = connectorBudgetKey(connectorId, day); const c = this.connectorCache.get(key); if (c && Date.now() - c.at < 10_000) return c.value; try { const v = Number((await this.redis.get(key)) ?? 0) || 0; this.connectorCache.set(key, { value: v, at: Date.now() }); return v; } catch (e) { this.fail(`connector ${connectorId}: ${(e as Error).message}`); return c?.value ?? 0; } } /** Account premium credits to a connector (fire-and-forget). */ spendConnector(connectorId: string, credits: number, day = budgetDay()): void { if (!(credits > 0)) return; const key = connectorBudgetKey(connectorId, day); const c = this.connectorCache.get(key); this.connectorCache.set(key, { value: (c?.value ?? 0) + credits, at: c?.at ?? 0 }); incr(this.redis, key, credits).then((v) => this.connectorCache.set(key, { value: v, at: Date.now() }), (e: Error) => this.fail(`spend connector ${key}: ${e.message}`)); } /** Snapshot for /healthz: today's shared usage per provider. */ snapshot(): Record { const day = budgetDay(); return Object.fromEntries(PROVIDERS.map((p) => [p, this.cache.get(budgetKey(p, day)) ?? 0])) as Record; } async start(): Promise { await this.refresh(); this.timer = setInterval(() => void this.refresh(), this.refreshMs); this.timer.unref(); return this; } stop(): void { if (this.timer) clearInterval(this.timer); this.timer = null; } private fail(msg: string): void { // log once per outage, not once per fetch if (this.lastError !== msg.split(":")[0]) this.log(`${msg} (budgets fall back to in-process counters)`); this.lastError = msg.split(":")[0] ?? msg; } } let installed: RedisBudgetStore | null = null; /** Create, start and register the Redis store with the fetchers. Idempotent. */ export async function installBudgetStore(redis: Redis, refreshMs?: number): Promise { if (installed) return installed; installed = await new RedisBudgetStore(redis, refreshMs).start(); setBudgetStore(installed); return installed; } export function budgetStore(): RedisBudgetStore | null { return installed; } export function uninstallBudgetStore(): void { installed?.stop(); installed = null; setBudgetStore(null); }