refactor(api): purge cache by canonical media prefix (#2802)

This commit is contained in:
Hampus
2026-09-16 02:32:36 +02:00
committed by GitHub
parent 570c8776c4
commit 7412ec3395
13 changed files with 293 additions and 329 deletions
+1 -1
View File
@@ -29,7 +29,7 @@ export interface IKVSubscription {
}
export interface KVPurgeBatchResult {
urls: Array<string>;
entries: Array<string>;
tokensConsumed: number;
}
+11 -11
View File
@@ -217,7 +217,7 @@ local refillIntervalMs = tonumber(ARGV[5])
local queueSize = redis.call('SCARD', queueKey)
if queueSize == 0 then
return '{"urls":[],"tokens":0}'
return '{"entries":[],"tokens":0}'
end
local tokens = maxTokens
@@ -237,14 +237,14 @@ end
local toPop = math.min(maxItems, math.floor(tokens), queueSize)
if toPop <= 0 then
redis.call('SET', bucketKey, cjson.encode({tokens = tokens, lastRefill = lastRefill}), 'EX', 3600)
return '{"urls":[],"tokens":0}'
return '{"entries":[],"tokens":0}'
end
local urls = redis.call('SPOP', queueKey, toPop)
tokens = tokens - #urls
local entries = redis.call('SPOP', queueKey, toPop)
tokens = tokens - #entries
redis.call('SET', bucketKey, cjson.encode({tokens = tokens, lastRefill = lastRefill}), 'EX', 3600)
return cjson.encode({urls = urls, tokens = #urls})
return cjson.encode({entries = entries, tokens = #entries})
`;
const CLAIM_BULK_DELETION_SCRIPT = `
local score = redis.call('ZSCORE', KEYS[1], ARGV[1])
@@ -902,17 +902,17 @@ function parseRateLimitResult(value: unknown): KVRateLimitResult {
function parsePurgeBatchResult(value: unknown, maxItems: number): KVPurgeBatchResult {
const command = 'dequeuePurgeBatch';
if (!isJsonObject(value)) throw createInvalidResponseError(command, 'a purge batch object');
const {urls, tokens} = value;
const {entries, tokens} = value;
if (
!Array.isArray(urls) ||
!urls.every((url): url is string => typeof url === 'string') ||
!Array.isArray(entries) ||
!entries.every((entry): entry is string => typeof entry === 'string') ||
!isNonNegativeSafeInteger(tokens) ||
tokens !== urls.length ||
urls.length > maxItems
tokens !== entries.length ||
entries.length > maxItems
) {
throw createInvalidResponseError(command, 'a bounded string array and matching token count');
}
return {urls, tokensConsumed: tokens};
return {entries, tokensConsumed: tokens};
}
function isJsonObject(value: unknown): value is Record<string, unknown> {
@@ -153,7 +153,7 @@ describe('KVClient script execution', () => {
},
{
name: 'dequeuePurgeBatch',
reply: JSON.stringify({urls: ['https://fluxer.test/a.png'], tokens: 1}),
reply: JSON.stringify({entries: ['/attachments/1/2/a'], tokens: 1}),
keyCount: 2,
run: async (client) => client.dequeuePurgeBatch('queue:key', 'bucket:key', 10, 10, 1, 1000),
},
@@ -5,18 +5,13 @@ import {createHttpCachePurgeAdapter} from '@app/api/infrastructure/HttpCachePurg
import {createNoneCachePurgeAdapter} from '@app/api/infrastructure/NoneCachePurgeAdapter';
import type {CachePurgeAdapterName} from '@fluxer/config/src/MasterConfig';
export interface CachePurgeBatch {
readonly exact: ReadonlyArray<string>;
readonly prefix: ReadonlyArray<string>;
}
export type CachePurgeOutcome =
| {readonly kind: 'purged'}
| {readonly kind: 'invalid_entries'; readonly status: number}
| {readonly kind: 'rejected'; readonly status: number}
| {readonly kind: 'failed'; readonly status: number | null; readonly error: unknown};
export interface CachePurgeAdapter {
purge(batch: CachePurgeBatch): Promise<CachePurgeOutcome>;
purge(prefixes: ReadonlyArray<string>): Promise<CachePurgeOutcome>;
}
const CACHE_PURGE_ADAPTERS = {
@@ -0,0 +1,70 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Config} from '@app/api/Config';
import {canonicalizePurgeUrl} from '@app/api/infrastructure/CachePurgePaths';
import {afterEach, beforeEach, describe, expect, it} from 'vitest';
const MEDIA = 'https://media.test';
describe('canonicalizePurgeUrl', () => {
let previousMedia: string;
beforeEach(() => {
previousMedia = Config.endpoints.media;
Config.endpoints.media = MEDIA;
});
afterEach(() => {
Config.endpoints.media = previousMedia;
});
it('qualifies every prefix with the media host and no scheme', () => {
expect(canonicalizePurgeUrl(`${MEDIA}/attachments/1/2/photo.png`)).toEqual(['media.test/attachments/1/2/photo']);
});
it('drops the file extension so every served format is covered', () => {
expect(canonicalizePurgeUrl(`${MEDIA}/emojis/1531309058777157632.webp`)).toEqual([
'media.test/emojis/1531309058777157632',
]);
});
it('pairs an asset hash with its animated spelling in both directions', () => {
expect(canonicalizePurgeUrl(`${MEDIA}/avatars/1/a_b35cc3d3`)).toEqual([
'media.test/avatars/1/a_b35cc3d3',
'media.test/avatars/1/b35cc3d3',
]);
expect(canonicalizePurgeUrl(`${MEDIA}/icons/1/808b11fa.webp`)).toEqual([
'media.test/icons/1/808b11fa',
'media.test/icons/1/a_808b11fa',
]);
});
it('leaves a non-hash stem alone', () => {
expect(canonicalizePurgeUrl(`${MEDIA}/attachments/1/2/holiday.jpg`)).toEqual([
'media.test/attachments/1/2/holiday',
]);
});
it('drops the query and fragment and decodes the path', () => {
expect(canonicalizePurgeUrl(`${MEDIA}/attachments/1/2/ação.png?size=128#x`)).toEqual([
'media.test/attachments/1/2/ação',
]);
});
it('keeps a base path when the media endpoint carries one', () => {
Config.endpoints.media = `${MEDIA}/media`;
expect(canonicalizePurgeUrl(`${MEDIA}/media/avatars/1/b35cc3d3`)).toEqual([
'media.test/media/avatars/1/b35cc3d3',
'media.test/media/avatars/1/a_b35cc3d3',
]);
expect(canonicalizePurgeUrl(`${MEDIA}/avatars/1/b35cc3d3`)).toEqual([]);
});
it('refuses anything that is not a media CDN object', () => {
expect(canonicalizePurgeUrl('https://elsewhere.test/avatars/1/b35cc3d3')).toEqual([]);
expect(canonicalizePurgeUrl(`${MEDIA}/avatars`)).toEqual([]);
expect(canonicalizePurgeUrl(`${MEDIA}/`)).toEqual([]);
expect(canonicalizePurgeUrl(`${MEDIA}/emojis/.webp`)).toEqual([]);
expect(canonicalizePurgeUrl('not-a-url')).toEqual([]);
});
});
@@ -0,0 +1,53 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Config} from '@app/api/Config';
const EXTENSION_PATTERN = /\.[A-Za-z0-9]{1,8}$/u;
const ASSET_HASH_PATTERN = /^[0-9a-f]{8}$/u;
const TRAILING_SLASHES_PATTERN = /\/+$/u;
const ANIMATED_PREFIX = 'a_';
const MIN_PATH_SEGMENTS = 2;
function mediaBase(): URL | null {
return URL.parse(Config.endpoints.media);
}
function decodePathname(pathname: string): string {
try {
return decodeURIComponent(pathname);
} catch {
return pathname;
}
}
function animationVariants(stem: string): Array<string> {
if (stem.startsWith(ANIMATED_PREFIX)) {
const bare = stem.slice(ANIMATED_PREFIX.length);
return ASSET_HASH_PATTERN.test(bare) ? [stem, bare] : [stem];
}
return ASSET_HASH_PATTERN.test(stem) ? [stem, `${ANIMATED_PREFIX}${stem}`] : [stem];
}
export function canonicalizePurgeUrl(url: string): Array<string> {
const base = mediaBase();
const parsed = URL.parse(url);
if (base === null || parsed === null || parsed.origin !== base.origin) {
return [];
}
const basePath = base.pathname.replace(TRAILING_SLASHES_PATTERN, '');
if (!parsed.pathname.startsWith(`${basePath}/`)) {
return [];
}
const segments = decodePathname(parsed.pathname.slice(basePath.length))
.split('/')
.filter((segment) => segment !== '');
if (segments.length < MIN_PATH_SEGMENTS) {
return [];
}
const stem = segments[segments.length - 1]!.replace(EXTENSION_PATTERN, '');
if (stem === '') {
return [];
}
const directory = segments.slice(0, -1).join('/');
return animationVariants(stem).map((variant) => `${base.host}${basePath}/${directory}/${variant}`);
}
@@ -1,6 +1,6 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {CachePurgeBatch} from '@app/api/infrastructure/CachePurgeAdapter';
import {canonicalizePurgeUrl} from '@app/api/infrastructure/CachePurgePaths';
import {Logger} from '@app/api/Logger';
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
@@ -8,112 +8,69 @@ export interface IPurgeQueue {
addUrls(urls: Array<string>): Promise<void>;
}
const EXACT_QUEUE_KEY = 'cache_purge:exact';
const PREFIX_QUEUE_KEY = 'cache_purge:prefix';
const EXACT_BUCKET_KEY = 'cache_purge:budget:exact';
const PREFIX_BUCKET_KEY = 'cache_purge:budget:prefix';
const QUEUE_KEY = 'cache_purge:queue';
const BUDGET_KEY = 'cache_purge:budget';
const REJECTED_KEY = 'cache_purge:rejected';
const EXACT_MAX_TOKENS = 120;
const EXACT_REFILL_RATE = 5;
const EXACT_REFILL_INTERVAL_MS = 1000;
const PREFIX_MAX_TOKENS = 20;
const PREFIX_REFILL_RATE = 1;
const PREFIX_REFILL_INTERVAL_MS = 2000;
function isPrefix(url: string): boolean {
return url.endsWith('*') || url.endsWith('/');
}
const MAX_TOKENS = 120;
const REFILL_RATE = 5;
const REFILL_INTERVAL_MS = 1000;
export class CachePurgeQueue implements IPurgeQueue {
private readonly kvClient: IKVProvider;
constructor(kvClient: IKVProvider) {
this.kvClient = kvClient;
}
constructor(private readonly kvClient: IKVProvider) {}
async addUrls(urls: Array<string>): Promise<void> {
if (urls.length === 0) {
return;
}
const exactUrls: Array<string> = [];
const prefixUrls: Array<string> = [];
const prefixes = new Set<string>();
for (const url of urls) {
const trimmed = url.trim();
if (trimmed === '') {
const canonical = canonicalizePurgeUrl(url);
if (canonical.length === 0) {
Logger.warn({url}, 'Skipped a cache purge URL that is not a media CDN object');
continue;
}
if (isPrefix(trimmed)) {
prefixUrls.push(trimmed);
} else {
exactUrls.push(trimmed);
for (const prefix of canonical) {
prefixes.add(prefix);
}
}
if (prefixes.size === 0) {
return;
}
try {
await this.addToSets({exact: exactUrls, prefix: prefixUrls});
Logger.debug({exact: exactUrls.length, prefix: prefixUrls.length}, 'Added URLs to cache purge queue');
await this.kvClient.sadd(QUEUE_KEY, ...prefixes);
Logger.debug({prefixes: prefixes.size}, 'Added prefixes to the cache purge queue');
} catch (error) {
Logger.error(
{error, exact: exactUrls.length, prefix: prefixUrls.length},
'Failed to add URLs to cache purge queue',
);
Logger.error({error, prefixes: prefixes.size}, 'Failed to add prefixes to the cache purge queue');
throw error;
}
}
async dequeueBatch(): Promise<CachePurgeBatch> {
const exact = await this.kvClient.dequeuePurgeBatch(
EXACT_QUEUE_KEY,
EXACT_BUCKET_KEY,
EXACT_MAX_TOKENS,
EXACT_MAX_TOKENS,
EXACT_REFILL_RATE,
EXACT_REFILL_INTERVAL_MS,
async dequeueBatch(): Promise<Array<string>> {
const batch = await this.kvClient.dequeuePurgeBatch(
QUEUE_KEY,
BUDGET_KEY,
MAX_TOKENS,
MAX_TOKENS,
REFILL_RATE,
REFILL_INTERVAL_MS,
);
return batch.entries;
}
async requeue(prefixes: ReadonlyArray<string>): Promise<void> {
if (prefixes.length === 0) {
return;
}
try {
const prefix = await this.kvClient.dequeuePurgeBatch(
PREFIX_QUEUE_KEY,
PREFIX_BUCKET_KEY,
PREFIX_MAX_TOKENS,
PREFIX_MAX_TOKENS,
PREFIX_REFILL_RATE,
PREFIX_REFILL_INTERVAL_MS,
);
return {exact: exact.urls, prefix: prefix.urls};
await this.kvClient.sadd(QUEUE_KEY, ...prefixes);
} catch (error) {
await this.requeue({exact: exact.urls, prefix: []});
Logger.error({error, prefixes: prefixes.length}, 'Failed to requeue cache purge prefixes');
throw error;
}
}
async requeue(batch: CachePurgeBatch): Promise<void> {
try {
await this.addToSets(batch);
} catch (error) {
Logger.error(
{error, exact: batch.exact.length, prefix: batch.prefix.length},
'Failed to requeue cache purge entries',
);
throw error;
async reject(prefixes: ReadonlyArray<string>): Promise<void> {
if (prefixes.length === 0) {
return;
}
}
async reject(batch: CachePurgeBatch): Promise<void> {
await this.kvClient.sadd(REJECTED_KEY, ...batch.exact, ...batch.prefix);
Logger.error(
{exact: batch.exact.length, prefix: batch.prefix.length, key: REJECTED_KEY},
'Set aside cache purge entries the endpoint rejected',
);
}
private async addToSets(batch: CachePurgeBatch): Promise<void> {
const ops: Array<Promise<number>> = [];
if (batch.exact.length > 0) {
ops.push(this.kvClient.sadd(EXACT_QUEUE_KEY, ...batch.exact));
}
if (batch.prefix.length > 0) {
ops.push(this.kvClient.sadd(PREFIX_QUEUE_KEY, ...batch.prefix));
}
await Promise.all(ops);
await this.kvClient.sadd(REJECTED_KEY, ...prefixes);
}
}
@@ -1,26 +1,18 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {APICachePurgeConfig} from '@app/api/config/APIConfig';
import type {CachePurgeAdapter, CachePurgeBatch, CachePurgeOutcome} from '@app/api/infrastructure/CachePurgeAdapter';
import type {CachePurgeAdapter, CachePurgeOutcome} from '@app/api/infrastructure/CachePurgeAdapter';
import * as FetchUtils from '@app/api/utils/FetchUtils';
function stripOneTrailingAsterisk(entry: string): string {
return entry.endsWith('*') ? entry.slice(0, -1) : entry;
}
export function createHttpCachePurgeAdapter(config: APICachePurgeConfig): CachePurgeAdapter {
return {
async purge(batch: CachePurgeBatch): Promise<CachePurgeOutcome> {
const headers: Record<string, string> = {'Content-Type': 'application/json'};
if (config.http.token !== '') {
headers.Authorization = `Bearer ${config.http.token}`;
}
async purge(prefixes: ReadonlyArray<string>): Promise<CachePurgeOutcome> {
let response: Response;
try {
response = await fetch(config.http.endpoint, {
method: 'POST',
headers,
body: JSON.stringify({exact: batch.exact, prefix: batch.prefix.map(stripOneTrailingAsterisk)}),
headers: {'Content-Type': 'application/json', Authorization: `Bearer ${config.http.token}`},
body: JSON.stringify({prefixes}),
redirect: 'manual',
signal: AbortSignal.timeout(config.http.timeoutMs),
});
@@ -32,7 +24,7 @@ export function createHttpCachePurgeAdapter(config: APICachePurgeConfig): CacheP
return {kind: 'purged'};
}
if (response.status === 400 || response.status === 422) {
return {kind: 'invalid_entries', status: response.status};
return {kind: 'rejected', status: response.status};
}
return {kind: 'failed', status: response.status, error: null};
},
@@ -704,7 +704,7 @@ export class MockKVProvider implements IKVProvider {
refillRate: number,
refillIntervalMs: number,
): Promise<{
urls: Array<string>;
entries: Array<string>;
tokensConsumed: number;
}> {
this.dequeuePurgeBatchSpy(queueKey, bucketKey, maxItems, maxTokens, refillRate, refillIntervalMs);
@@ -734,13 +734,13 @@ export class MockKVProvider implements IKVProvider {
if (toPop <= 0) {
this.stringStore.set(bucketKey, JSON.stringify({tokens, lastRefill}));
this.expiries.set(bucketKey, now + 3600 * 1000);
return {urls: [], tokensConsumed: 0};
return {entries: [], tokensConsumed: 0};
}
const urls = await this.spop(queueKey, toPop);
tokens -= urls.length;
const entries = await this.spop(queueKey, toPop);
tokens -= entries.length;
this.stringStore.set(bucketKey, JSON.stringify({tokens, lastRefill}));
this.expiries.set(bucketKey, now + 3600 * 1000);
return {urls, tokensConsumed: urls.length};
return {entries, tokensConsumed: entries.length};
}
async evalScript(
@@ -1,104 +1,53 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {Config} from '@app/api/Config';
import {
type CachePurgeBatch,
type CachePurgeOutcome,
createCachePurgeAdapter,
} from '@app/api/infrastructure/CachePurgeAdapter';
import {type CachePurgeOutcome, createCachePurgeAdapter} from '@app/api/infrastructure/CachePurgeAdapter';
import {CachePurgeQueue} from '@app/api/infrastructure/CachePurgeQueue';
import {Logger} from '@app/api/Logger';
import {getWorkerDependencies} from '@app/api/worker/WorkerContext';
import type {CachePurgeAdapterName} from '@fluxer/config/src/MasterConfig';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
const RUN_DEADLINE_MS = 10_000;
type CachePurgeFailure = Extract<CachePurgeOutcome, {kind: 'failed'}>;
function describeBatch(adapter: CachePurgeAdapterName, batch: CachePurgeBatch) {
return {adapter, exact: batch.exact.length, prefix: batch.prefix.length};
}
function logFailure(adapter: CachePurgeAdapterName, requeued: CachePurgeBatch, failure: CachePurgeFailure): void {
function logFailure(context: Record<string, unknown>, failure: CachePurgeFailure): void {
const {status, error} = failure;
const context = {...describeBatch(adapter, requeued), status, error};
if (status === null || status === 408 || status === 429 || status >= 500) {
Logger.warn(context, 'Cache purge request failed, requeued its entries');
Logger.warn({...context, status, error}, 'Cache purge request failed, requeued its prefixes');
return;
}
Logger.error(context, 'Cache purge endpoint refused the request, requeued its entries');
}
function untriedFrom(singles: ReadonlyArray<CachePurgeBatch>, index: number): CachePurgeBatch {
const untried = singles.slice(index);
return {exact: untried.flatMap((single) => single.exact), prefix: untried.flatMap((single) => single.prefix)};
Logger.error({...context, status, error}, 'Cache purge endpoint refused the request, requeued its prefixes');
}
const processCachePurgeQueue: WorkerTaskHandler = async (_payload, _helpers) => {
const adapterName = Config.cachePurge.adapter;
if (adapterName === 'none') {
const adapter = Config.cachePurge.adapter;
if (adapter === 'none') {
Logger.warn('Skipped a cache purge run because the adapter is none');
return;
}
const adapter = createCachePurgeAdapter(Config.cachePurge);
const queue = new CachePurgeQueue(getWorkerDependencies().kvClient);
const startedAt = Date.now();
const batch = await queue.dequeueBatch();
if (batch.exact.length === 0 && batch.prefix.length === 0) {
const prefixes = await queue.dequeueBatch();
if (prefixes.length === 0) {
return;
}
const context = {adapter, prefixes: prefixes.length};
let outcome: CachePurgeOutcome;
try {
outcome = await adapter.purge(batch);
outcome = await createCachePurgeAdapter(Config.cachePurge).purge(prefixes);
} catch (error) {
await queue.requeue(batch);
await queue.requeue(prefixes);
throw error;
}
if (outcome.kind === 'purged') {
Logger.debug(describeBatch(adapterName, batch), 'Purged a batch of cache entries');
Logger.debug(context, 'Purged a batch of cache prefixes');
return;
}
if (outcome.kind === 'failed') {
await queue.requeue(batch);
logFailure(adapterName, batch, outcome);
if (outcome.kind === 'rejected') {
await queue.reject(prefixes);
Logger.error({...context, status: outcome.status}, 'Cache purge endpoint rejected the batch, set it aside');
return;
}
Logger.warn(
{...describeBatch(adapterName, batch), status: outcome.status},
'Cache purge endpoint rejected the batch, retrying each entry on its own',
);
const singles: Array<CachePurgeBatch> = [
...batch.exact.map((entry) => ({exact: [entry], prefix: []})),
...batch.prefix.map((entry) => ({exact: [], prefix: [entry]})),
];
for (const [index, single] of singles.entries()) {
if (Date.now() - startedAt >= RUN_DEADLINE_MS) {
const untried = untriedFrom(singles, index);
await queue.requeue(untried);
Logger.warn(
describeBatch(adapterName, untried),
'Cache purge run reached its deadline, requeued untried entries',
);
return;
}
let singleOutcome: CachePurgeOutcome;
try {
singleOutcome = await adapter.purge(single);
if (singleOutcome.kind === 'invalid_entries') {
await queue.reject(single);
}
} catch (error) {
await queue.requeue(untriedFrom(singles, index));
throw error;
}
if (singleOutcome.kind === 'failed') {
const untried = untriedFrom(singles, index);
await queue.requeue(untried);
logFailure(adapterName, untried, singleOutcome);
return;
}
}
await queue.requeue(prefixes);
logFailure(context, outcome);
};
export default processCachePurgeQueue;
@@ -13,15 +13,14 @@ import {delay, HttpResponse, http} from 'msw';
import {afterEach, beforeEach, describe, expect, it, vi} from 'vitest';
const HELPERS = {logger: new NoopLogger()} as unknown as WorkerTaskHelpers;
const ENDPOINT = 'https://cache-purge.test/purge';
const ENDPOINT = 'https://cache-purge.test/__cache/purge';
const MEDIA = 'https://media.test';
const EXACT_KEY = 'cache_purge:exact';
const PREFIX_KEY = 'cache_purge:prefix';
const QUEUE_KEY = 'cache_purge:queue';
const BUDGET_KEY = 'cache_purge:budget';
const REJECTED_KEY = 'cache_purge:rejected';
interface PurgeBody {
exact: Array<string>;
prefix: Array<string>;
prefixes: Array<string>;
}
interface RecordedRequest {
@@ -52,10 +51,6 @@ function recordPurges(respond: (body: PurgeBody) => Response | Promise<Response>
return requests;
}
function entryCount(body: PurgeBody): number {
return body.exact.length + body.prefix.length;
}
async function members(kvClient: MockKVProvider, key: string): Promise<Array<string>> {
return (await kvClient.smembers(key)).sort();
}
@@ -66,9 +61,12 @@ function sorted(values: Array<string>): Array<string> {
describe('processCachePurgeQueue', () => {
let previousCachePurge: APICachePurgeConfig;
let previousMedia: string;
beforeEach(() => {
previousCachePurge = {adapter: Config.cachePurge.adapter, http: Config.cachePurge.http};
previousMedia = Config.endpoints.media;
Config.endpoints.media = MEDIA;
Config.cachePurge.adapter = 'http';
Config.cachePurge.http = {endpoint: ENDPOINT, token: 'test-token', timeoutMs: 50};
});
@@ -76,15 +74,18 @@ describe('processCachePurgeQueue', () => {
afterEach(() => {
Config.cachePurge.adapter = previousCachePurge.adapter;
Config.cachePurge.http = previousCachePurge.http;
Config.endpoints.media = previousMedia;
clearWorkerDependencies();
vi.useRealTimers();
});
it('sends queued exact and prefix entries in one request and drains the queue', async () => {
it('sends host qualified prefixes in one request and drains the queue', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/a_abc`, `${MEDIA}/emojis/9.webp`, `${MEDIA}/attachments/1/2/ação.png`];
const prefix = [`${MEDIA}/some/dir/`];
await queue.addUrls([...exact, ...prefix]);
await queue.addUrls([
`${MEDIA}/avatars/1/a_b35cc3d3`,
`${MEDIA}/emojis/9.webp`,
`${MEDIA}/attachments/1/2/ação.png`,
]);
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
@@ -92,55 +93,69 @@ describe('processCachePurgeQueue', () => {
expect(requests).toHaveLength(1);
expect(requests[0]!.authorization).toBe('Bearer test-token');
expect(requests[0]!.contentType).toBe('application/json');
expect(sorted(requests[0]!.body.exact)).toEqual(sorted(exact));
expect(requests[0]!.body.prefix).toEqual(prefix);
expect(await members(kvClient, EXACT_KEY)).toEqual([]);
expect(await members(kvClient, PREFIX_KEY)).toEqual([]);
expect(sorted(requests[0]!.body.prefixes)).toEqual(
sorted([
'media.test/avatars/1/a_b35cc3d3',
'media.test/avatars/1/b35cc3d3',
'media.test/emojis/9',
'media.test/attachments/1/2/ação',
]),
);
expect(await members(kvClient, QUEUE_KEY)).toEqual([]);
});
it('strips a trailing asterisk from a prefix entry and keeps a trailing slash', async () => {
const {kvClient, queue} = createHarness();
await queue.addUrls([`${MEDIA}/stickers/1**`, `${MEDIA}/emojis/`]);
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(requests[0]!.body.exact).toEqual([]);
expect(sorted(requests[0]!.body.prefix)).toEqual([`${MEDIA}/emojis/`, `${MEDIA}/stickers/1*`]);
expect(await members(kvClient, PREFIX_KEY)).toEqual([]);
});
it('sends no authorization header when no token is configured', async () => {
Config.cachePurge.http.token = '';
it('collapses every extension of one asset into a single prefix', async () => {
const {queue} = createHarness();
await queue.addUrls([`${MEDIA}/avatars/1/abc`]);
await queue.addUrls([`${MEDIA}/stickers/7.webp`, `${MEDIA}/stickers/7.gif`, `${MEDIA}/stickers/7.png`]);
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(requests[0]!.authorization).toBeNull();
expect(requests[0]!.body.prefixes).toEqual(['media.test/stickers/7']);
});
it.each([200, 202, 204])('treats %i as accepted and drops the batch', async (status) => {
const {kvClient, queue} = createHarness();
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const requests = recordPurges(() => new HttpResponse(null, {status}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(await members(kvClient, QUEUE_KEY)).toEqual([]);
expect(await members(kvClient, REJECTED_KEY)).toEqual([]);
});
it('never queues a URL that is not a media CDN object', async () => {
const {kvClient, queue} = createHarness();
await queue.addUrls(['https://elsewhere.test/avatars/1/b35cc3d3', `${MEDIA}/avatars`, 'not-a-url']);
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toEqual([]);
expect(await members(kvClient, QUEUE_KEY)).toEqual([]);
});
it('requeues the whole batch after a server error', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`];
const prefix = [`${MEDIA}/some/dir/`];
await queue.addUrls([...exact, ...prefix]);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`, `${MEDIA}/attachments/1/3/b.png`]);
const requests = recordPurges(() => new HttpResponse(null, {status: 503}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact));
expect(await members(kvClient, PREFIX_KEY)).toEqual(prefix);
expect(await members(kvClient, QUEUE_KEY)).toEqual([
'media.test/attachments/1/2/a',
'media.test/attachments/1/3/b',
]);
expect(await members(kvClient, REJECTED_KEY)).toEqual([]);
});
it('requeues the batch when the endpoint does not answer in time', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`];
await queue.addUrls(exact);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const requests = recordPurges(async () => {
await delay('infinite');
return new HttpResponse(null, {status: 204});
@@ -149,25 +164,23 @@ describe('processCachePurgeQueue', () => {
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact));
expect(await members(kvClient, QUEUE_KEY)).toEqual(['media.test/attachments/1/2/a']);
});
it('requeues the batch when the endpoint is unreachable', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`];
await queue.addUrls(exact);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const requests = recordPurges(() => HttpResponse.error());
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact));
expect(await members(kvClient, QUEUE_KEY)).toEqual(['media.test/attachments/1/2/a']);
});
it('requeues the batch on a redirect without following it', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`];
await queue.addUrls(exact);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const redirectTarget = 'https://cache-purge.test/elsewhere';
const redirectedRequests: Array<string> = [];
server.use(
@@ -182,122 +195,51 @@ describe('processCachePurgeQueue', () => {
expect(requests).toHaveLength(1);
expect(redirectedRequests).toEqual([]);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact));
expect(await members(kvClient, QUEUE_KEY)).toEqual(['media.test/attachments/1/2/a']);
});
it('requeues the batch on an authorisation failure without splitting it', async () => {
it('requeues rather than discards when the bearer token is refused', async () => {
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`, `${MEDIA}/emojis/9.webp`];
await queue.addUrls(exact);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const requests = recordPurges(() => new HttpResponse(null, {status: 401}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toHaveLength(1);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact));
expect(await members(kvClient, QUEUE_KEY)).toEqual(['media.test/attachments/1/2/a']);
expect(await members(kvClient, REJECTED_KEY)).toEqual([]);
});
it('sets aside only the entry the endpoint rejects on its own', async () => {
it.each([400, 422])('sets the batch aside when the endpoint answers %i', async (status) => {
const {kvClient, queue} = createHarness();
const poison = `${MEDIA}/attachments/1/2/poison.png`;
await queue.addUrls([`${MEDIA}/avatars/1/abc`, poison, `${MEDIA}/emojis/9.webp`]);
const requests = recordPurges((body) =>
body.exact.includes(poison) ? new HttpResponse(null, {status: 422}) : new HttpResponse(null, {status: 204}),
);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`, `${MEDIA}/emojis/9.webp`]);
const requests = recordPurges(() => new HttpResponse(null, {status}));
await processCachePurgeQueue({}, HELPERS);
expect(requests.map((request) => entryCount(request.body))).toEqual([3, 1, 1, 1]);
expect(await members(kvClient, REJECTED_KEY)).toEqual([poison]);
expect(await members(kvClient, EXACT_KEY)).toEqual([]);
});
it('splits the batch when the endpoint answers 400', async () => {
const {kvClient, queue} = createHarness();
const poison = `${MEDIA}/attachments/1/2/poison.png`;
await queue.addUrls([`${MEDIA}/avatars/1/abc`, poison, `${MEDIA}/emojis/9.webp`]);
const requests = recordPurges((body) =>
body.exact.includes(poison) ? new HttpResponse(null, {status: 400}) : new HttpResponse(null, {status: 204}),
);
await processCachePurgeQueue({}, HELPERS);
expect(requests.map((request) => entryCount(request.body))).toEqual([3, 1, 1, 1]);
expect(await members(kvClient, REJECTED_KEY)).toEqual([poison]);
expect(await members(kvClient, EXACT_KEY)).toEqual([]);
});
it('keeps untried exact and prefix entries queued when a single-entry retry fails', async () => {
const {kvClient, queue} = createHarness();
const first = `${MEDIA}/avatars/1/abc`;
const second = `${MEDIA}/banners/1/def`;
const third = `${MEDIA}/emojis/9.webp`;
const directory = `${MEDIA}/some/dir/`;
await queue.addUrls([first, second, third, directory]);
let singleRequests = 0;
const requests = recordPurges((body) => {
if (entryCount(body) > 1) {
return new HttpResponse(null, {status: 422});
}
singleRequests++;
return new HttpResponse(null, {status: singleRequests === 1 ? 204 : 503});
});
await processCachePurgeQueue({}, HELPERS);
expect(requests.map((request) => request.body.exact)).toEqual([[first, second, third], [first], [second]]);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted([second, third]));
expect(await members(kvClient, PREFIX_KEY)).toEqual([directory]);
expect(await members(kvClient, REJECTED_KEY)).toEqual([]);
});
it('requeues untried entries when the fallback runs out of time', async () => {
vi.useFakeTimers({toFake: ['Date']});
vi.setSystemTime(new Date('2026-09-13T00:00:00.000Z'));
const {kvClient, queue} = createHarness();
const entries = [1, 2, 3, 4, 5].map((index) => `${MEDIA}/avatars/${index}/abc`);
await queue.addUrls(entries);
const requests = recordPurges((body) => {
if (entryCount(body) > 1) {
return new HttpResponse(null, {status: 422});
}
vi.setSystemTime(Date.now() + 4_000);
return new HttpResponse(null, {status: 204});
});
await processCachePurgeQueue({}, HELPERS);
expect(requests.map((request) => request.body.exact)).toEqual([entries, [entries[0]], [entries[1]], [entries[2]]]);
expect(await members(kvClient, EXACT_KEY)).toEqual(sorted([entries[3]!, entries[4]!]));
expect(await members(kvClient, REJECTED_KEY)).toEqual([]);
expect(requests).toHaveLength(1);
expect(await members(kvClient, QUEUE_KEY)).toEqual([]);
expect(await members(kvClient, REJECTED_KEY)).toEqual(['media.test/attachments/1/2/a', 'media.test/emojis/9']);
});
it('leaves the queue untouched when the adapter is none', async () => {
Config.cachePurge.adapter = 'none';
const {kvClient, queue} = createHarness();
const exact = [`${MEDIA}/avatars/1/abc`];
const prefix = [`${MEDIA}/some/dir/`];
await queue.addUrls([...exact, ...prefix]);
await queue.addUrls([`${MEDIA}/attachments/1/2/a.png`]);
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
expect(requests).toEqual([]);
expect(await members(kvClient, EXACT_KEY)).toEqual(exact);
expect(await members(kvClient, PREFIX_KEY)).toEqual(prefix);
expect(await kvClient.get('cache_purge:budget:exact')).toBeNull();
expect(await kvClient.get('cache_purge:budget:prefix')).toBeNull();
expect(await members(kvClient, QUEUE_KEY)).toEqual(['media.test/attachments/1/2/a']);
expect(await kvClient.get(BUDGET_KEY)).toBeNull();
});
it('sends no more than the token bucket allows and refills at the configured rate', async () => {
vi.useFakeTimers({toFake: ['Date']});
vi.setSystemTime(new Date('2026-09-13T00:00:00.000Z'));
vi.setSystemTime(new Date('2026-09-16T00:00:00.000Z'));
const {kvClient, queue} = createHarness();
await queue.addUrls([
...Array.from({length: 200}, (_, index) => `${MEDIA}/avatars/${index}/abc`),
...Array.from({length: 30}, (_, index) => `${MEDIA}/dirs/${index}/`),
]);
await queue.addUrls(Array.from({length: 200}, (_, index) => `${MEDIA}/attachments/1/${index}/a.png`));
const requests = recordPurges(() => new HttpResponse(null, {status: 204}));
await processCachePurgeQueue({}, HELPERS);
@@ -305,11 +247,7 @@ describe('processCachePurgeQueue', () => {
vi.setSystemTime(Date.now() + 10_000);
await processCachePurgeQueue({}, HELPERS);
expect(requests.map((request) => [request.body.exact.length, request.body.prefix.length])).toEqual([
[120, 20],
[50, 5],
]);
expect(await kvClient.scard(EXACT_KEY)).toBe(30);
expect(await kvClient.scard(PREFIX_KEY)).toBe(5);
expect(requests.map((request) => request.body.prefixes.length)).toEqual([120, 50]);
expect(await kvClient.scard(QUEUE_KEY)).toBe(30);
});
});
+2 -1
View File
@@ -484,7 +484,8 @@ function validateCachePurgeConfig(config: MasterConfig): void {
) {
throw new Error('FLUXER_CACHE_PURGE_HTTP_ENDPOINT must be an absolute http or https URL without credentials');
}
if (!/^[\x21-\x7e]*$/u.test(cachePurge.http.token)) {
requireString(cachePurge.http.token, 'FLUXER_CACHE_PURGE_HTTP_TOKEN');
if (!/^[\x21-\x7e]+$/u.test(cachePurge.http.token)) {
throw new Error('FLUXER_CACHE_PURGE_HTTP_TOKEN must contain only visible ASCII characters');
}
assertIntegerInRange(cachePurge.http.timeout_ms, 'FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS', 1_000, 10_000);
@@ -572,6 +572,14 @@ describe('ConfigLoader', () => {
await expect(loadConfig()).rejects.toThrow('FLUXER_CACHE_PURGE_HTTP_ENDPOINT is required');
});
test('rejects the http cache purge adapter without a token', async () => {
stubMinimalEnv({
FLUXER_CACHE_PURGE_ADAPTER: 'http',
FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'https://purge.internal/purge',
});
await expect(loadConfig()).rejects.toThrow('FLUXER_CACHE_PURGE_HTTP_TOKEN is required');
});
test('rejects a cache purge endpoint that is not an absolute http URL', async () => {
for (const endpoint of ['/purge', 'purge.internal/purge', 'ftp://purge.internal/purge']) {
stubMinimalEnv({FLUXER_CACHE_PURGE_ADAPTER: 'http', FLUXER_CACHE_PURGE_HTTP_ENDPOINT: endpoint});
@@ -596,6 +604,7 @@ describe('ConfigLoader', () => {
stubMinimalEnv({
FLUXER_CACHE_PURGE_ADAPTER: 'http',
FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'https://purge.internal/purge',
FLUXER_CACHE_PURGE_HTTP_TOKEN: 'purge-token',
FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS: timeout,
});
await expect(loadConfig()).rejects.toThrow(