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%
5.7 KB · 115 lines typescript
Raw Blame History
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