mirror of
https://github.com/HeyPuter/puter.git
synced 2026-10-11 14:21:51 +00:00
refactor(util): add BoundedTtlMap for core caches and log throttles
One size-capped, LRU-evicting map with an optional TTL and a once-per-window shouldEmit helper. Replaces the hand-rolled maps in the DynamoDB throttle and repair loggers (the repair one was unbounded), the Slack repeat throttle, the auth probe's reauth log throttle, and the card-verification status memo.
This commit is contained in:
1 parent
491729e060
commit
a014bec72a
7 files changed
+268
-62
No files matched your search
@@ -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<string, number>();
|
||||
|
||||
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<string, true>({
|
||||
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);
|
||||
|
||||
@@ -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<string, number>();
|
||||
const repairWarnings = new BoundedTtlMap<string, true>({
|
||||
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}`,
|
||||
);
|
||||
|
||||
@@ -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<string, { at: number; folded: number }>();
|
||||
// No TTL: an entry past its interval still carries the count to report.
|
||||
const lastLogged = new BoundedTtlMap<string, { at: number; folded: number }>({
|
||||
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
|
||||
|
||||
@@ -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<string, number>();
|
||||
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<string, true>({
|
||||
maxEntries: REAUTH_LOG_MAX_KEYS,
|
||||
ttlMs: REAUTH_LOG_WINDOW_MS,
|
||||
});
|
||||
|
||||
return async (req, _res, next): Promise<void> => {
|
||||
// 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 ?? '-'}`,
|
||||
);
|
||||
|
||||
@@ -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 <https://www.gnu.org/licenses/>.
|
||||
*/
|
||||
|
||||
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<string, boolean | null>({
|
||||
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<string, number>({
|
||||
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<string, number>({ 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<string, number>({ 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<string, number>({ 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<string, number>({ 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<string, true>({
|
||||
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<string, true>({
|
||||
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<string, true>({
|
||||
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<string, true>({
|
||||
maxEntries: 10,
|
||||
ttlMs: 60_000,
|
||||
});
|
||||
expect(map.shouldEmit('k')).toBe(true);
|
||||
map.delete('k');
|
||||
expect(map.shouldEmit('k')).toBe(true);
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -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 <https://www.gnu.org/licenses/>.
|
||||
*/
|
||||
|
||||
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<V> {
|
||||
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<K, V> {
|
||||
readonly #maxEntries: number;
|
||||
readonly #ttlMs: number | undefined;
|
||||
readonly #entries = new Map<K, Entry<V>>();
|
||||
|
||||
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<K, true>, key: K): boolean {
|
||||
if (this.has(key)) return false;
|
||||
this.set(key, true);
|
||||
return true;
|
||||
}
|
||||
|
||||
#live(key: K): Entry<V> | 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;
|
||||
}
|
||||
}
|
||||
@@ -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<boolean | null> {
|
||||
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;
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user