import type { PacingRule } from './rate-state-backend'; export type ResolvedPacing = { provider: string; rules: PacingRule[]; }; export type AdaptiveAdmissionSnapshot = { provider: string; currentRequestsPerSecond: number; latencyMsEwma: number | null; successCount: number; rateLimitCount: number; retryAfterUntilMs: number | null; }; export interface RuntimeAdaptiveAdmission { resolvePacing(input: { orgId: string; toolId: string; declaredPacing: ResolvedPacing | null; fallbackPacing: ResolvedPacing; }): ResolvedPacing; observeProviderSuccess(input: { orgId: string; provider: string; latencyMs?: number | null; }): void; observeProviderBackpressure(input: { orgId: string; provider: string; retryAfterMs: number; }): void; snapshot(): Record; } export type InMemoryAdaptiveAdmissionOptions = { initialRequestsPerSecond: number; minRequestsPerSecond?: number; maxRequestsPerSecond: number; additiveIncreasePerWindow?: number; cleanWindowMs?: number; now?: () => number; }; type ProviderState = { provider: string; currentRequestsPerSecond: number; successCount: number; successCountAtLastIncrease: number; rateLimitCount: number; lastIncreaseAtMs: number; retryAfterUntilMs: number | null; latencyMsEwma: number | null; fallbackRules: PacingRule[]; hasSeenBackpressure: boolean; }; function finitePositive(value: number, fallback: number): number { return Number.isFinite(value) && value > 0 ? value : fallback; } function stateKey(orgId: string, provider: string): string { return `${orgId}:${provider}`; } function cloneRulesWithRequestsPerSecond( rules: readonly PacingRule[], requestsPerSecond: number, ): PacingRule[] { const rps = Math.max(1, Math.floor(requestsPerSecond)); return rules.map((rule) => ({ ...rule, requestsPerWindow: rule.windowMs === 1_000 ? rps : Math.max(1, Math.floor((rps * rule.windowMs) / 1_000)), })); } export function createInMemoryAdaptiveAdmission( options: InMemoryAdaptiveAdmissionOptions, ): RuntimeAdaptiveAdmission { const now = options.now ?? Date.now; const minRequestsPerSecond = finitePositive( options.minRequestsPerSecond ?? 1, 1, ); const initialRequestsPerSecond = Math.max( minRequestsPerSecond, finitePositive(options.initialRequestsPerSecond, minRequestsPerSecond), ); const maxRequestsPerSecond = Math.max( initialRequestsPerSecond, finitePositive(options.maxRequestsPerSecond, initialRequestsPerSecond), ); const additiveIncreasePerWindow = finitePositive( options.additiveIncreasePerWindow ?? 1, 1, ); const cleanWindowMs = finitePositive(options.cleanWindowMs ?? 1_000, 1_000); const states = new Map(); function getOrCreateState(input: { orgId: string; provider: string; fallbackRules: readonly PacingRule[]; }): ProviderState { const key = stateKey(input.orgId, input.provider); let state = states.get(key); if (!state) { state = { provider: input.provider, currentRequestsPerSecond: initialRequestsPerSecond, successCount: 0, successCountAtLastIncrease: 0, rateLimitCount: 0, lastIncreaseAtMs: now(), retryAfterUntilMs: null, latencyMsEwma: null, fallbackRules: [...input.fallbackRules], hasSeenBackpressure: false, }; states.set(key, state); } return state; } function getExistingState( orgId: string, provider: string, ): ProviderState | null { return states.get(stateKey(orgId, provider)) ?? null; } return { resolvePacing(input) { if (input.declaredPacing && input.declaredPacing.rules.length > 0) { return input.declaredPacing; } const state = getOrCreateState({ orgId: input.orgId, provider: input.fallbackPacing.provider, fallbackRules: input.fallbackPacing.rules, }); return { provider: input.fallbackPacing.provider, rules: cloneRulesWithRequestsPerSecond( state.fallbackRules, state.currentRequestsPerSecond, ), }; }, observeProviderSuccess(input) { const state = getExistingState(input.orgId, input.provider); if (!state) return; const observedLatency = input.latencyMs; if (observedLatency != null && Number.isFinite(observedLatency)) { const latency = Math.max(0, observedLatency); state.latencyMsEwma = state.latencyMsEwma == null ? latency : state.latencyMsEwma * 0.8 + latency * 0.2; } state.successCount += 1; const currentTime = now(); if ( state.retryAfterUntilMs != null && state.retryAfterUntilMs > currentTime ) { return; } state.retryAfterUntilMs = null; const successesSinceIncrease = state.successCount - state.successCountAtLastIncrease; if ( currentTime - state.lastIncreaseAtMs >= cleanWindowMs && successesSinceIncrease >= Math.max(1, state.currentRequestsPerSecond) ) { state.currentRequestsPerSecond = Math.min( maxRequestsPerSecond, state.hasSeenBackpressure ? state.currentRequestsPerSecond + additiveIncreasePerWindow : state.currentRequestsPerSecond * 2, ); state.lastIncreaseAtMs = currentTime; state.successCountAtLastIncrease = state.successCount; } }, observeProviderBackpressure(input) { const state = getExistingState(input.orgId, input.provider); if (!state) return; const currentTime = now(); state.rateLimitCount += 1; const nextRetryAfterUntilMs = currentTime + Math.max(0, input.retryAfterMs); if ( state.retryAfterUntilMs != null && state.retryAfterUntilMs > currentTime ) { state.retryAfterUntilMs = Math.max( state.retryAfterUntilMs, nextRetryAfterUntilMs, ); return; } state.hasSeenBackpressure = true; state.currentRequestsPerSecond = Math.max( minRequestsPerSecond, Math.floor(state.currentRequestsPerSecond / 2), ); state.retryAfterUntilMs = nextRetryAfterUntilMs; state.lastIncreaseAtMs = currentTime; state.successCountAtLastIncrease = state.successCount; }, snapshot() { return Object.fromEntries( [...states.entries()].map(([key, state]) => [ key, { provider: state.provider, currentRequestsPerSecond: state.currentRequestsPerSecond, latencyMsEwma: state.latencyMsEwma, successCount: state.successCount, rateLimitCount: state.rateLimitCount, retryAfterUntilMs: state.retryAfterUntilMs, }, ]), ); }, }; }