File size: 10,261 Bytes
20f83d9
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
// Per-account REST rate-limit layer for #3199 (Phase 1). Two limit types:
//
//   1. Per-minute burst  β€” infra/abuse protection. Hard limit; the gateway
//      returns 429 on violation (in enforce mode). Value from the catalog
//      `apiRateLimit` (user keys) or a hardcoded constant (enterprise).
//   2. Daily usage meter β€” the commercial "included allowance". Hard-rejects
//      (429 in enforce mode) at the sold allowance (#4635; was a 10Γ— ceiling).
//
// This module is decision-only: it never builds a Response and never reads the
// enforce flag β€” the gateway (the single chokepoint at server/gateway.ts:1034)
// owns enforce-vs-shadow, per-IP bypass, and Response construction. That keeps
// the burst/meter math unit-testable in isolation (inject the pipeline + date;
// stub this module's decisions in the gateway-wiring test).
//
// Patterns cloned: api/mcp/quota.ts (INCR-first meter + DECR rollback),
// api/_rate-limit.js (lazy Upstash singleton, NODE_TEST_CONTEXT retry skip,
// X-RateLimit-* header shape), server/_shared/pro-mcp-token.ts UTC helpers.

import { Ratelimit } from '@upstash/ratelimit';
import { Redis } from '@upstash/redis';

import { getKeyPrefix } from './redis';
import { secondsUntilUtcMidnight } from './pro-mcp-token';

/** Hardcoded per-minute burst for enterprise env keys β€” they carry no Convex
 *  entitlement (gateway.ts:1006-1007 skips checkEntitlement), so this cannot
 *  be sourced from `features.apiRateLimit`. Mirrors ENTERPRISE_FEATURES. */
export const ENTERPRISE_API_RATE_LIMIT = 1000;


// One Redis client shared across every per-minute Ratelimit instance; one
// Ratelimit per distinct numeric limit (60, 300, 1000) cached in the Map so two
// Starter accounts share a limiter *config* but get separate buckets via the
// per-account identifier passed to `.limit()`.
let redisSingleton: Redis | null = null;
const burstLimiters = new Map<number, Ratelimit>();

function getRedis(): Redis | null {
  if (redisSingleton) return redisSingleton;
  const url = process.env.UPSTASH_REDIS_REST_URL;
  const token = process.env.UPSTASH_REDIS_REST_TOKEN;
  if (!url || !token) return null;
  // Skip the @upstash/redis retry backoff under the node test runner so
  // fail-open tests pointed at a fake host degrade immediately; production
  // (env unset) keeps the resilient default. Mirrors api/_rate-limit.js.
  // `retry: false` must stay a literal (not a spread) or it widens to
  // `boolean` and fails tsconfig.api.json's RetryConfig type.
  redisSingleton = process.env.NODE_TEST_CONTEXT
    ? new Redis({ url, token, retry: false })
    : new Redis({ url, token });
  return redisSingleton;
}

/**
 * The per-minute burst limiter for `perMinute` requests / 60s, cached by limit.
 * Returns null when Upstash is not configured (caller fail-opens).
 */
export function getBurstLimiter(perMinute: number): Ratelimit | null {
  const existing = burstLimiters.get(perMinute);
  if (existing) return existing;
  const redis = getRedis();
  if (!redis) return null;
  const limiter = new Ratelimit({
    redis,
    limiter: Ratelimit.slidingWindow(perMinute, '60 s'),
    // Env-scope the prefix exactly like the daily meter (runRedisPipeline's
    // prefixKey) so a preview deployment sharing one Upstash database doesn't
    // consume/pollute the production burst namespace. Empty in production.
    prefix: `${getKeyPrefix()}rl:apikey:min`,
    analytics: false,
  });
  burstLimiters.set(perMinute, limiter);
  return limiter;
}

export type BurstDecision =
  | { ok: true }
  | { ok: false; limit: number; reset: number };

/**
 * Evaluate the per-minute burst window for `identity`. Fail-OPEN: a missing
 * Upstash config or any Redis error resolves to `{ ok: true }` so a paying
 * customer is never 429'd for our outage (mirrors api/_rate-limit.js).
 */
export async function checkBurst(perMinute: number, identity: string): Promise<BurstDecision> {
  const limiter = getBurstLimiter(perMinute);
  if (!limiter) return { ok: true };
  try {
    const { success, limit, reset } = await limiter.limit(identity);
    if (!success) return { ok: false, limit, reset };
    return { ok: true };
  } catch {
    return { ok: true };
  }
}

/** Plain (un-prefixed) daily-meter key β€” `runRedisPipeline` applies the
 *  deployment/env prefix. UTC calendar day so the daily meter resets at midnight.
 *  `date` is injectable for deterministic tests. */
export function apiKeyDailyKey(userId: string, date?: Date): string {
  if (!userId) return '';
  const d = date ?? new Date();
  const yyyy = d.getUTCFullYear();
  const mm = String(d.getUTCMonth() + 1).padStart(2, '0');
  const dd = String(d.getUTCDate()).padStart(2, '0');
  return `rl:apikey:day:${userId}:${yyyy}-${mm}-${dd}`;
}

/** 48h TTL: covers UTC-midnight rollover + an inspection window. Mirrors
 *  PRO_DAILY_QUOTA_TTL_SECONDS. */
export const API_DAILY_TTL_SECONDS = 172_800;

/** Minimal pipeline contract β€” the subset of `runRedisPipeline` this module
 *  needs. The gateway passes `(cmds) => runRedisPipeline(cmds)`; tests inject a
 *  mock. Returns `[]` on failure (fail-open), never throws by contract. */
