First public release
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+177
@@ -0,0 +1,177 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user