spb/datacenterindex
Public
HTML 53.9%
TypeScript 44.5%
JavaScript 0.6%
SQL 0.5%
1/**2 * Shared premium-credit budgets in Redis, so that daily caps hold across worker processes and restarts.3 *4 * Keys (UTC day `YYYY-MM-DD`, values are floats written with INCRBYFLOAT, 3-day TTL):5 * dci:budget:<provider>:<day> credits spent today per provider (scrapfly | firecrawl)6 * dci:budget:connector:<connectorId>:<day> credits spent today per connector (enforces fetch.maxCreditsPerDay)7 *8 * The API (`/api/admin/ops`) can read them with plain GET / MGET; `budgetKey()` / `connectorBudgetKey()` build the9 * names. `RedisBudgetStore` implements the connectors package's pluggable `BudgetStore`: `spend()` is fire-and-forget10 * (never blocks a fetch, never throws), `used()` returns the last value read by the periodic refresh (every11 * DCI_BUDGET_REFRESH_MS) — the in-process counter of the fetchers guarantees the value never lags below what this12 * process spent itself.13 */14import type { Redis } from "ioredis";15import { budgetDay, setBudgetStore, type BudgetStore, type PremiumProvider } from "@dci/connectors";1617export const BUDGET_KEY_PREFIX = "dci:budget";18export const BUDGET_TTL_SECONDS = 3 * 86_400;19export const PROVIDERS: readonly PremiumProvider[] = ["scrapfly", "firecrawl"];2021export function budgetKey(provider: PremiumProvider, day = budgetDay()): string { return `${BUDGET_KEY_PREFIX}:${provider}:${day}`; }22export function connectorBudgetKey(connectorId: string, day = budgetDay()): string { return `${BUDGET_KEY_PREFIX}:connector:${connectorId}:${day}`; }2324async function incr(redis: Redis, key: string, credits: number): Promise<number> {25 const res = (await redis.multi().incrbyfloat(key, credits).expire(key, BUDGET_TTL_SECONDS).exec()) ?? [];26 const first = res[0];27 if (!first) throw new Error("empty MULTI reply");28 if (first[0]) throw first[0];29 return Number(first[1]) || 0;30}3132export class RedisBudgetStore implements BudgetStore {33 private readonly cache = new Map<string, number>();34 private readonly connectorCache = new Map<string, { value: number; at: number }>();35 private timer: NodeJS.Timeout | null = null;36 private lastError: string | null = null;3738 constructor(private readonly redis: Redis, private readonly refreshMs = 15_000, private readonly log: (msg: string) => void = (m) => console.error(`[budget] ${m}`)) {}3940 used(provider: PremiumProvider, day: string): number | null {41 return this.cache.get(budgetKey(provider, day)) ?? null;42 }4344 spend(provider: PremiumProvider, day: string, credits: number): void {45 if (!(credits > 0)) return;46 const key = budgetKey(provider, day);47 this.cache.set(key, (this.cache.get(key) ?? 0) + credits); // optimistic, corrected by the next refresh48 incr(this.redis, key, credits).then((v) => this.cache.set(key, v), (e: Error) => this.fail(`spend ${key}: ${e.message}`));49 }5051 /** Re-read today's provider counters (also called by start()). Never throws. */52 async refresh(): Promise<void> {53 const day = budgetDay();54 const keys = PROVIDERS.map((p) => budgetKey(p, day));55 try {56 const vals = await this.redis.mget(...keys);57 keys.forEach((k, i) => this.cache.set(k, Number(vals[i] ?? 0) || 0));58 if (this.lastError) { this.log("redis budget counters reachable again"); this.lastError = null; }59 } catch (e) { this.fail(`refresh: ${(e as Error).message}`); }60 }6162 /** Credits a connector spent today (cached ≤ 10 s; 0 when Redis is unreachable — the global caps still apply). */63 async connectorUsed(connectorId: string, day = budgetDay()): Promise<number> {64 const key = connectorBudgetKey(connectorId, day);65 const c = this.connectorCache.get(key);66 if (c && Date.now() - c.at < 10_000) return c.value;67 try {68 const v = Number((await this.redis.get(key)) ?? 0) || 0;69 this.connectorCache.set(key, { value: v, at: Date.now() });70 return v;71 } catch (e) { this.fail(`connector ${connectorId}: ${(e as Error).message}`); return c?.value ?? 0; }72 }7374 /** Account premium credits to a connector (fire-and-forget). */75 spendConnector(connectorId: string, credits: number, day = budgetDay()): void {76 if (!(credits > 0)) return;77 const key = connectorBudgetKey(connectorId, day);78 const c = this.connectorCache.get(key);79 this.connectorCache.set(key, { value: (c?.value ?? 0) + credits, at: c?.at ?? 0 });80 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}`));81 }8283 /** Snapshot for /healthz: today's shared usage per provider. */84 snapshot(): Record<PremiumProvider, number> {85 const day = budgetDay();86 return Object.fromEntries(PROVIDERS.map((p) => [p, this.cache.get(budgetKey(p, day)) ?? 0])) as Record<PremiumProvider, number>;87 }8889 async start(): Promise<this> {90 await this.refresh();91 this.timer = setInterval(() => void this.refresh(), this.refreshMs);92 this.timer.unref();93 return this;94 }9596 stop(): void { if (this.timer) clearInterval(this.timer); this.timer = null; }9798 private fail(msg: string): void {99 // log once per outage, not once per fetch100 if (this.lastError !== msg.split(":")[0]) this.log(`${msg} (budgets fall back to in-process counters)`);101 this.lastError = msg.split(":")[0] ?? msg;102 }103}104105let installed: RedisBudgetStore | null = null;106/** Create, start and register the Redis store with the fetchers. Idempotent. */107export async function installBudgetStore(redis: Redis, refreshMs?: number): Promise<RedisBudgetStore> {108 if (installed) return installed;109 installed = await new RedisBudgetStore(redis, refreshMs).start();110 setBudgetStore(installed);111 return installed;112}113export function budgetStore(): RedisBudgetStore | null { return installed; }114export function uninstallBudgetStore(): void { installed?.stop(); installed = null; setBudgetStore(null); }115