diff --git a/src/backend/clients/alarm/slack.ts b/src/backend/clients/alarm/slack.ts index 3a371f1bb..487048154 100644 --- a/src/backend/clients/alarm/slack.ts +++ b/src/backend/clients/alarm/slack.ts @@ -18,6 +18,7 @@ */ import type { ISlackAlertConfig, PagerSeverity } from '../../types'; +import { BoundedTtlMap } from '../../util/boundedTtlMap'; import type { AlertHandler, AlertPayload } from './types'; const REQUEST_TIMEOUT_MS = 5000; @@ -123,30 +124,14 @@ export function createSlackAlertHandler( ): AlertHandler { const webhookUrl = conf.webhookUrl as string; const throttleMs = conf.repeatThrottleMs ?? DEFAULT_REPEAT_THROTTLE_MS; - const lastPosted = new Map(); - - const shouldPost = (alert: AlertPayload): boolean => { - if (throttleMs <= 0) return true; - const now = Date.now(); - const previous = lastPosted.get(alert.id); - if (previous !== undefined && now - previous < throttleMs) return false; - - if (lastPosted.size >= MAX_THROTTLE_ENTRIES) { - for (const [id, at] of lastPosted) { - if (now - at >= throttleMs) lastPosted.delete(id); - } - // Still full of live entries — drop the oldest to stay bounded. - if (lastPosted.size >= MAX_THROTTLE_ENTRIES) { - const oldest = lastPosted.keys().next().value; - if (oldest !== undefined) lastPosted.delete(oldest); - } - } - lastPosted.set(alert.id, now); - return true; - }; + // A zero or negative throttle expires every entry on arrival. + const lastPosted = new BoundedTtlMap({ + maxEntries: MAX_THROTTLE_ENTRIES, + ttlMs: throttleMs, + }); return async (alert) => { - if (!shouldPost(alert)) return; + if (!lastPosted.shouldEmit(alert.id)) return; const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), REQUEST_TIMEOUT_MS); diff --git a/src/backend/clients/dynamodb/DDBClient.ts b/src/backend/clients/dynamodb/DDBClient.ts index feabff3ad..d9463fe13 100644 --- a/src/backend/clients/dynamodb/DDBClient.ts +++ b/src/backend/clients/dynamodb/DDBClient.ts @@ -43,6 +43,7 @@ import { once } from 'node:events'; import { Agent as httpsAgent } from 'node:https'; import { HttpError } from '../../core/http'; import type { IConfig, IDynamoConfig } from '../../types'; +import { BoundedTtlMap } from '../../util/boundedTtlMap.js'; import { Span } from '../../util/span.js'; import { PuterClient } from '../types'; import { @@ -124,15 +125,13 @@ const sleep = async (ms: number) => { // One caller writing an out-of-range number usually keeps doing it, so the // warning is per-target and throttled rather than one line per write. -const REPAIR_WARNING_INTERVAL_MS = 60_000; -const lastRepairWarningAt = new Map(); +const repairWarnings = new BoundedTtlMap({ + maxEntries: 1_000, + ttlMs: 60_000, +}); const warnRepaired = (target: string, message: string): void => { - const now = Date.now(); - const lastWarnedAt = lastRepairWarningAt.get(target) ?? 0; - if (now - lastWarnedAt < REPAIR_WARNING_INTERVAL_MS) return; - - lastRepairWarningAt.set(target, now); + if (!repairWarnings.shouldEmit(target)) return; console.warn( `[ddb] clamped an out-of-range value in ${target}: ${message}`, ); diff --git a/src/backend/clients/dynamodb/throttleLog.ts b/src/backend/clients/dynamodb/throttleLog.ts index eee6bde07..d93b04b56 100644 --- a/src/backend/clients/dynamodb/throttleLog.ts +++ b/src/backend/clients/dynamodb/throttleLog.ts @@ -18,6 +18,7 @@ */ import type { DynamoDBClient } from '@aws-sdk/client-dynamodb'; +import { BoundedTtlMap } from '../../util/boundedTtlMap.js'; const THROTTLE_ERROR_NAMES = new Set([ 'ProvisionedThroughputExceededException', @@ -39,7 +40,10 @@ interface Target { keys: string; } -const lastLogged = new Map(); +// No TTL: an entry past its interval still carries the count to report. +const lastLogged = new BoundedTtlMap({ + maxEntries: MAX_TRACKED_TARGETS, +}); const logOnce = (target: string, detail?: string): void => { const now = Date.now(); @@ -48,13 +52,6 @@ const logOnce = (target: string, detail?: string): void => { previous.folded += 1; return; } - - if (!previous && lastLogged.size >= MAX_TRACKED_TARGETS) { - for (const [key, entry] of lastLogged) { - if (now - entry.at >= LOG_INTERVAL_MS) lastLogged.delete(key); - } - if (lastLogged.size >= MAX_TRACKED_TARGETS) lastLogged.clear(); - } lastLogged.set(target, { at: now, folded: 0 }); const folded = previous?.folded diff --git a/src/backend/core/http/middleware/authProbe.ts b/src/backend/core/http/middleware/authProbe.ts index b58d9df5f..040ab6310 100644 --- a/src/backend/core/http/middleware/authProbe.ts +++ b/src/backend/core/http/middleware/authProbe.ts @@ -22,6 +22,7 @@ import type { AuthService, ReauthReason, } from '../../../services/auth/AuthService'; +import { BoundedTtlMap } from '../../../util/boundedTtlMap'; import { assertResolvedActor } from '../../actor'; import type { TokenSource } from '../types'; @@ -68,23 +69,10 @@ const REAUTH_LOG_MAX_KEYS = 1024; export const createAuthProbe = (opts: AuthProbeOptions): RequestHandler => { const { authService, cookieName } = opts; - const reauthLoggedAt = new Map(); - const shouldLogReauth = (key: string): boolean => { - const now = Date.now(); - const last = reauthLoggedAt.get(key); - if (last !== undefined && now - last < REAUTH_LOG_WINDOW_MS) { - return false; - } - if (reauthLoggedAt.size >= REAUTH_LOG_MAX_KEYS) { - // Map iterates in insertion order and every log re-inserts, so - // the first key is the least recently logged. - const oldest = reauthLoggedAt.keys().next().value; - if (oldest !== undefined) reauthLoggedAt.delete(oldest); - } - reauthLoggedAt.delete(key); - reauthLoggedAt.set(key, now); - return true; - }; + const reauthLogged = new BoundedTtlMap({ + maxEntries: REAUTH_LOG_MAX_KEYS, + ttlMs: REAUTH_LOG_WINDOW_MS, + }); return async (req, _res, next): Promise => { // If something upstream already attached an actor, respect it. @@ -144,7 +132,7 @@ export const createAuthProbe = (opts: AuthProbeOptions): RequestHandler => { return signedToken; }, }; - if (shouldLogReauth(`${reason}:${auth_id ?? '-'}`)) { + if (reauthLogged.shouldEmit(`${reason}:${auth_id ?? '-'}`)) { console.info( `[auth-v2] reauth reason=${reason} auth_id=${auth_id ?? '-'}`, ); diff --git a/src/backend/util/boundedTtlMap.test.ts b/src/backend/util/boundedTtlMap.test.ts new file mode 100644 index 000000000..722358282 --- /dev/null +++ b/src/backend/util/boundedTtlMap.test.ts @@ -0,0 +1,133 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { BoundedTtlMap } from './boundedTtlMap.ts'; + +describe('BoundedTtlMap', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + afterEach(() => { + vi.useRealTimers(); + }); + + it('stores and returns values, including falsy ones', () => { + const map = new BoundedTtlMap({ + maxEntries: 10, + }); + map.set('a', null).set('b', false); + expect(map.get('a')).toBeNull(); + expect(map.has('a')).toBe(true); + expect(map.get('b')).toBe(false); + expect(map.get('missing')).toBeUndefined(); + expect(map.has('missing')).toBe(false); + }); + + it('expires entries after the TTL', () => { + const map = new BoundedTtlMap({ + maxEntries: 10, + ttlMs: 1000, + }); + map.set('a', 1); + vi.advanceTimersByTime(999); + expect(map.get('a')).toBe(1); + vi.advanceTimersByTime(1); + expect(map.get('a')).toBeUndefined(); + expect(map.size).toBe(0); + }); + + it('keeps entries without a TTL until evicted', () => { + const map = new BoundedTtlMap({ maxEntries: 10 }); + map.set('a', 1); + vi.advanceTimersByTime(365 * 24 * 60 * 60 * 1000); + expect(map.get('a')).toBe(1); + }); + + it('evicts the least recently used entry once full', () => { + const map = new BoundedTtlMap({ maxEntries: 3 }); + map.set('a', 1).set('b', 2).set('c', 3); + // Reading `a` makes `b` the least recently used. + expect(map.get('a')).toBe(1); + map.set('d', 4); + expect(map.size).toBe(3); + expect(map.has('b')).toBe(false); + expect([map.get('a'), map.get('c'), map.get('d')]).toEqual([1, 3, 4]); + }); + + it('overwriting a key does not evict another', () => { + const map = new BoundedTtlMap({ maxEntries: 2 }); + map.set('a', 1).set('b', 2).set('a', 3); + expect(map.size).toBe(2); + expect(map.get('a')).toBe(3); + expect(map.get('b')).toBe(2); + }); + + it('supports delete and clear', () => { + const map = new BoundedTtlMap({ maxEntries: 5 }); + map.set('a', 1).set('b', 2); + expect(map.delete('a')).toBe(true); + expect(map.has('a')).toBe(false); + map.clear(); + expect(map.size).toBe(0); + }); + + describe('shouldEmit', () => { + it('is true once per window per key', () => { + const map = new BoundedTtlMap({ + maxEntries: 10, + ttlMs: 60_000, + }); + expect(map.shouldEmit('k')).toBe(true); + expect(map.shouldEmit('k')).toBe(false); + expect(map.shouldEmit('other')).toBe(true); + vi.advanceTimersByTime(60_000); + expect(map.shouldEmit('k')).toBe(true); + expect(map.shouldEmit('k')).toBe(false); + }); + + it('is always true with a zero TTL', () => { + const map = new BoundedTtlMap({ + maxEntries: 10, + ttlMs: 0, + }); + expect(map.shouldEmit('k')).toBe(true); + expect(map.shouldEmit('k')).toBe(true); + }); + + it('stays bounded under a flood of distinct keys', () => { + const map = new BoundedTtlMap({ + maxEntries: 100, + ttlMs: 60_000, + }); + for (let i = 0; i < 10_000; i++) map.shouldEmit(`k${i}`); + expect(map.size).toBe(100); + }); + + it('emits again after the key is deleted', () => { + const map = new BoundedTtlMap({ + maxEntries: 10, + ttlMs: 60_000, + }); + expect(map.shouldEmit('k')).toBe(true); + map.delete('k'); + expect(map.shouldEmit('k')).toBe(true); + }); + }); +}); diff --git a/src/backend/util/boundedTtlMap.ts b/src/backend/util/boundedTtlMap.ts new file mode 100644 index 000000000..1029ec44c --- /dev/null +++ b/src/backend/util/boundedTtlMap.ts @@ -0,0 +1,102 @@ +/* + * Copyright (C) 2024-present Puter Technologies Inc. + * + * This file is part of Puter. + * + * Puter is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published + * by the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see . + */ + +export interface BoundedTtlMapOptions { + /** Past this many entries, the least recently used one is evicted. */ + maxEntries: number; + /** Entry lifetime. Omit for entries that only leave by eviction. */ + ttlMs?: number; +} + +interface Entry { + value: V; + expiresAt: number; +} + +/** + * In-process map with a size cap and an optional per-entry TTL, for caches, + * memos and log throttles. Expired entries read as missing; reads refresh an + * entry's recency, so eviction drops the least recently used one. + */ +export class BoundedTtlMap { + readonly #maxEntries: number; + readonly #ttlMs: number | undefined; + readonly #entries = new Map>(); + + constructor({ maxEntries, ttlMs }: BoundedTtlMapOptions) { + this.#maxEntries = Math.max(1, maxEntries); + this.#ttlMs = ttlMs; + } + + get size(): number { + return this.#entries.size; + } + + get(key: K): V | undefined { + return this.#live(key)?.value; + } + + has(key: K): boolean { + return this.#live(key) !== undefined; + } + + set(key: K, value: V): this { + this.#entries.delete(key); + if (this.#entries.size >= this.#maxEntries) { + const oldest = this.#entries.keys().next(); + if (!oldest.done) this.#entries.delete(oldest.value); + } + this.#entries.set(key, { + value, + expiresAt: + this.#ttlMs === undefined ? Infinity : Date.now() + this.#ttlMs, + }); + return this; + } + + delete(key: K): boolean { + return this.#entries.delete(key); + } + + clear(): void { + this.#entries.clear(); + } + + /** + * True at most once per TTL window for `key`, for "log/post this at most + * once a window" throttles. Calls inside the window return false. + */ + shouldEmit(this: BoundedTtlMap, key: K): boolean { + if (this.has(key)) return false; + this.set(key, true); + return true; + } + + #live(key: K): Entry | undefined { + const entry = this.#entries.get(key); + if (!entry) return undefined; + if (entry.expiresAt <= Date.now()) { + this.#entries.delete(key); + return undefined; + } + this.#entries.delete(key); + this.#entries.set(key, entry); + return entry; + } +} diff --git a/src/backend/util/cardFallback.ts b/src/backend/util/cardFallback.ts index 3626f6771..3da02119c 100644 --- a/src/backend/util/cardFallback.ts +++ b/src/backend/util/cardFallback.ts @@ -30,6 +30,7 @@ import type { EventClient } from '../clients/event/EventClient'; import type { PreludeClient } from '../clients/prelude/PreludeClient'; import type { IConfig } from '../types'; +import { BoundedTtlMap } from './boundedTtlMap'; /** * `/send-confirm-phone` route rate limit: how many verification texts one @@ -133,11 +134,14 @@ export function cardFallbackDepsFrom( * switch still takes effect immediately. */ const CARD_STATUS_TTL_MS = 60_000; -let cardStatusCache: { at: number; enabled: boolean | null } | null = null; +const cardStatusCache = new BoundedTtlMap<'status', boolean | null>({ + maxEntries: 1, + ttlMs: CARD_STATUS_TTL_MS, +}); /** Drops the memoized probe answer. For tests, and for a config reload. */ export function resetCardVerificationStatusCache(): void { - cardStatusCache = null; + cardStatusCache.clear(); } /** @@ -147,10 +151,8 @@ export function resetCardVerificationStatusCache(): void { export async function isCardVerificationEnabled( deps: CardFallbackDeps, ): Promise { - const now = Date.now(); - if (cardStatusCache && now - cardStatusCache.at < CARD_STATUS_TTL_MS) { - return cardStatusCache.enabled; - } + const cached = cardStatusCache.get('status'); + if (cached !== undefined) return cached; let enabled: boolean | null = null; try { enabled = await deps.probeCardVerification(); @@ -159,7 +161,7 @@ export async function isCardVerificationEnabled( // not work is worse than not offering one. console.warn('[card-verification] status probe failed:', e); } - cardStatusCache = { at: now, enabled }; + cardStatusCache.set('status', enabled); return enabled; }