export type RateLimitPipeline = (
  commands: Array<Array<string | number>>,
) => Promise<Array<{ result?: unknown }>>;

export interface MeterResult {
  /** Post-INCR count for this UTC day (0 when not metered). */
  count: number;
  /** True when count exceeded the sold daily allowance. */
  overLimit: boolean;
  /** False when Redis was unavailable (fail-open: serve uncounted). */
  metered: boolean;
  /** Seconds until UTC midnight β€” the daily 429 `Retry-After`. */
  retryAfterSec: number;
  /** Idempotent DECR rollback. The gateway calls this only when it actually
   *  rejects (enforce + overLimit); in shadow the request is served, so the
   *  increment stands and reflects true demand. */
  rollback: () => Promise<void>;
}

/**
 * Increment the per-account daily meter and report whether the sold daily
 * allowance is now exceeded. INCR-first (atomic; no check-then-incr race), mirroring
 * api/mcp/quota.ts::reserveQuota.
 *
 * - `allowance < 0` (unlimited, e.g. enterprise) β†’ no Redis call; never metered.
 * - Redis unavailable / pipeline failure β†’ `metered:false`, `overLimit:false`
 *   (fail-open: the gateway serves uncounted).
 */
export async function reserveDailyMeter(opts: {
  userId: string;
  allowance: number;
  pipeline: RateLimitPipeline;
  date?: Date;
}): Promise<MeterResult> {
  const { userId, allowance, pipeline, date } = opts;
  const noop = async (): Promise<void> => {};
  const retryAfterSec = secondsUntilUtcMidnight(date);

  // No daily limit: `-1` is unlimited (enterprise); `0` is a misconfiguration
  // (positive burst but zero allowance) that we fail OPEN on rather than 429
  // every request (a 0 allowance would reject request #1). This guard also keeps
  // unlimited (-1) from ever reaching `count > allowance` (always-true otherwise).
  // Callers already gate eligibility on apiRateLimit > 0, so this is defensive.
  if (allowance <= 0) {
    return { count: 0, overLimit: false, metered: false, retryAfterSec, rollback: noop };
  }

  const key = apiKeyDailyKey(userId, date);
  if (!key) {
    return { count: 0, overLimit: false, metered: false, retryAfterSec, rollback: noop };
  }

  let pipeResult: Array<{ result?: unknown }> | null;
  try {
    pipeResult = await pipeline([
      ['INCR', key],
      ['EXPIRE', key, API_DAILY_TTL_SECONDS],
    ]);
  } catch {
    pipeResult = null;
  }

  // Fail-open: couldn't meter β†’ serve uncounted (never punish a paying
  // customer for our Redis outage).
  if (!pipeResult || !Array.isArray(pipeResult) || pipeResult.length === 0) {
    return { count: 0, overLimit: false, metered: false, retryAfterSec, rollback: noop };
  }

  const incrRaw = pipeResult[0]?.result;
  const count = typeof incrRaw === 'number' ? incrRaw : Number(incrRaw);
  if (!Number.isFinite(count) || count < 1) {
    return { count: 0, overLimit: false, metered: false, retryAfterSec, rollback: noop };
  }

  let rolledBack = false;
  const rollback = async (): Promise<void> => {
    if (rolledBack) return;
    rolledBack = true;
    try {
      await pipeline([['DECR', key]]);
    } catch {
      // Best-effort: a failed DECR overshoots the meter by 1, the
      // cost-protection-correct direction.
    }
  };

  // Enforce at the SOLD allowance (#4635): the customer's plan limit is the
  // limit. (Was a 10Γ— safety ceiling; dropped so the sold cap is authoritative.)
  return { count, overLimit: count > allowance, metered: true, retryAfterSec, rollback };
}

/**
 * Standard rate-limit response headers for a 429. Emits the IETF RateLimit
 * fields (draft-ietf-httpapi-ratelimit-headers) β€” RateLimit-Policy advertises
 * the quota/window, the combined RateLimit member carries live remaining +
 * delta-seconds reset β€” alongside the legacy X-RateLimit-* set for back-compat,
 * so customers get a uniform self-throttle contract across the per-IP and
 * per-account limiters. Mirrors api/_rate-limit.js. The gateway merges these
 * with corsHeaders.
 *
 * `resetMs` is a Unix epoch in MILLISECONDS; the IETF reset (`t` /
 * RateLimit-Reset) is delta-SECONDS, so it is derived here. `windowSec` is the
 * policy window in seconds (defaults to the 60 s burst window).
 */
export function rateLimitHeaders(opts: {
  limit: number;
  remaining: number;
  resetMs: number;
  retryAfterSec: number;
  windowSec?: number;
}): Record<string, string> {
  const remaining = Math.max(0, opts.remaining);
  const resetSeconds = Math.max(0, Math.ceil((opts.resetMs - Date.now()) / 1000));
  const windowSec = opts.windowSec ?? 60;
  return {
    // IETF RateLimit fields.
    'RateLimit-Policy': `"default";q=${opts.limit};w=${windowSec}`,
    'RateLimit-Limit': String(opts.limit),
    'RateLimit-Remaining': String(remaining),
    'RateLimit-Reset': String(resetSeconds),
    RateLimit: `"default";r=${remaining};t=${resetSeconds}`,
    // Legacy X-RateLimit-* retained for back-compat (Reset is epoch-ms).
    'X-RateLimit-Limit': String(opts.limit),
    'X-RateLimit-Remaining': String(remaining),
    'X-RateLimit-Reset': String(opts.resetMs),
    'Retry-After': String(Math.max(1, opts.retryAfterSec)),
  };
}