178 lines
6.2 KiB
TypeScript
178 lines
6.2 KiB
TypeScript
import { planWriteQuota, type QuotaOptions } from './quota.js'
|
||
|
||
/**
|
||
* Sliding-window write limiter, decoupled from both the HTTP framework and the
|
||
* counter store.
|
||
*
|
||
* WHY THE STORE IS PLUGGABLE: the original implementation held buckets in a
|
||
* module-level Map. That is correct on ONE instance and silently wrong on more
|
||
* than one — each instance keeps its own counters, so N instances grant N×
|
||
* the quota while every dashboard still reports the configured limit. The
|
||
* limiter does not fail closed, it fails *quietly generous*, which is the
|
||
* hardest kind of limit failure to notice.
|
||
*
|
||
* MemoryStore is still the default because it is right for a single instance
|
||
* and needs no infrastructure. Passing a shared store (Redis, Postgres) is what
|
||
* makes the limit true across a fleet.
|
||
*/
|
||
|
||
export interface LimiterStore {
|
||
/**
|
||
* Atomically increment the counter for `key`, creating it with the given TTL
|
||
* if absent. Returns the new count and when the window resets (epoch ms).
|
||
*
|
||
* MUST be atomic in a shared store: a read-then-write pair races under
|
||
* concurrency and undercounts exactly when the limiter matters most.
|
||
*/
|
||
hit: (key: string, windowMs: number) => { count: number, resetAt: number } | Promise<{ count: number, resetAt: number }>
|
||
}
|
||
|
||
interface Bucket { count: number, resetAt: number }
|
||
|
||
/** In-process store. Correct for a single instance only — see above. */
|
||
export class MemoryStore implements LimiterStore {
|
||
private buckets = new Map<string, Bucket>()
|
||
private gc: ReturnType<typeof setInterval> | null = null
|
||
|
||
hit(key: string, windowMs: number): { count: number, resetAt: number } {
|
||
this.startGC()
|
||
const now = Date.now()
|
||
const b = this.buckets.get(key)
|
||
if (!b || b.resetAt < now) {
|
||
const fresh = { count: 1, resetAt: now + windowMs }
|
||
this.buckets.set(key, fresh)
|
||
return fresh
|
||
}
|
||
b.count++
|
||
return b
|
||
}
|
||
|
||
private startGC() {
|
||
if (this.gc) return
|
||
this.gc = setInterval(() => {
|
||
const now = Date.now()
|
||
for (const [k, v] of this.buckets) if (v.resetAt < now) this.buckets.delete(k)
|
||
}, 5 * 60 * 1000)
|
||
// Never hold the process open just to expire counters.
|
||
this.gc.unref?.()
|
||
}
|
||
|
||
/** Test helper — drop all counters. */
|
||
reset() { this.buckets.clear() }
|
||
}
|
||
|
||
export interface LimitInput {
|
||
product: string
|
||
/** The account this budget belongs to. Empty means unauthenticated. */
|
||
userId: string
|
||
plan: string | null | undefined
|
||
windowSec?: number
|
||
}
|
||
|
||
export interface LimitDecision {
|
||
/** False when the caller should be rejected. */
|
||
allowed: boolean
|
||
/** Configured ceiling. Infinity for unlimited plans. */
|
||
limit: number
|
||
/** Requests left in this window. Never negative. */
|
||
remaining: number
|
||
/** Seconds until the window resets. 0 when unlimited. */
|
||
retryAfter: number
|
||
/** Ready-made message for a 429 body. */
|
||
message?: string
|
||
}
|
||
|
||
const UNLIMITED: LimitDecision = {
|
||
allowed: true, limit: Infinity, remaining: Infinity, retryAfter: 0,
|
||
}
|
||
|
||
/** Store whose `hit` is synchronous. MemoryStore satisfies this. */
|
||
export interface SyncLimiterStore {
|
||
hit: (key: string, windowMs: number) => { count: number, resetAt: number }
|
||
}
|
||
|
||
export interface LimiterOptions extends QuotaOptions {
|
||
store?: LimiterStore
|
||
}
|
||
|
||
export interface SyncLimiterOptions extends QuotaOptions {
|
||
store?: SyncLimiterStore
|
||
}
|
||
|
||
/** Shared decision logic. Pure — no I/O, no store, no framework. */
|
||
function decide(max: number, count: number, resetAt: number): LimitDecision {
|
||
const remaining = Math.max(0, max - count)
|
||
if (count > max) {
|
||
const retryAfter = Math.max(1, Math.ceil((resetAt - Date.now()) / 1000))
|
||
return {
|
||
allowed: false,
|
||
limit: max,
|
||
remaining: 0,
|
||
retryAfter,
|
||
message: `You've hit your plan's limit of ${max} changes/minute. Try again in ${retryAfter}s, or upgrade for a higher limit.`,
|
||
}
|
||
}
|
||
return { allowed: true, limit: max, remaining, retryAfter: 0 }
|
||
}
|
||
|
||
/**
|
||
* True when this call needs no counting at all: an unlimited plan, or no
|
||
* principal to bill the write to. Unauthenticated writes are gated by auth,
|
||
* not by quota — counting them under a shared empty key would let one
|
||
* anonymous caller exhaust the budget for every other one.
|
||
*/
|
||
function isExempt(max: number, userId: string): boolean {
|
||
return !max || max <= 0 || !userId
|
||
}
|
||
|
||
function bucketKey(input: LimitInput): string {
|
||
return `w:${input.product}:${input.userId}`
|
||
}
|
||
|
||
/**
|
||
* Build an async limiter. Returns a DECISION rather than throwing, so the
|
||
* caller owns the HTTP shape — the same limiter serves an H3 handler, a queue
|
||
* worker, or a gRPC interceptor.
|
||
*
|
||
* Use this when the store is shared (Redis, Postgres) and therefore async.
|
||
*/
|
||
export function createLimiter(options: LimiterOptions = {}) {
|
||
const store = options.store || new MemoryStore()
|
||
|
||
return async function check(input: LimitInput): Promise<LimitDecision> {
|
||
const max = planWriteQuota(input.plan, options)
|
||
if (isExempt(max, input.userId)) return UNLIMITED
|
||
|
||
const windowMs = (input.windowSec || 60) * 1000
|
||
const { count, resetAt } = await store.hit(bucketKey(input), windowMs)
|
||
return decide(max, count, resetAt)
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Build a SYNCHRONOUS limiter.
|
||
*
|
||
* Exists because callers that gate a request on the result must not be forced
|
||
* to become async. An existing sync call site that silently keeps calling an
|
||
* async limiter without awaiting turns every rejection into an unhandled
|
||
* promise: the limit appears configured, reports correctly, and enforces
|
||
* nothing. Offering a sync path removes that trap instead of relying on every
|
||
* call site remembering to await.
|
||
*
|
||
* Only usable with a synchronous store, which in practice means in-process —
|
||
* so this is single-instance by construction. For a fleet, use `createLimiter`
|
||
* with a shared store and await it.
|
||
*/
|
||
export function createSyncLimiter(options: SyncLimiterOptions = {}) {
|
||
const store = options.store || new MemoryStore()
|
||
|
||
return function check(input: LimitInput): LimitDecision {
|
||
const max = planWriteQuota(input.plan, options)
|
||
if (isExempt(max, input.userId)) return UNLIMITED
|
||
|
||
const windowMs = (input.windowSec || 60) * 1000
|
||
const { count, resetAt } = store.hit(bucketKey(input), windowMs)
|
||
return decide(max, count, resetAt)
|
||
}
|
||
}
|