From ca7ddd92725e02ee909a486af68d2241d9ec345f Mon Sep 17 00:00:00 2001 From: Hampus Date: Sun, 13 Sep 2026 20:28:56 +0200 Subject: [PATCH] refactor(api): drop Bunny for pluggable cache purge adapters (#2742) --- fluxer_api/src/api/Config.ts | 11 +- .../AdminSyntheticUserServiceGuards.test.ts | 2 +- .../services/AttachmentUploadService.ts | 8 +- .../channel/services/ChannelDataService.ts | 2 +- .../api/channel/services/ChannelService.ts | 2 +- .../api/channel/services/MessageService.ts | 2 +- .../channel_data/ChannelUtilsService.ts | 2 +- .../services/message/MessageDeleteService.ts | 2 +- .../services/message/MessageHelpers.test.ts | 112 +++++ .../services/message/MessageHelpers.ts | 12 +- .../message/UserMessageDeletionService.ts | 2 +- fluxer_api/src/api/config/APIConfig.ts | 16 +- .../src/api/csam/NcmecSubmissionService.ts | 2 +- .../src/api/infrastructure/BunnyPurgeQueue.ts | 124 ----- .../api/infrastructure/CachePurgeAdapter.ts | 29 ++ .../src/api/infrastructure/CachePurgeQueue.ts | 122 +++++ .../infrastructure/HttpCachePurgeAdapter.ts | 40 ++ .../infrastructure/NoneCachePurgeAdapter.ts | 9 + .../middleware/GuildStackServiceFactory.ts | 2 +- .../src/api/middleware/ServiceSingletons.ts | 4 +- .../test/msw/handlers/BunnyEdgeHandlers.ts | 27 - fluxer_api/src/api/test/msw/server.ts | 2 - .../api/user/services/UserDeletionService.ts | 2 +- .../src/api/utils/ExternalResponseLimits.ts | 1 - .../src/api/worker/WorkerDependencies.ts | 2 +- fluxer_api/src/api/worker/WorkerLaneConfig.ts | 2 +- fluxer_api/src/api/worker/WorkerMain.ts | 6 +- .../src/api/worker/WorkerTaskRegistry.ts | 4 +- .../worker/tasks/ProcessAssetDeletionQueue.ts | 2 +- .../worker/tasks/ProcessBunnyPurgeQueue.ts | 152 ------ .../worker/tasks/ProcessCachePurgeQueue.ts | 104 ++++ .../tests/ProcessAssetDeletionQueue.test.ts | 2 +- .../tests/ProcessCachePurgeQueue.test.ts | 315 ++++++++++++ .../content/docs/operator/configuration.mdx | 30 +- fluxer_media_proxy/src/bunny_ip_gate.rs | 469 ------------------ fluxer_media_proxy/src/config/mod.rs | 25 +- fluxer_media_proxy/src/config/parse.rs | 19 - fluxer_media_proxy/src/config/tests/mod.rs | 33 +- fluxer_media_proxy/src/lib.rs | 1 - fluxer_media_proxy/src/server/runtime.rs | 57 +-- fluxer_media_proxy/src/storage/tests/mod.rs | 3 - packages/config/src/ConfigLoader.ts | 36 +- packages/config/src/MasterConfig.ts | 13 +- .../config/src/__tests__/ConfigLoader.test.ts | 84 ++++ .../src/config_loader/EnvironmentOverrides.ts | 10 +- 45 files changed, 929 insertions(+), 977 deletions(-) create mode 100644 fluxer_api/src/api/channel/services/message/MessageHelpers.test.ts delete mode 100644 fluxer_api/src/api/infrastructure/BunnyPurgeQueue.ts create mode 100644 fluxer_api/src/api/infrastructure/CachePurgeAdapter.ts create mode 100644 fluxer_api/src/api/infrastructure/CachePurgeQueue.ts create mode 100644 fluxer_api/src/api/infrastructure/HttpCachePurgeAdapter.ts create mode 100644 fluxer_api/src/api/infrastructure/NoneCachePurgeAdapter.ts delete mode 100644 fluxer_api/src/api/test/msw/handlers/BunnyEdgeHandlers.ts delete mode 100644 fluxer_api/src/api/worker/tasks/ProcessBunnyPurgeQueue.ts create mode 100644 fluxer_api/src/api/worker/tasks/ProcessCachePurgeQueue.ts create mode 100644 fluxer_api/src/api/worker/tests/ProcessCachePurgeQueue.test.ts delete mode 100644 fluxer_media_proxy/src/bunny_ip_gate.rs diff --git a/fluxer_api/src/api/Config.ts b/fluxer_api/src/api/Config.ts index de099e1b3..d98d4ed9f 100644 --- a/fluxer_api/src/api/Config.ts +++ b/fluxer_api/src/api/Config.ts @@ -409,10 +409,13 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig { : undefined, legacyPrices: master.integrations.stripe.legacy_prices, }, - bunny: { - purgeEnabled: master.integrations.bunny.purge_enabled, - apiKey: master.integrations.bunny.api_key, - pullZoneId: master.integrations.bunny.pull_zone_id, + cachePurge: { + adapter: master.integrations.cache_purge.adapter, + http: { + endpoint: master.integrations.cache_purge.http.endpoint, + token: master.integrations.cache_purge.http.token, + timeoutMs: master.integrations.cache_purge.http.timeout_ms, + }, }, clamav: { enabled: master.integrations.clamav.enabled, diff --git a/fluxer_api/src/api/admin/tests/AdminSyntheticUserServiceGuards.test.ts b/fluxer_api/src/api/admin/tests/AdminSyntheticUserServiceGuards.test.ts index 13cceb774..526cd520f 100644 --- a/fluxer_api/src/api/admin/tests/AdminSyntheticUserServiceGuards.test.ts +++ b/fluxer_api/src/api/admin/tests/AdminSyntheticUserServiceGuards.test.ts @@ -5,7 +5,7 @@ import type {IChannelRepository} from '@app/api/channel/IChannelRepository'; import {UserMessageDeletionService} from '@app/api/channel/services/message/UserMessageDeletionService'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; import {GuildMemberOperationsService} from '@app/api/guild/services/member/GuildMemberOperationsService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import {DELETED_USER_ID} from '@fluxer/constants/src/UserConstants'; diff --git a/fluxer_api/src/api/channel/services/AttachmentUploadService.ts b/fluxer_api/src/api/channel/services/AttachmentUploadService.ts index fcbc9d6e3..4afe4d660 100644 --- a/fluxer_api/src/api/channel/services/AttachmentUploadService.ts +++ b/fluxer_api/src/api/channel/services/AttachmentUploadService.ts @@ -21,7 +21,7 @@ import { } from '@app/api/channel/services/message/MessageHelpers'; import {applyUploadRelayDecision, resolveUploadRelayDecision} from '@app/api/channel/services/UploadRelay'; import {SYSTEM_USER_ID} from '@app/api/constants/Core'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import type {LimitConfigService} from '@app/api/limits/LimitConfigService'; @@ -387,10 +387,8 @@ export class AttachmentUploadService { } const cdnKey = makeAttachmentCdnKey(message.channelId, attachment.id, attachment.filename); await this.storageService.deleteObject(Config.s3.buckets.cdn, cdnKey); - if (Config.bunny.purgeEnabled) { - const cdnUrl = makeAttachmentCdnUrl(message.channelId, attachment.id, attachment.filename); - await this.purgeQueue.addUrls([cdnUrl]); - } + const cdnUrl = makeAttachmentCdnUrl(message.channelId, attachment.id, attachment.filename); + await this.purgeQueue.addUrls([cdnUrl]); const updatedAttachments = message.attachments.filter((a: Attachment) => a.id !== attachmentId); const updatedRowData = { ...message.toRow(), diff --git a/fluxer_api/src/api/channel/services/ChannelDataService.ts b/fluxer_api/src/api/channel/services/ChannelDataService.ts index 18005746b..7816c7e1a 100644 --- a/fluxer_api/src/api/channel/services/ChannelDataService.ts +++ b/fluxer_api/src/api/channel/services/ChannelDataService.ts @@ -12,7 +12,7 @@ import type {MessagePersistenceService} from '@app/api/channel/services/message/ import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; import type {AvatarService} from '@app/api/infrastructure/AvatarService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {ILiveKitService} from '@app/api/infrastructure/ILiveKitService'; import type {ISnowflakeService} from '@app/api/infrastructure/ISnowflakeService'; diff --git a/fluxer_api/src/api/channel/services/ChannelService.ts b/fluxer_api/src/api/channel/services/ChannelService.ts index c5b559d8d..7ccdcede9 100644 --- a/fluxer_api/src/api/channel/services/ChannelService.ts +++ b/fluxer_api/src/api/channel/services/ChannelService.ts @@ -16,7 +16,7 @@ import type {IFavoriteMemeRepository} from '@app/api/favorite_meme/IFavoriteMeme import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; import type {AvatarService} from '@app/api/infrastructure/AvatarService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {EmbedService} from '@app/api/infrastructure/EmbedService'; import type {ILiveKitService} from '@app/api/infrastructure/ILiveKitService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; diff --git a/fluxer_api/src/api/channel/services/MessageService.ts b/fluxer_api/src/api/channel/services/MessageService.ts index 110d149d6..cc92001da 100644 --- a/fluxer_api/src/api/channel/services/MessageService.ts +++ b/fluxer_api/src/api/channel/services/MessageService.ts @@ -20,7 +20,7 @@ import {MessageValidationService} from '@app/api/channel/services/message/Messag import type {IFavoriteMemeRepository} from '@app/api/favorite_meme/IFavoriteMemeRepository'; import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IMediaService} from '@app/api/infrastructure/IMediaService'; import type {ISnowflakeService} from '@app/api/infrastructure/ISnowflakeService'; diff --git a/fluxer_api/src/api/channel/services/channel_data/ChannelUtilsService.ts b/fluxer_api/src/api/channel/services/channel_data/ChannelUtilsService.ts index 1600a9b84..665b7adc3 100644 --- a/fluxer_api/src/api/channel/services/channel_data/ChannelUtilsService.ts +++ b/fluxer_api/src/api/channel/services/channel_data/ChannelUtilsService.ts @@ -6,7 +6,7 @@ import type {IChannelRepositoryAggregate} from '@app/api/channel/repositories/IC import {dispatchChannelEvent} from '@app/api/channel/services/ChannelGatewayDispatch'; import {dispatchMessageCreateBroadcast} from '@app/api/channel/services/message/MessageGatewayDispatch'; import {purgeMessageAttachments} from '@app/api/channel/services/message/MessageHelpers'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import type {UserCacheService} from '@app/api/infrastructure/UserCacheService'; diff --git a/fluxer_api/src/api/channel/services/message/MessageDeleteService.ts b/fluxer_api/src/api/channel/services/message/MessageDeleteService.ts index 313e2f007..9b0898e2d 100644 --- a/fluxer_api/src/api/channel/services/message/MessageDeleteService.ts +++ b/fluxer_api/src/api/channel/services/message/MessageDeleteService.ts @@ -9,7 +9,7 @@ import {isOperationDisabled, purgeMessageAttachments} from '@app/api/channel/ser import type {MessageSearchService} from '@app/api/channel/services/message/MessageSearchService'; import type {MessageValidationService} from '@app/api/channel/services/message/MessageValidationService'; import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import type {RequestCache} from '@app/api/middleware/RequestCacheMiddleware'; diff --git a/fluxer_api/src/api/channel/services/message/MessageHelpers.test.ts b/fluxer_api/src/api/channel/services/message/MessageHelpers.test.ts new file mode 100644 index 000000000..74113ca10 --- /dev/null +++ b/fluxer_api/src/api/channel/services/message/MessageHelpers.test.ts @@ -0,0 +1,112 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createAttachmentID, createChannelID, createMessageID, createUserID} from '@app/api/BrandedTypes'; +import {Config} from '@app/api/Config'; +import {purgeMessageAttachments} from '@app/api/channel/services/message/MessageHelpers'; +import type {MessageEmbed} from '@app/api/database/types/MessageTypes'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; +import type {IStorageService} from '@app/api/infrastructure/IStorageService'; +import {Message} from '@app/api/models/Message'; +import {MessageTypes} from '@fluxer/constants/src/ChannelConstants'; +import {describe, expect, it} from 'vitest'; + +const CHANNEL_ID = createChannelID(10n); +const ATTACHMENT_KEY = 'attachments/10/200/ação.png'; +const OTHER_MESSAGE_ATTACHMENT_KEY = 'attachments/11/300/photo.jpg'; + +function imageEmbed(key: string): MessageEmbed { + return { + type: 'image', + title: null, + description: null, + url: null, + timestamp: null, + color: null, + author: null, + provider: null, + thumbnail: null, + image: { + url: `${Config.endpoints.media}/${key}`, + width: 16, + height: 16, + description: null, + content_type: 'image/png', + content_hash: null, + placeholder: null, + flags: 0, + duration: null, + }, + video: null, + footer: null, + fields: null, + nsfw: null, + }; +} + +function makeMessageWithMedia(): Message { + return new Message({ + channel_id: CHANNEL_ID, + bucket: 0, + message_id: createMessageID(100n), + author_id: createUserID(3n), + type: MessageTypes.DEFAULT, + webhook_id: null, + webhook_name: null, + webhook_avatar_hash: null, + content: '', + edited_timestamp: null, + pinned_timestamp: null, + flags: 0, + mention_everyone: false, + mention_users: null, + mention_roles: null, + mention_channels: null, + attachments: [ + { + attachment_id: createAttachmentID(200n), + filename: 'ação.png', + size: 1024n, + title: null, + description: null, + width: 16, + height: 16, + content_type: 'image/png', + content_hash: null, + placeholder: null, + flags: 0, + duration: null, + nsfw: null, + waveform: null, + }, + ], + embeds: [imageEmbed(ATTACHMENT_KEY), imageEmbed(OTHER_MESSAGE_ATTACHMENT_KEY)], + sticker_items: null, + message_reference: null, + message_snapshots: null, + call: null, + has_reaction: null, + version: 1, + }); +} + +describe('purgeMessageAttachments', () => { + it('queues each stored media URL of a deleted message once and leaves media of other messages alone', async () => { + const deletedObjects: Array = []; + const queuedUrls: Array = []; + const storageService = { + deleteObject: async (bucket: string, key: string) => { + deletedObjects.push(`${bucket}/${key}`); + }, + } as unknown as IStorageService; + const purgeQueue: IPurgeQueue = { + addUrls: async (urls) => { + queuedUrls.push(...urls); + }, + }; + + await purgeMessageAttachments(makeMessageWithMedia(), storageService, purgeQueue); + + expect(deletedObjects).toEqual([`${Config.s3.buckets.cdn}/${ATTACHMENT_KEY}`]); + expect(queuedUrls).toEqual([`${Config.endpoints.media}/${ATTACHMENT_KEY}`]); + }); +}); diff --git a/fluxer_api/src/api/channel/services/message/MessageHelpers.ts b/fluxer_api/src/api/channel/services/message/MessageHelpers.ts index bf1b6ef62..0fe6edd65 100644 --- a/fluxer_api/src/api/channel/services/message/MessageHelpers.ts +++ b/fluxer_api/src/api/channel/services/message/MessageHelpers.ts @@ -7,7 +7,7 @@ import type { MessageSnapshot as CassandraMessageSnapshot, MessageAttachment, } from '@app/api/database/types/MessageTypes'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {ISnowflakeService} from '@app/api/infrastructure/ISnowflakeService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import {Logger} from '@app/api/Logger'; @@ -353,21 +353,17 @@ export async function purgeMessageAttachments( continue; } cdnKeys.add(cdnKey); - if (Config.bunny.purgeEnabled) { - cdnUrls.push(makeAttachmentCdnUrl(message.channelId, attachment.id, attachment.filename)); - } + cdnUrls.push(makeAttachmentCdnUrl(message.channelId, attachment.id, attachment.filename)); } for (const embedKey of collectEmbedReferencedAttachmentCdnKeys(message, ownedCdnKeys)) { if (cdnKeys.has(embedKey)) { continue; } cdnKeys.add(embedKey); - if (Config.bunny.purgeEnabled) { - cdnUrls.push(`${Config.endpoints.media}/${embedKey}`); - } + cdnUrls.push(`${Config.endpoints.media}/${embedKey}`); } await Promise.all([...cdnKeys].map((cdnKey) => storageService.deleteObject(Config.s3.buckets.cdn, cdnKey))); - if (Config.bunny.purgeEnabled && cdnUrls.length > 0) { + if (cdnUrls.length > 0) { await purgeQueue.addUrls(cdnUrls); } } diff --git a/fluxer_api/src/api/channel/services/message/UserMessageDeletionService.ts b/fluxer_api/src/api/channel/services/message/UserMessageDeletionService.ts index f887db807..b3cf03f38 100644 --- a/fluxer_api/src/api/channel/services/message/UserMessageDeletionService.ts +++ b/fluxer_api/src/api/channel/services/message/UserMessageDeletionService.ts @@ -11,7 +11,7 @@ import { type SelfMessageFilter, } from '@app/api/channel/services/message/SelfMessageFilter'; import {assertMutableUserId} from '@app/api/constants/Core'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import {Logger} from '@app/api/Logger'; diff --git a/fluxer_api/src/api/config/APIConfig.ts b/fluxer_api/src/api/config/APIConfig.ts index b77971ea7..ce48863b4 100644 --- a/fluxer_api/src/api/config/APIConfig.ts +++ b/fluxer_api/src/api/config/APIConfig.ts @@ -1,6 +1,7 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import type {WorkerTaskName} from '@app/api/worker/WorkerLaneConfig'; +import type {CachePurgeAdapterName} from '@fluxer/config/src/MasterConfig'; import type {ResolvedDownloadsProvider} from '@fluxer/config/src/S3DownloadsProvider'; export type APIWorkerMode = 'all_lanes' | 'single_lane' | 'single_task'; @@ -14,6 +15,15 @@ export interface PushProviderAppConfig { projectId?: string; } +export interface APICachePurgeConfig { + adapter: CachePurgeAdapterName; + http: { + endpoint: string; + token: string; + timeoutMs: number; + }; +} + interface APIGeoipFilesystemConfig { mode: 'filesystem'; maxmindDbPath?: string; @@ -262,11 +272,7 @@ export interface APIConfig { }; legacyPrices?: Record | undefined>; }; - bunny: { - purgeEnabled: boolean; - apiKey?: string; - pullZoneId?: number; - }; + cachePurge: APICachePurgeConfig; clamav: { enabled: boolean; host: string; diff --git a/fluxer_api/src/api/csam/NcmecSubmissionService.ts b/fluxer_api/src/api/csam/NcmecSubmissionService.ts index f0ceced13..1becd0913 100644 --- a/fluxer_api/src/api/csam/NcmecSubmissionService.ts +++ b/fluxer_api/src/api/csam/NcmecSubmissionService.ts @@ -29,7 +29,7 @@ import type {NcmecRepository} from '@app/api/csam/NcmecRepository'; import type {AttachmentUploadTraceByAttachmentRow} from '@app/api/database/types/AttachmentUploadTypes'; import type {NcmecAttachmentSubmissionRow, NcmecUserWorkflowRow} from '@app/api/database/types/CsamTypes'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import type {KVAccountDeletionQueueService} from '@app/api/infrastructure/KVAccountDeletionQueueService'; diff --git a/fluxer_api/src/api/infrastructure/BunnyPurgeQueue.ts b/fluxer_api/src/api/infrastructure/BunnyPurgeQueue.ts deleted file mode 100644 index cf7ad62d3..000000000 --- a/fluxer_api/src/api/infrastructure/BunnyPurgeQueue.ts +++ /dev/null @@ -1,124 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -import {Logger} from '@app/api/Logger'; -import type {IKVProvider, KVPurgeBatchResult} from '@pkgs/kv_client/src/IKVProvider'; - -export interface IPurgeQueue { - addUrls(urls: Array): Promise; - getQueueSize(): Promise; - clear(): Promise; -} - -const EXACT_QUEUE_KEY = 'bunny:purge:exact'; -const PREFIX_QUEUE_KEY = 'bunny:purge:prefix'; -const EXACT_BUCKET_KEY = 'bunny:purge:budget:exact'; -const PREFIX_BUCKET_KEY = 'bunny:purge:budget:prefix'; -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('/'); -} - -export class BunnyPurgeQueue implements IPurgeQueue { - private readonly kvClient: IKVProvider; - - constructor(kvClient: IKVProvider) { - this.kvClient = kvClient; - } - - async addUrls(urls: Array): Promise { - if (urls.length === 0) { - return; - } - const exactUrls: Array = []; - const prefixUrls: Array = []; - for (const url of urls) { - const trimmed = url.trim(); - if (trimmed === '') { - continue; - } - if (isPrefix(trimmed)) { - prefixUrls.push(trimmed); - } else { - exactUrls.push(trimmed); - } - } - try { - const ops: Array> = []; - if (exactUrls.length > 0) { - ops.push(this.kvClient.sadd(EXACT_QUEUE_KEY, ...exactUrls)); - } - if (prefixUrls.length > 0) { - ops.push(this.kvClient.sadd(PREFIX_QUEUE_KEY, ...prefixUrls)); - } - await Promise.all(ops); - Logger.debug({exact: exactUrls.length, prefix: prefixUrls.length}, 'Added URLs to CDN purge queue'); - } catch (error) { - Logger.error( - {error, exact: exactUrls.length, prefix: prefixUrls.length}, - 'Failed to add URLs to CDN purge queue', - ); - throw error; - } - } - - async getQueueSize(): Promise { - try { - const [exactSize, prefixSize] = await Promise.all([ - this.kvClient.scard(EXACT_QUEUE_KEY), - this.kvClient.scard(PREFIX_QUEUE_KEY), - ]); - return exactSize + prefixSize; - } catch (error) { - Logger.error({error}, 'Failed to get CDN purge queue size'); - throw error; - } - } - - async clear(): Promise { - try { - await this.kvClient.del(EXACT_QUEUE_KEY, PREFIX_QUEUE_KEY); - Logger.debug('Cleared CDN purge queue'); - } catch (error) { - Logger.error({error}, 'Failed to clear CDN purge queue'); - throw error; - } - } - - async dequeueExactBatch(maxItems: number): Promise { - return this.kvClient.dequeuePurgeBatch( - EXACT_QUEUE_KEY, - EXACT_BUCKET_KEY, - maxItems, - EXACT_MAX_TOKENS, - EXACT_REFILL_RATE, - EXACT_REFILL_INTERVAL_MS, - ); - } - - async dequeuePrefixBatch(maxItems: number): Promise { - return this.kvClient.dequeuePurgeBatch( - PREFIX_QUEUE_KEY, - PREFIX_BUCKET_KEY, - maxItems, - PREFIX_MAX_TOKENS, - PREFIX_REFILL_RATE, - PREFIX_REFILL_INTERVAL_MS, - ); - } -} - -export class NoopPurgeQueue implements IPurgeQueue { - async addUrls(_urls: Array): Promise {} - - async getQueueSize(): Promise { - return 0; - } - - async clear(): Promise {} -} diff --git a/fluxer_api/src/api/infrastructure/CachePurgeAdapter.ts b/fluxer_api/src/api/infrastructure/CachePurgeAdapter.ts new file mode 100644 index 000000000..71b96dd40 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/CachePurgeAdapter.ts @@ -0,0 +1,29 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {APICachePurgeConfig} from '@app/api/config/APIConfig'; +import {createHttpCachePurgeAdapter} from '@app/api/infrastructure/HttpCachePurgeAdapter'; +import {createNoneCachePurgeAdapter} from '@app/api/infrastructure/NoneCachePurgeAdapter'; +import type {CachePurgeAdapterName} from '@fluxer/config/src/MasterConfig'; + +export interface CachePurgeBatch { + readonly exact: ReadonlyArray; + readonly prefix: ReadonlyArray; +} + +export type CachePurgeOutcome = + | {readonly kind: 'purged'} + | {readonly kind: 'invalid_entries'; readonly status: number} + | {readonly kind: 'failed'; readonly status: number | null; readonly error: unknown}; + +export interface CachePurgeAdapter { + purge(batch: CachePurgeBatch): Promise; +} + +const CACHE_PURGE_ADAPTERS = { + none: createNoneCachePurgeAdapter, + http: createHttpCachePurgeAdapter, +} satisfies Record CachePurgeAdapter>; + +export function createCachePurgeAdapter(config: APICachePurgeConfig): CachePurgeAdapter { + return CACHE_PURGE_ADAPTERS[config.adapter](config); +} diff --git a/fluxer_api/src/api/infrastructure/CachePurgeQueue.ts b/fluxer_api/src/api/infrastructure/CachePurgeQueue.ts new file mode 100644 index 000000000..cdffe0b76 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/CachePurgeQueue.ts @@ -0,0 +1,122 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {CachePurgeBatch} from '@app/api/infrastructure/CachePurgeAdapter'; +import {Logger} from '@app/api/Logger'; +import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider'; + +export interface IPurgeQueue { + addUrls(urls: Array): Promise; +} + +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 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('/'); +} + +export class CachePurgeQueue implements IPurgeQueue { + private readonly kvClient: IKVProvider; + + constructor(kvClient: IKVProvider) { + this.kvClient = kvClient; + } + + async addUrls(urls: Array): Promise { + if (urls.length === 0) { + return; + } + const exactUrls: Array = []; + const prefixUrls: Array = []; + for (const url of urls) { + const trimmed = url.trim(); + if (trimmed === '') { + continue; + } + if (isPrefix(trimmed)) { + prefixUrls.push(trimmed); + } else { + exactUrls.push(trimmed); + } + } + try { + await this.addToSets({exact: exactUrls, prefix: prefixUrls}); + Logger.debug({exact: exactUrls.length, prefix: prefixUrls.length}, 'Added URLs to cache purge queue'); + } catch (error) { + Logger.error( + {error, exact: exactUrls.length, prefix: prefixUrls.length}, + 'Failed to add URLs to cache purge queue', + ); + throw error; + } + } + + async dequeueBatch(): Promise { + 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, + ); + 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}; + } catch (error) { + await this.requeue({exact: exact.urls, prefix: []}); + throw error; + } + } + + async requeue(batch: CachePurgeBatch): Promise { + 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(batch: CachePurgeBatch): Promise { + 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 { + const ops: Array> = []; + 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); + } +} + +export class NoopPurgeQueue implements IPurgeQueue { + async addUrls(_urls: Array): Promise {} +} diff --git a/fluxer_api/src/api/infrastructure/HttpCachePurgeAdapter.ts b/fluxer_api/src/api/infrastructure/HttpCachePurgeAdapter.ts new file mode 100644 index 000000000..507821922 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/HttpCachePurgeAdapter.ts @@ -0,0 +1,40 @@ +// 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 * 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 { + const headers: Record = {'Content-Type': 'application/json'}; + if (config.http.token !== '') { + headers.Authorization = `Bearer ${config.http.token}`; + } + let response: Response; + try { + response = await fetch(config.http.endpoint, { + method: 'POST', + headers, + body: JSON.stringify({exact: batch.exact, prefix: batch.prefix.map(stripOneTrailingAsterisk)}), + redirect: 'manual', + signal: AbortSignal.timeout(config.http.timeoutMs), + }); + } catch (error) { + return {kind: 'failed', status: null, error}; + } + FetchUtils.discardResponseBody(response.body, response.status); + if (response.ok) { + return {kind: 'purged'}; + } + if (response.status === 400 || response.status === 422) { + return {kind: 'invalid_entries', status: response.status}; + } + return {kind: 'failed', status: response.status, error: null}; + }, + }; +} diff --git a/fluxer_api/src/api/infrastructure/NoneCachePurgeAdapter.ts b/fluxer_api/src/api/infrastructure/NoneCachePurgeAdapter.ts new file mode 100644 index 000000000..d9ccf08cd --- /dev/null +++ b/fluxer_api/src/api/infrastructure/NoneCachePurgeAdapter.ts @@ -0,0 +1,9 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {CachePurgeAdapter} from '@app/api/infrastructure/CachePurgeAdapter'; + +export function createNoneCachePurgeAdapter(): CachePurgeAdapter { + return { + purge: async () => ({kind: 'purged'}), + }; +} diff --git a/fluxer_api/src/api/middleware/GuildStackServiceFactory.ts b/fluxer_api/src/api/middleware/GuildStackServiceFactory.ts index 8fad152ee..79ef0f4fa 100644 --- a/fluxer_api/src/api/middleware/GuildStackServiceFactory.ts +++ b/fluxer_api/src/api/middleware/GuildStackServiceFactory.ts @@ -9,7 +9,7 @@ import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; import type {IGuildRepositoryAggregate} from '@app/api/guild/repositories/IGuildRepositoryAggregate'; import {GuildService} from '@app/api/guild/services/GuildService'; import type {AvatarService} from '@app/api/infrastructure/AvatarService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {EmbedService} from '@app/api/infrastructure/EmbedService'; import type {EntityAssetService} from '@app/api/infrastructure/EntityAssetService'; import type {IAssetDeletionQueue} from '@app/api/infrastructure/IAssetDeletionQueue'; diff --git a/fluxer_api/src/api/middleware/ServiceSingletons.ts b/fluxer_api/src/api/middleware/ServiceSingletons.ts index a754f66c0..5dcc3af5c 100644 --- a/fluxer_api/src/api/middleware/ServiceSingletons.ts +++ b/fluxer_api/src/api/middleware/ServiceSingletons.ts @@ -31,7 +31,7 @@ import {GuildRepository} from '@app/api/guild/repositories/GuildRepository'; import {GuildDiscoveryService} from '@app/api/guild/services/GuildDiscoveryService'; import {AssetDeletionQueue} from '@app/api/infrastructure/AssetDeletionQueue'; import {AvatarService} from '@app/api/infrastructure/AvatarService'; -import {BunnyPurgeQueue, type IPurgeQueue, NoopPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import {CachePurgeQueue, type IPurgeQueue, NoopPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import {DisabledVirusScanService} from '@app/api/infrastructure/DisabledVirusScanService'; import {DiscriminatorService} from '@app/api/infrastructure/DiscriminatorService'; import {EmailDnsValidationService} from '@app/api/infrastructure/EmailDnsValidationService'; @@ -240,7 +240,7 @@ export const getLimitConfigService = singleton( }, ); export const getPurgeQueue: () => IPurgeQueue = singleton(() => - Config.bunny.purgeEnabled ? new BunnyPurgeQueue(getKVClient()) : new NoopPurgeQueue(), + Config.cachePurge.adapter === 'none' ? new NoopPurgeQueue() : new CachePurgeQueue(getKVClient()), ); export const getAssetDeletionQueue: () => IAssetDeletionQueue = singleton(() => new AssetDeletionQueue(getKVClient())); diff --git a/fluxer_api/src/api/test/msw/handlers/BunnyEdgeHandlers.ts b/fluxer_api/src/api/test/msw/handlers/BunnyEdgeHandlers.ts deleted file mode 100644 index d2c5f5662..000000000 --- a/fluxer_api/src/api/test/msw/handlers/BunnyEdgeHandlers.ts +++ /dev/null @@ -1,27 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -import {HttpResponse, http} from 'msw'; - -const BUNNY_EDGE_IPV4_LIST = ['198.51.100.10', '198.51.100.11'].join('\n'); -const BUNNY_EDGE_IPV6_LIST = ['2001:db8::10', '2001:db8::11'].join('\n'); - -export function createBunnyEdgeHandlers() { - return [ - http.get( - 'https://bunnycdn.com/api/system/edgeserverlist/plain', - () => - new HttpResponse(BUNNY_EDGE_IPV4_LIST, { - status: 200, - headers: {'content-type': 'text/plain; charset=utf-8'}, - }), - ), - http.get( - 'https://bunnycdn.com/api/system/edgeserverlist/IPv6/plain', - () => - new HttpResponse(BUNNY_EDGE_IPV6_LIST, { - status: 200, - headers: {'content-type': 'text/plain; charset=utf-8'}, - }), - ), - ]; -} diff --git a/fluxer_api/src/api/test/msw/server.ts b/fluxer_api/src/api/test/msw/server.ts index 79407da08..8fc13b296 100644 --- a/fluxer_api/src/api/test/msw/server.ts +++ b/fluxer_api/src/api/test/msw/server.ts @@ -1,6 +1,5 @@ // SPDX-License-Identifier: AGPL-3.0-or-later -import {createBunnyEdgeHandlers} from '@app/api/test/msw/handlers/BunnyEdgeHandlers'; import {createIpInfoLookupHandler} from '@app/api/test/msw/handlers/IpInfoHandlers'; import {createNcmecHandlers} from '@app/api/test/msw/handlers/NcmecHandlers'; import {createOnionooDetailsHandler} from '@app/api/test/msw/handlers/OnionooHandlers'; @@ -9,7 +8,6 @@ import {createPwnedPasswordsRangeHandler} from '@app/api/test/msw/handlers/Pwned import {setupServer} from 'msw/node'; export const server = setupServer( - ...createBunnyEdgeHandlers(), ...createNcmecHandlers(), ...createOpenNsfwHandlers(), createIpInfoLookupHandler(), diff --git a/fluxer_api/src/api/user/services/UserDeletionService.ts b/fluxer_api/src/api/user/services/UserDeletionService.ts index 38ec74e7f..a570beea6 100644 --- a/fluxer_api/src/api/user/services/UserDeletionService.ts +++ b/fluxer_api/src/api/user/services/UserDeletionService.ts @@ -8,7 +8,7 @@ import type {ChannelRepository} from '@app/api/channel/ChannelRepository'; import type {IConnectionRepository} from '@app/api/connection/IConnectionRepository'; import type {FavoriteMemeRepository} from '@app/api/favorite_meme/FavoriteMemeRepository'; import type {GuildRepository} from '@app/api/guild/repositories/GuildRepository'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {DiscriminatorService} from '@app/api/infrastructure/DiscriminatorService'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; import type {ISnowflakeService} from '@app/api/infrastructure/ISnowflakeService'; diff --git a/fluxer_api/src/api/utils/ExternalResponseLimits.ts b/fluxer_api/src/api/utils/ExternalResponseLimits.ts index 909f62c2f..2f618ca4b 100644 --- a/fluxer_api/src/api/utils/ExternalResponseLimits.ts +++ b/fluxer_api/src/api/utils/ExternalResponseLimits.ts @@ -21,7 +21,6 @@ export const EXTERNAL_RESPONSE_LIMITS = { pwnedPasswordsBytes: 1024 * 1024, rdapBytes: 512 * 1024, externalTemplateBytes: 512 * 1024, - bunnyErrorBytes: 16 * 1024, ncmecResponseBytes: 64 * 1024, fileBlocklistBytes: 25 * 1024 * 1024, urlBlocklistBytes: 25 * 1024 * 1024, diff --git a/fluxer_api/src/api/worker/WorkerDependencies.ts b/fluxer_api/src/api/worker/WorkerDependencies.ts index 458afc28e..d054a7718 100644 --- a/fluxer_api/src/api/worker/WorkerDependencies.ts +++ b/fluxer_api/src/api/worker/WorkerDependencies.ts @@ -17,7 +17,7 @@ import type {GuildAuditLogService} from '@app/api/guild/GuildAuditLogService'; import type {GuildRepository} from '@app/api/guild/repositories/GuildRepository'; import type {GuildService} from '@app/api/guild/services/GuildService'; import type {AvatarService} from '@app/api/infrastructure/AvatarService'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {DiscriminatorService} from '@app/api/infrastructure/DiscriminatorService'; import type {EmbedService} from '@app/api/infrastructure/EmbedService'; import type {IAssetDeletionQueue} from '@app/api/infrastructure/IAssetDeletionQueue'; diff --git a/fluxer_api/src/api/worker/WorkerLaneConfig.ts b/fluxer_api/src/api/worker/WorkerLaneConfig.ts index f1fa2c2f9..fff9c7241 100644 --- a/fluxer_api/src/api/worker/WorkerLaneConfig.ts +++ b/fluxer_api/src/api/worker/WorkerLaneConfig.ts @@ -70,7 +70,7 @@ const LANE_CONFIG = { 'indexChannelMessages', 'indexGuildMembers', 'processAssetDeletionQueue', - 'processBunnyPurgeQueue', + 'processCachePurgeQueue', 'processExpiredPremiumSweep', 'processInactivityDeletions', 'processPendingBulkMessageDeletions', diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index 21cc46543..2ec2ca23f 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -54,8 +54,8 @@ const SEARCH_REQUIRED_TASKS = new Set([ function registerCronJobs(cron: CronScheduler): void { cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *', {ledger: false}); - if (Config.bunny.purgeEnabled) { - cron.upsert('processBunnyPurgeQueue', 'processBunnyPurgeQueue', {}, '*/10 * * * * *', {ledger: false}); + if (Config.cachePurge.adapter !== 'none') { + cron.upsert('processCachePurgeQueue', 'processCachePurgeQueue', {}, '*/10 * * * * *', {ledger: false}); } cron.upsert('processPendingBulkMessageDeletions', 'processPendingBulkMessageDeletions', {}, '0 */10 * * * *', { ledger: false, @@ -80,7 +80,7 @@ function registerCronJobs(cron: CronScheduler): void { Logger.info( { blocklistFeeds: Config.blocklistFeeds.enabled, - bunnyPurge: Config.bunny.purgeEnabled, + cachePurgeAdapter: Config.cachePurge.adapter, selfHosted: Config.instance.selfHosted, }, 'Cron jobs registered successfully', diff --git a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts index 6c0a927be..bf29dc4c0 100644 --- a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts +++ b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts @@ -25,7 +25,7 @@ import indexChannelMessages from '@app/api/worker/tasks/IndexChannelMessages'; import indexGuildMembers from '@app/api/worker/tasks/IndexGuildMembers'; import messageShred from '@app/api/worker/tasks/MessageShred'; import processAssetDeletionQueue from '@app/api/worker/tasks/ProcessAssetDeletionQueue'; -import processBunnyPurgeQueue from '@app/api/worker/tasks/ProcessBunnyPurgeQueue'; +import processCachePurgeQueue from '@app/api/worker/tasks/ProcessCachePurgeQueue'; import processExpiredPremiumSweep from '@app/api/worker/tasks/ProcessExpiredPremiumSweep'; import processInactivityDeletions from '@app/api/worker/tasks/ProcessInactivityDeletions'; import processPendingBulkMessageDeletions from '@app/api/worker/tasks/ProcessPendingBulkMessageDeletions'; @@ -69,7 +69,7 @@ export const workerTasks: Record = { indexGuildMembers, messageShred, processAssetDeletionQueue, - processBunnyPurgeQueue, + processCachePurgeQueue, processStripeWebhook, processExpiredPremiumSweep, processInactivityDeletions, diff --git a/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts index 75c422de6..6696f24ac 100644 --- a/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts +++ b/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts @@ -3,7 +3,7 @@ import {createGuildID, createUserID} from '@app/api/BrandedTypes'; import {Config} from '@app/api/Config'; import type {GuildRepository} from '@app/api/guild/repositories/GuildRepository'; -import type {IPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import type {IPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type { IAssetDeletionQueue, QueuedAssetDeletion, diff --git a/fluxer_api/src/api/worker/tasks/ProcessBunnyPurgeQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessBunnyPurgeQueue.ts deleted file mode 100644 index a7d69e5da..000000000 --- a/fluxer_api/src/api/worker/tasks/ProcessBunnyPurgeQueue.ts +++ /dev/null @@ -1,152 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -import {Config} from '@app/api/Config'; -import type {BunnyPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; -import {Logger} from '@app/api/Logger'; -import {EXTERNAL_RESPONSE_LIMITS} from '@app/api/utils/ExternalResponseLimits'; -import * as FetchUtils from '@app/api/utils/FetchUtils'; -import {getWorkerDependencies} from '@app/api/worker/WorkerContext'; -import {formatUrlForDiagnostics} from '@pkgs/http_client/src/HttpClientDiagnostics'; -import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask'; -import {ms} from 'itty-time'; - -const EXACT_BATCH_SIZE = 120; -const PREFIX_BATCH_SIZE = 20; - -type PurgeLabel = 'exact' | 'prefix'; - -interface PurgeStats { - attempted: number; - purged: number; - requeued: number; - failures: number; - rateLimited: number; -} - -function createEmptyPurgeStats(): PurgeStats { - return { - attempted: 0, - purged: 0, - requeued: 0, - failures: 0, - rateLimited: 0, - }; -} - -function mergePurgeStats(target: PurgeStats, source: PurgeStats): void { - target.attempted += source.attempted; - target.purged += source.purged; - target.requeued += source.requeued; - target.failures += source.failures; - target.rateLimited += source.rateLimited; -} - -async function purgeUrls( - urls: Array, - apiKey: string, - queue: BunnyPurgeQueue, - label: PurgeLabel, -): Promise { - const stats = createEmptyPurgeStats(); - for (let index = 0; index < urls.length; index++) { - const url = urls[index]!; - stats.attempted++; - try { - const response = await fetch(`https://api.bunny.net/purge?url=${encodeURIComponent(url)}&async=true`, { - method: 'POST', - headers: { - AccessKey: apiKey, - }, - signal: AbortSignal.timeout(ms('30 seconds')), - }); - if (response.status === 429) { - FetchUtils.discardResponseBody(response.body, response.status); - const remaining = urls.slice(index); - await queue.addUrls(remaining); - stats.requeued += remaining.length; - stats.rateLimited += remaining.length; - Logger.warn({label, rateLimitedUrls: remaining.length}, 'Rate limited by Bunny CDN, re-queued remaining URLs'); - break; - } - if (!response.ok) { - stats.failures++; - const errorBody = await FetchUtils.streamToBufferWithLimit(response.body, { - maxBytes: EXTERNAL_RESPONSE_LIMITS.bunnyErrorBytes, - headers: response.headers, - url: response.url, - description: 'Bunny CDN purge response', - }); - Logger.error( - { - status: response.status, - responseBodyBytes: errorBody.byteLength, - url: formatUrlForDiagnostics(url), - label, - }, - 'Failed to purge URL via Bunny CDN', - ); - await queue.addUrls([url]); - stats.requeued++; - continue; - } - FetchUtils.discardResponseBody(response.body, response.status); - stats.purged++; - } catch (error) { - stats.failures++; - Logger.error({error, url: formatUrlForDiagnostics(url), label}, 'Error purging URL via Bunny CDN'); - await queue.addUrls([url]); - stats.requeued++; - } - } - return stats; -} - -const processBunnyPurgeQueue: WorkerTaskHandler = async (_payload, _helpers) => { - if (!Config.bunny.purgeEnabled) { - Logger.debug('Bunny CDN cache purge is disabled, skipping queue processing'); - return; - } - if (!Config.bunny.apiKey || !Config.bunny.pullZoneId) { - Logger.error('Bunny CDN cache purge is enabled but credentials are missing'); - return; - } - const deps = getWorkerDependencies(); - const queue = deps.purgeQueue as BunnyPurgeQueue; - const apiKey = Config.bunny.apiKey; - const totalStats = createEmptyPurgeStats(); - try { - const queueSizeBefore = await queue.getQueueSize(); - if (queueSizeBefore === 0) { - Logger.debug('CDN purge queue is empty'); - return; - } - Logger.debug({queueSize: queueSizeBefore}, 'Processing CDN purge queue'); - const exactBatch = await queue.dequeueExactBatch(EXACT_BATCH_SIZE); - if (exactBatch.urls.length > 0) { - const exactStats = await purgeUrls(exactBatch.urls, apiKey, queue, 'exact'); - mergePurgeStats(totalStats, exactStats); - } - const prefixBatch = await queue.dequeuePrefixBatch(PREFIX_BATCH_SIZE); - if (prefixBatch.urls.length > 0) { - const prefixStats = await purgeUrls(prefixBatch.urls, apiKey, queue, 'prefix'); - mergePurgeStats(totalStats, prefixStats); - } - const remainingQueueSize = await queue.getQueueSize(); - Logger.debug( - { - totalAttempted: totalStats.attempted, - totalPurged: totalStats.purged, - totalRequeued: totalStats.requeued, - totalFailures: totalStats.failures, - totalRateLimited: totalStats.rateLimited, - remainingQueueSize, - }, - 'Finished processing CDN purge queue', - ); - } catch (error) { - Logger.error({error}, 'Error processing CDN purge queue'); - throw error; - } -}; - -export default processBunnyPurgeQueue; diff --git a/fluxer_api/src/api/worker/tasks/ProcessCachePurgeQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessCachePurgeQueue.ts new file mode 100644 index 000000000..94880b951 --- /dev/null +++ b/fluxer_api/src/api/worker/tasks/ProcessCachePurgeQueue.ts @@ -0,0 +1,104 @@ +// 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 {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; + +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 { + 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'); + return; + } + Logger.error(context, 'Cache purge endpoint refused the request, requeued its entries'); +} + +function untriedFrom(singles: ReadonlyArray, index: number): CachePurgeBatch { + const untried = singles.slice(index); + return {exact: untried.flatMap((single) => single.exact), prefix: untried.flatMap((single) => single.prefix)}; +} + +const processCachePurgeQueue: WorkerTaskHandler = async (_payload, _helpers) => { + const adapterName = Config.cachePurge.adapter; + if (adapterName === '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) { + return; + } + let outcome: CachePurgeOutcome; + try { + outcome = await adapter.purge(batch); + } catch (error) { + await queue.requeue(batch); + throw error; + } + if (outcome.kind === 'purged') { + Logger.debug(describeBatch(adapterName, batch), 'Purged a batch of cache entries'); + return; + } + if (outcome.kind === 'failed') { + await queue.requeue(batch); + logFailure(adapterName, batch, outcome); + return; + } + Logger.warn( + {...describeBatch(adapterName, batch), status: outcome.status}, + 'Cache purge endpoint rejected the batch, retrying each entry on its own', + ); + const singles: Array = [ + ...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; + } + } +}; + +export default processCachePurgeQueue; diff --git a/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts b/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts index ecb27459b..5dd2d99af 100644 --- a/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts +++ b/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts @@ -1,7 +1,7 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import {AssetDeletionQueue} from '@app/api/infrastructure/AssetDeletionQueue'; -import {NoopPurgeQueue} from '@app/api/infrastructure/BunnyPurgeQueue'; +import {NoopPurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; import type {IStorageService} from '@app/api/infrastructure/IStorageService'; import {MockKVProvider} from '@app/api/test/mocks/MockKVProvider'; import {NoopLogger} from '@app/api/test/mocks/NoopLogger'; diff --git a/fluxer_api/src/api/worker/tests/ProcessCachePurgeQueue.test.ts b/fluxer_api/src/api/worker/tests/ProcessCachePurgeQueue.test.ts new file mode 100644 index 000000000..b628de50c --- /dev/null +++ b/fluxer_api/src/api/worker/tests/ProcessCachePurgeQueue.test.ts @@ -0,0 +1,315 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {Config} from '@app/api/Config'; +import type {APICachePurgeConfig} from '@app/api/config/APIConfig'; +import {CachePurgeQueue} from '@app/api/infrastructure/CachePurgeQueue'; +import {MockKVProvider} from '@app/api/test/mocks/MockKVProvider'; +import {NoopLogger} from '@app/api/test/mocks/NoopLogger'; +import {server} from '@app/api/test/msw/server'; +import processCachePurgeQueue from '@app/api/worker/tasks/ProcessCachePurgeQueue'; +import {clearWorkerDependencies, setWorkerDependenciesForTest} from '@app/api/worker/WorkerContext'; +import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask'; +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 MEDIA = 'https://media.test'; +const EXACT_KEY = 'cache_purge:exact'; +const PREFIX_KEY = 'cache_purge:prefix'; +const REJECTED_KEY = 'cache_purge:rejected'; + +interface PurgeBody { + exact: Array; + prefix: Array; +} + +interface RecordedRequest { + authorization: string | null; + contentType: string | null; + body: PurgeBody; +} + +function createHarness() { + const kvClient = new MockKVProvider(); + setWorkerDependenciesForTest({kvClient}); + return {kvClient, queue: new CachePurgeQueue(kvClient)}; +} + +function recordPurges(respond: (body: PurgeBody) => Response | Promise): Array { + const requests: Array = []; + server.use( + http.post(ENDPOINT, async ({request}) => { + const body = (await request.json()) as PurgeBody; + requests.push({ + authorization: request.headers.get('authorization'), + contentType: request.headers.get('content-type'), + body, + }); + return respond(body); + }), + ); + return requests; +} + +function entryCount(body: PurgeBody): number { + return body.exact.length + body.prefix.length; +} + +async function members(kvClient: MockKVProvider, key: string): Promise> { + return (await kvClient.smembers(key)).sort(); +} + +function sorted(values: Array): Array { + return [...values].sort(); +} + +describe('processCachePurgeQueue', () => { + let previousCachePurge: APICachePurgeConfig; + + beforeEach(() => { + previousCachePurge = {adapter: Config.cachePurge.adapter, http: Config.cachePurge.http}; + Config.cachePurge.adapter = 'http'; + Config.cachePurge.http = {endpoint: ENDPOINT, token: 'test-token', timeoutMs: 50}; + }); + + afterEach(() => { + Config.cachePurge.adapter = previousCachePurge.adapter; + Config.cachePurge.http = previousCachePurge.http; + clearWorkerDependencies(); + vi.useRealTimers(); + }); + + it('sends queued exact and prefix entries 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]); + const requests = recordPurges(() => new HttpResponse(null, {status: 204})); + + await processCachePurgeQueue({}, HELPERS); + + 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([]); + }); + + 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 = ''; + const {queue} = createHarness(); + await queue.addUrls([`${MEDIA}/avatars/1/abc`]); + const requests = recordPurges(() => new HttpResponse(null, {status: 204})); + + await processCachePurgeQueue({}, HELPERS); + + expect(requests).toHaveLength(1); + expect(requests[0]!.authorization).toBeNull(); + }); + + 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]); + 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); + }); + + 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); + const requests = recordPurges(async () => { + await delay('infinite'); + return new HttpResponse(null, {status: 204}); + }); + + await processCachePurgeQueue({}, HELPERS); + + expect(requests).toHaveLength(1); + expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact)); + }); + + 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); + const requests = recordPurges(() => HttpResponse.error()); + + await processCachePurgeQueue({}, HELPERS); + + expect(requests).toHaveLength(1); + expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact)); + }); + + 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); + const redirectTarget = 'https://cache-purge.test/elsewhere'; + const redirectedRequests: Array = []; + server.use( + http.all(redirectTarget, ({request}) => { + redirectedRequests.push(request.method); + return new HttpResponse(null, {status: 204}); + }), + ); + const requests = recordPurges(() => new HttpResponse(null, {status: 302, headers: {Location: redirectTarget}})); + + await processCachePurgeQueue({}, HELPERS); + + expect(requests).toHaveLength(1); + expect(redirectedRequests).toEqual([]); + expect(await members(kvClient, EXACT_KEY)).toEqual(sorted(exact)); + }); + + it('requeues the batch on an authorisation failure without splitting it', async () => { + const {kvClient, queue} = createHarness(); + const exact = [`${MEDIA}/avatars/1/abc`, `${MEDIA}/banners/1/def`, `${MEDIA}/emojis/9.webp`]; + await queue.addUrls(exact); + 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, REJECTED_KEY)).toEqual([]); + }); + + it('sets aside only the entry the endpoint rejects on its own', 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: 422}) : 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('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([]); + }); + + 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]); + 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(); + }); + + 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')); + const {kvClient, queue} = createHarness(); + await queue.addUrls([ + ...Array.from({length: 200}, (_, index) => `${MEDIA}/avatars/${index}/abc`), + ...Array.from({length: 30}, (_, index) => `${MEDIA}/dirs/${index}/`), + ]); + const requests = recordPurges(() => new HttpResponse(null, {status: 204})); + + await processCachePurgeQueue({}, HELPERS); + await processCachePurgeQueue({}, HELPERS); + 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); + }); +}); diff --git a/fluxer_docs/src/content/docs/operator/configuration.mdx b/fluxer_docs/src/content/docs/operator/configuration.mdx index 37af0d0c0..595898438 100644 --- a/fluxer_docs/src/content/docs/operator/configuration.mdx +++ b/fluxer_docs/src/content/docs/operator/configuration.mdx @@ -1187,17 +1187,21 @@ Default empty. The Klipy GIF provider key. GIF search reports itself unavailable Default empty. The YouTube Data API key. Also read by `unfurl`, which accepts `YOUTUBE_API_KEY` as a fallback. -#### `FLUXER_BUNNY_PURGE_ENABLED` +#### `FLUXER_CACHE_PURGE_ADAPTER` -Default `false`. Bunny CDN cache purging. Needs `FLUXER_BUNNY_API_KEY` and `FLUXER_BUNNY_PULL_ZONE_ID`. +Default `none`. How cached copies of deleted or replaced media are purged. `none` or `http`. Anything else fails startup. With `none`, nothing is queued or sent. -#### `FLUXER_BUNNY_API_KEY` +#### `FLUXER_CACHE_PURGE_HTTP_ENDPOINT` -Default empty. The Bunny API key. Paired with the pull zone. +Default empty. Required when the adapter is `http`. Must be an absolute `http` or `https` URL without credentials, or startup fails. While purges are queued, the worker posts JSON `{"exact": [...], "prefix": [...]}` here every 10 seconds. Redirects are not followed. Any 2xx marks the batch done. The endpoint purges every proxy that caches media. A `400` or `422` makes the worker resend each URL alone, and a URL refused alone moves to the `cache_purge:rejected` key. Any other failure is retried on a later run. -#### `FLUXER_BUNNY_PULL_ZONE_ID` +#### `FLUXER_CACHE_PURGE_HTTP_TOKEN` -Default `0`. The Bunny pull zone. Integer. +Default empty. The bearer token sent in the `Authorization` header of each purge request. Empty sends no header. When the adapter is `http`, a token with spaces, line breaks or non-ASCII characters fails startup. + +#### `FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS` + +Default `10000`. Purge request timeout. Accepts 1000 to 10000. The range is checked only when the adapter is `http`. #### `FLUXER_GEOIP_DB_PATH` @@ -1397,18 +1401,6 @@ Default empty. An external NSFW classifier. No such service ships with the stack Default `0.85`. The media-side NSFW cutoff. Accepts 0.0 to 1.0 and finite. Distinct from the API threshold. -#### `FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_ENABLED` - -Default `false`. Restricts reads to Bunny edge IPs. - -#### `FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES` - -Default empty. Extra trusted addresses for that gate. Every entry must parse as an IP or the process exits. - -#### `FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_REFRESH_SECS` - -Default `3600`. How often the edge IP list refreshes. Accepts 60 to 86400. - `media-proxy` range-checks these values at startup and then reads them nowhere: `FLUXER_MEDIA_PROXY_UPLOAD_RELAY_BUFFERED_RETRY_BYTES` and `FLUXER_MEDIA_PROXY_UPLOAD_RELAY_BUFFERED_RETRY_TOTAL_BYTES`. A bad value still fails the boot. ## Gateway settings @@ -2004,7 +1996,7 @@ CORS origins are exactly the app and marketing endpoints. Serving the client fro | edge-config | The edge's own state | No | | meilisearch-data | The search index, rebuildable | No | -Losing `valkey-data` can delay scheduled account and bulk-message deletions while their queues are rebuilt. Pending asset deletions and CDN purges can be lost, so do not treat this volume as disposable cache. +Losing `valkey-data` can delay scheduled account and bulk-message deletions while their queues are rebuilt. Pending asset deletions and cache purges can be lost, so do not treat this volume as disposable cache. `nats-data` retains pending jobs in `JOBS` for up to 7 days and failed jobs in `JOBS_DLQ` for up to 30 days. A full jobs stream rejects new work. A full dead-letter stream drops its oldest entries, so investigate failures promptly. If dead-letter storage is unavailable, failed jobs remain in `JOBS` only until they expire. Losing this volume loses queued work, which is not automatically recovered from the database. diff --git a/fluxer_media_proxy/src/bunny_ip_gate.rs b/fluxer_media_proxy/src/bunny_ip_gate.rs deleted file mode 100644 index 5897defd1..000000000 --- a/fluxer_media_proxy/src/bunny_ip_gate.rs +++ /dev/null @@ -1,469 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -use axum::{ - body::Body, - extract::{ConnectInfo, State}, - http::{Request, StatusCode}, - middleware::Next, - response::{IntoResponse, Response}, -}; -use parking_lot::RwLock; -use std::{ - collections::HashSet, - net::{IpAddr, SocketAddr}, - sync::Arc, - time::Duration, -}; -use tracing::{info, warn}; - -pub const BUNNY_IPV4_URL: &str = "https://api.bunny.net/system/edgeserverlist"; -pub const BUNNY_IPV6_URL: &str = "https://api.bunny.net/system/edgeserverlist/IPv6"; - -#[derive(Clone, Debug, Default)] -pub struct Allowlist { - ips: Arc>, -} - -impl Allowlist { - fn is_empty(&self) -> bool { - self.ips.is_empty() - } - - fn len(&self) -> usize { - self.ips.len() - } - - pub fn contains(&self, ip: &IpAddr) -> bool { - self.ips.contains(ip) - } - - #[cfg(test)] - pub fn from_ips>(iter: I) -> Self { - Self { - ips: Arc::new(iter.into_iter().collect()), - } - } -} - -pub struct BunnyIpGate { - inner: RwLock, - trusted_proxies: HashSet, - client: reqwest::Client, - ipv4_url: String, - ipv6_url: String, -} - -impl BunnyIpGate { - pub fn new(client: reqwest::Client, trusted_proxies: Vec) -> Self { - Self::with_urls( - client, - trusted_proxies, - BUNNY_IPV4_URL.to_owned(), - BUNNY_IPV6_URL.to_owned(), - ) - } - - pub fn with_urls( - client: reqwest::Client, - trusted_proxies: Vec, - ipv4_url: String, - ipv6_url: String, - ) -> Self { - Self { - inner: RwLock::new(Allowlist::default()), - trusted_proxies: trusted_proxies.into_iter().collect(), - client, - ipv4_url, - ipv6_url, - } - } - - pub fn snapshot(&self) -> Allowlist { - self.inner.read().clone() - } - - pub async fn refresh_once(&self) -> anyhow::Result { - let allow = fetch_allowlist(&self.client, &self.ipv4_url, &self.ipv6_url).await?; - let len = allow.len(); - *self.inner.write() = allow; - Ok(len) - } - - pub fn spawn_background_refresher(self: Arc, interval: Duration) { - tokio::spawn(async move { - let mut ticker = tokio::time::interval(interval); - ticker.tick().await; - loop { - ticker.tick().await; - match self.refresh_once().await { - Ok(count) => info!(count, "bunny ip allowlist refreshed"), - Err(err) => warn!( - error = %err, - "bunny ip allowlist refresh failed; serving previous list" - ), - } - } - }); - } - - fn is_trusted_proxy(&self, ip: &IpAddr) -> bool { - self.trusted_proxies.contains(ip) - } - - #[cfg(test)] - pub fn install_for_test(&self, allow: Allowlist) { - *self.inner.write() = allow; - } -} - -async fn fetch_allowlist( - client: &reqwest::Client, - ipv4_url: &str, - ipv6_url: &str, -) -> anyhow::Result { - let (v4, v6) = tokio::join!( - fetch_ip_list(client, ipv4_url), - fetch_ip_list(client, ipv6_url) - ); - let v4 = v4?; - let v6 = v6?; - let mut set: HashSet = HashSet::with_capacity(v4.len() + v6.len()); - set.extend(v4); - set.extend(v6); - anyhow::ensure!( - !set.is_empty(), - "bunny returned an empty edge ip list from {ipv4_url} + {ipv6_url}" - ); - Ok(Allowlist { ips: Arc::new(set) }) -} - -async fn fetch_ip_list(client: &reqwest::Client, url: &str) -> anyhow::Result> { - let raw: Vec = client - .get(url) - .send() - .await? - .error_for_status()? - .json() - .await?; - Ok(parse_ip_list(&raw, url)) -} - -fn parse_ip_list(raw: &[String], url: &str) -> Vec { - let mut out = Vec::with_capacity(raw.len()); - for entry in raw { - match entry.trim().parse::() { - Ok(ip) => out.push(ip), - Err(_) => warn!(value = %entry, url, "skipping unparseable bunny ip"), - } - } - out -} - -pub async fn gate_middleware( - State(gate): State>, - ConnectInfo(peer): ConnectInfo, - request: Request, - next: Next, -) -> Response { - if is_exempt_path(request.uri().path()) { - return next.run(request).await; - } - let allow = gate.snapshot(); - if allow.is_empty() { - warn!("bunny ip allowlist is empty; refusing public request"); - return ( - StatusCode::SERVICE_UNAVAILABLE, - "bunny allowlist not loaded", - ) - .into_response(); - } - let client_ip = resolve_client_ip(&gate, &peer, request.headers()); - if !allow.contains(&client_ip) { - warn!(%client_ip, path = request.uri().path(), "rejecting non-bunny origin"); - return (StatusCode::FORBIDDEN, "origin not in bunny allowlist").into_response(); - } - next.run(request).await -} - -fn is_exempt_path(path: &str) -> bool { - matches!( - path, - "/_health" | "/_metrics" | "/_metadata" | "/_thumbnail" | "/_frames" - ) || path.starts_with("/v1/relay/") -} - -fn resolve_client_ip( - gate: &BunnyIpGate, - peer: &SocketAddr, - headers: &axum::http::HeaderMap, -) -> IpAddr { - let peer_ip = peer.ip(); - if !gate.is_trusted_proxy(&peer_ip) { - return peer_ip; - } - if let Some(xff) = headers.get("x-forwarded-for").and_then(|v| v.to_str().ok()) { - for hop in xff.split(',').rev() { - if let Ok(ip) = hop.trim().parse::() - && !gate.is_trusted_proxy(&ip) - { - return ip; - } - } - } - peer_ip -} - -pub fn build_refresh_client() -> reqwest::Result { - reqwest::Client::builder() - .connect_timeout(Duration::from_secs(5)) - .timeout(Duration::from_secs(20)) - .user_agent(crate::constants::OUTBOUND_USER_AGENT) - .build() -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::test_fixtures::ADVERSARIAL_TEXT_INPUTS; - use axum::{Router, middleware, routing::get}; - use std::net::{Ipv4Addr, Ipv6Addr}; - use tower::ServiceExt; - - fn ip4(a: u8, b: u8, c: u8, d: u8) -> IpAddr { - IpAddr::V4(Ipv4Addr::new(a, b, c, d)) - } - - fn sock(ip: IpAddr) -> SocketAddr { - SocketAddr::new(ip, 0) - } - - #[test] - fn parse_ip_list_skips_garbage_keeps_good() { - let raw = vec![ - "1.2.3.4".to_owned(), - "not-an-ip".to_owned(), - " 5.6.7.8 ".to_owned(), - "2a01:4f8::1".to_owned(), - "".to_owned(), - ]; - let parsed = parse_ip_list(&raw, "test"); - assert_eq!(parsed.len(), 3); - assert!(parsed.contains(&ip4(1, 2, 3, 4))); - assert!(parsed.contains(&ip4(5, 6, 7, 8))); - assert!(parsed.contains(&IpAddr::V6("2a01:4f8::1".parse::().unwrap()))); - } - - #[test] - fn untrusted_peer_xff_is_ignored() { - let gate = BunnyIpGate::new(reqwest::Client::new(), vec![ip4(10, 0, 0, 1)]); - let mut headers = axum::http::HeaderMap::new(); - headers.insert("x-forwarded-for", "9.9.9.9".parse().unwrap()); - let resolved = resolve_client_ip(&gate, &sock(ip4(8, 8, 8, 8)), &headers); - assert_eq!(resolved, ip4(8, 8, 8, 8)); - } - - #[test] - fn trusted_peer_uses_rightmost_untrusted_xff() { - let gate = BunnyIpGate::new( - reqwest::Client::new(), - vec![ip4(10, 0, 0, 1), ip4(10, 0, 0, 2)], - ); - let mut headers = axum::http::HeaderMap::new(); - headers.insert( - "x-forwarded-for", - "89.187.188.227, 10.0.0.2, 10.0.0.1".parse().unwrap(), - ); - let resolved = resolve_client_ip(&gate, &sock(ip4(10, 0, 0, 1)), &headers); - assert_eq!(resolved, ip4(89, 187, 188, 227)); - } - - #[test] - fn trusted_peer_falls_back_to_peer_when_no_xff() { - let gate = BunnyIpGate::new(reqwest::Client::new(), vec![ip4(10, 0, 0, 1)]); - let headers = axum::http::HeaderMap::new(); - let resolved = resolve_client_ip(&gate, &sock(ip4(10, 0, 0, 1)), &headers); - assert_eq!(resolved, ip4(10, 0, 0, 1)); - } - - #[test] - fn exempt_paths_cover_internal_routes() { - assert!(is_exempt_path("/_health")); - assert!(is_exempt_path("/_metrics")); - assert!(is_exempt_path("/_metadata")); - assert!(is_exempt_path("/_thumbnail")); - assert!(is_exempt_path("/_frames")); - assert!(is_exempt_path("/v1/relay/abc/def")); - assert!(!is_exempt_path("/external/some/url")); - assert!(!is_exempt_path("/some/asset.png")); - assert!(!is_exempt_path("/")); - } - - #[test] - fn snapshot_lookup_matches_loaded_ips() { - let gate = BunnyIpGate::new(reqwest::Client::new(), vec![]); - gate.install_for_test(Allowlist::from_ips([ip4(89, 187, 188, 227)])); - let snap = gate.snapshot(); - assert!(snap.contains(&ip4(89, 187, 188, 227))); - assert!(!snap.contains(&ip4(1, 1, 1, 1))); - assert_eq!(snap.len(), 1); - } - - fn ip6(text: &str) -> IpAddr { - IpAddr::V6(text.parse::().unwrap()) - } - - fn forwarded_for(gate: &BunnyIpGate, peer: IpAddr, xff: &str) -> IpAddr { - let mut headers = axum::http::HeaderMap::new(); - headers.insert("x-forwarded-for", xff.parse().unwrap()); - resolve_client_ip(gate, &sock(peer), &headers) - } - - #[test] - fn a_long_forwarded_chain_still_resolves_the_rightmost_untrusted_hop() { - let trusted: Vec = (1..=24).map(|last| ip4(10, 0, 0, last)).collect(); - let gate = BunnyIpGate::new(reqwest::Client::new(), trusted.clone()); - let mut chain = vec!["89.187.188.227".to_owned()]; - chain.extend(trusted.iter().map(ToString::to_string)); - assert_eq!(25, chain.len()); - assert_eq!( - ip4(89, 187, 188, 227), - forwarded_for(&gate, ip4(10, 0, 0, 24), &chain.join(", ")) - ); - assert_eq!( - ip4(10, 0, 0, 24), - forwarded_for(&gate, ip4(10, 0, 0, 24), &chain[1..].join(", ")) - ); - } - - #[test] - fn unparseable_hops_are_skipped_instead_of_ending_the_scan() { - let gate = BunnyIpGate::new( - reqwest::Client::new(), - vec![ip4(10, 0, 0, 1), ip4(10, 0, 0, 2)], - ); - assert_eq!( - ip4(89, 187, 188, 227), - forwarded_for( - &gate, - ip4(10, 0, 0, 1), - "89.187.188.227, unknown, 10.0.0.2, 10.0.0.1" - ) - ); - assert_eq!( - ip4(89, 187, 188, 227), - forwarded_for( - &gate, - ip4(10, 0, 0, 1), - "89.187.188.227, 10.0.0.1, , not-an-ip" - ) - ); - assert_eq!( - ip4(10, 0, 0, 1), - forwarded_for(&gate, ip4(10, 0, 0, 1), "unknown, also-unknown") - ); - } - - #[test] - fn forwarded_ipv6_hops_are_only_accepted_in_their_bare_form() { - let gate = BunnyIpGate::new(reqwest::Client::new(), vec![ip4(10, 0, 0, 1)]); - assert_eq!( - ip6("2a01:4f8::1"), - forwarded_for(&gate, ip4(10, 0, 0, 1), "2a01:4f8::1, 10.0.0.1") - ); - assert_eq!( - ip4(10, 0, 0, 1), - forwarded_for(&gate, ip4(10, 0, 0, 1), "[2a01:4f8::1]:443, 10.0.0.1") - ); - assert_eq!( - ip6("2a01:4f8::2"), - forwarded_for( - &gate, - ip4(10, 0, 0, 1), - "2a01:4f8::2, [2a01:4f8::1]:443, 10.0.0.1" - ) - ); - } - - #[test] - fn a_trusted_ipv6_proxy_is_skipped_like_any_other_hop() { - let gate = BunnyIpGate::new(reqwest::Client::new(), vec![ip6("2a01:4f8::1")]); - assert_eq!( - ip4(89, 187, 188, 227), - forwarded_for(&gate, ip6("2a01:4f8::1"), "89.187.188.227, 2a01:4f8::1") - ); - } - - #[test] - fn adversarial_paths_are_never_exempt_from_the_bunny_gate() { - for text in ADVERSARIAL_TEXT_INPUTS { - for path in [ - (*text).to_owned(), - format!("/{text}"), - format!("/avatars/1/{text}.png"), - format!("/external/{text}"), - format!("/_health{text}"), - format!("/v1/relay{text}"), - format!("/v1/relay/{text}"), - ] { - if !is_exempt_path(&path) { - continue; - } - assert!( - path == "/_health" || path.starts_with("/v1/relay/"), - "{path:?} is exempt from the bunny gate" - ); - } - } - } - - #[tokio::test] - async fn gate_after_successful_refresh_never_returns_503() { - let gate = Arc::new(BunnyIpGate::new(reqwest::Client::new(), vec![])); - gate.install_for_test(Allowlist::from_ips([ip4(89, 187, 188, 227)])); - let app = Router::new() - .route("/asset.png", get(|| async { "ok" })) - .layer(middleware::from_fn_with_state(gate, gate_middleware)); - let allowed = app - .clone() - .oneshot( - Request::builder() - .uri("/asset.png") - .extension(ConnectInfo(sock(ip4(89, 187, 188, 227)))) - .body(Body::empty()) - .unwrap(), - ) - .await - .unwrap(); - assert_eq!(StatusCode::OK, allowed.status()); - let foreign = app - .oneshot( - Request::builder() - .uri("/asset.png") - .extension(ConnectInfo(sock(ip4(1, 1, 1, 1)))) - .body(Body::empty()) - .unwrap(), - ) - .await - .unwrap(); - assert_eq!(StatusCode::FORBIDDEN, foreign.status()); - } - - #[tokio::test] - async fn an_empty_allowlist_refuses_a_public_request() { - let gate = Arc::new(BunnyIpGate::new(reqwest::Client::new(), vec![])); - let app = Router::new() - .route("/asset.png", get(|| async { "ok" })) - .layer(middleware::from_fn_with_state(gate, gate_middleware)); - let refused = app - .oneshot( - Request::builder() - .uri("/asset.png") - .extension(ConnectInfo(sock(ip4(89, 187, 188, 227)))) - .body(Body::empty()) - .unwrap(), - ) - .await - .unwrap(); - assert_eq!(StatusCode::SERVICE_UNAVAILABLE, refused.status()); - } -} diff --git a/fluxer_media_proxy/src/config/mod.rs b/fluxer_media_proxy/src/config/mod.rs index ea73c7809..301adc2dc 100644 --- a/fluxer_media_proxy/src/config/mod.rs +++ b/fluxer_media_proxy/src/config/mod.rs @@ -8,10 +8,10 @@ use crate::constants; use crate::secret::{SecretBytes, SecretString}; use parse::{ EnvMap, decode_upload_relay_secret, default_native_transform_concurrency, non_empty, - parse_bool, parse_bucket_style, parse_f32, parse_ip_list_env, parse_mode_env, - parse_storage_backend, parse_u16, parse_u64, parse_usize, validate_read_endpoint, + parse_bool, parse_bucket_style, parse_f32, parse_mode_env, parse_storage_backend, parse_u16, + parse_u64, parse_usize, validate_read_endpoint, }; -use std::{env, net::IpAddr, path::PathBuf}; +use std::{env, path::PathBuf}; #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum StorageBackend { @@ -92,9 +92,6 @@ pub struct Config { pub storage: StorageConfig, pub media: MediaServingConfig, pub upload_relay: UploadRelayConfig, - pub bunny_ip_gate_enabled: bool, - pub bunny_ip_gate_trusted_proxies: Vec, - pub bunny_ip_gate_refresh_secs: u64, } impl Config { @@ -170,22 +167,6 @@ impl Config { storage: StorageConfig::load(&env)?, media: MediaServingConfig::load(&env)?, upload_relay: UploadRelayConfig::load(&env, mode)?, - bunny_ip_gate_enabled: parse_bool( - "FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_ENABLED", - env.get("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_ENABLED"), - )? - .unwrap_or(false), - bunny_ip_gate_trusted_proxies: parse_ip_list_env( - "FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES", - env.get("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES"), - )?, - bunny_ip_gate_refresh_secs: parse_u64( - "FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_REFRESH_SECS", - env.get("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_REFRESH_SECS"), - 3_600, - 60, - 24 * 60 * 60, - )?, }) } } diff --git a/fluxer_media_proxy/src/config/parse.rs b/fluxer_media_proxy/src/config/parse.rs index 683120cfb..59c375860 100644 --- a/fluxer_media_proxy/src/config/parse.rs +++ b/fluxer_media_proxy/src/config/parse.rs @@ -3,7 +3,6 @@ use super::{BucketStyle, DeploymentMode, StorageBackend}; use crate::secret::SecretBytes; use base64::{Engine as _, engine::general_purpose}; -use std::net::IpAddr; #[derive(Debug, Default)] pub(super) struct EnvMap(Vec<(String, String)>); @@ -213,24 +212,6 @@ where Ok(parsed) } -pub(super) fn parse_ip_list_env(var_name: &str, raw: Option<&str>) -> anyhow::Result> { - let Some(raw) = raw.map(str::trim).filter(|s| !s.is_empty()) else { - return Ok(Vec::new()); - }; - let mut out = Vec::new(); - for entry in raw.split(',') { - let entry = entry.trim(); - if entry.is_empty() { - continue; - } - let ip = entry - .parse::() - .map_err(|_| anyhow::anyhow!("{var_name} contains invalid IP: {entry}"))?; - out.push(ip); - } - Ok(out) -} - pub(super) fn default_native_transform_concurrency() -> usize { std::thread::available_parallelism() .map(usize::from) diff --git a/fluxer_media_proxy/src/config/tests/mod.rs b/fluxer_media_proxy/src/config/tests/mod.rs index 4daf28e03..7f368c168 100644 --- a/fluxer_media_proxy/src/config/tests/mod.rs +++ b/fluxer_media_proxy/src/config/tests/mod.rs @@ -410,7 +410,7 @@ fn every_deployment_mode_and_storage_backend_variant_parses() { } #[test] -fn upload_relay_spool_and_bunny_ip_gate_keys_apply() { +fn upload_relay_spool_keys_apply() { let cfg = Config::load_from_iter([ ("FLUXER_MEDIA_PROXY_SECRET_KEY", "secret"), ( @@ -425,12 +425,6 @@ fn upload_relay_spool_and_bunny_ip_gate_keys_apply() { "FLUXER_MEDIA_PROXY_UPLOAD_RELAY_SPOOL_MAX_TOTAL_BYTES", "1073741824", ), - ("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_ENABLED", "yes"), - ( - "FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES", - "10.0.0.1, 2001:db8::1 ,", - ), - ("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_REFRESH_SECS", "900"), ]) .unwrap(); @@ -440,31 +434,6 @@ fn upload_relay_spool_and_bunny_ip_gate_keys_apply() { ); assert_eq!(2 << 20, cfg.upload_relay.spool_chunk_bytes); assert_eq!(1 << 30, cfg.upload_relay.spool_max_total_bytes); - assert!(cfg.bunny_ip_gate_enabled); - assert_eq!( - vec![ - "10.0.0.1".parse::().unwrap(), - "2001:db8::1".parse::().unwrap(), - ], - cfg.bunny_ip_gate_trusted_proxies - ); - assert_eq!(900, cfg.bunny_ip_gate_refresh_secs); -} - -#[test] -fn rejects_invalid_bunny_ip_gate_trusted_proxies() { - let err = Config::load_from_iter([ - ("FLUXER_MEDIA_PROXY_SECRET_KEY", "secret"), - ( - "FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES", - "10.0.0.1,not-an-ip", - ), - ]) - .unwrap_err(); - assert!( - err.to_string() - .contains("FLUXER_MEDIA_PROXY_BUNNY_IP_GATE_TRUSTED_PROXIES") - ); } #[test] diff --git a/fluxer_media_proxy/src/lib.rs b/fluxer_media_proxy/src/lib.rs index c6fde2827..fcaaa435e 100644 --- a/fluxer_media_proxy/src/lib.rs +++ b/fluxer_media_proxy/src/lib.rs @@ -4,7 +4,6 @@ mod aggregate_error; pub mod asset_hash; mod asset_size; pub mod aws_sigv4; -pub mod bunny_ip_gate; mod byte_budget; mod byte_cache; pub mod cli; diff --git a/fluxer_media_proxy/src/server/runtime.rs b/fluxer_media_proxy/src/server/runtime.rs index 24147bd74..3c35fa629 100644 --- a/fluxer_media_proxy/src/server/runtime.rs +++ b/fluxer_media_proxy/src/server/runtime.rs @@ -7,13 +7,7 @@ use super::{ routes, state::AppState, }; -use crate::{ - aggregate_error::aggregate_results, - bunny_ip_gate::{self, BunnyIpGate}, - config::Config, - media_process, request_log, -}; -use anyhow::Context as _; +use crate::{aggregate_error::aggregate_results, config::Config, media_process, request_log}; use axum::{ Router, middleware, routing::{any, get, post, put}, @@ -26,7 +20,6 @@ pub async fn run(cfg: Config) -> anyhow::Result<()> { media_process::warmup_vips()?; let addr: SocketAddr = format!("{}:{}", cfg.bind_host, cfg.port).parse()?; let state = Arc::new(AppState::try_new(cfg)?); - let bunny_gate = start_bunny_ip_gate(&state).await?; if let Some(read_endpoint) = state.cfg.storage.s3_read_endpoint.as_deref() { info!( endpoint = read_endpoint, @@ -37,7 +30,7 @@ pub async fn run(cfg: Config) -> anyhow::Result<()> { ); } let drain = HttpRequestDrain::new(); - let app = build_router(Arc::clone(&state), bunny_gate, drain.clone()); + let app = build_router(Arc::clone(&state), drain.clone()); let listener = TcpListener::bind(addr).await?; info!(%addr, "media proxy listening"); axum::serve( @@ -49,12 +42,8 @@ pub async fn run(cfg: Config) -> anyhow::Result<()> { drain_and_shutdown(&state, &drain).await } -fn build_router( - state: Arc, - bunny_gate: Option>, - drain: HttpRequestDrain, -) -> Router { - let mut router = Router::new() +fn build_router(state: Arc, drain: HttpRequestDrain) -> Router { + Router::new() .route("/_health", get(routes::ops::health)) .route("/_metrics", get(routes::ops::metrics_handler)) .route("/_metadata", post(routes::internal::metadata_handler)) @@ -69,14 +58,7 @@ fn build_router( .layer(middleware::from_fn_with_state( state.metrics.request(), request_log::trace, - )); - if let Some(gate) = bunny_gate { - router = router.layer(middleware::from_fn_with_state( - gate, - bunny_ip_gate::gate_middleware, - )); - } - router + )) .layer(middleware::from_fn_with_state( state.cfg.mode, add_security_header_middleware, @@ -85,29 +67,6 @@ fn build_router( .with_state(state) } -async fn start_bunny_ip_gate(state: &Arc) -> anyhow::Result>> { - if !state.cfg.bunny_ip_gate_enabled { - return Ok(None); - } - let gate = Arc::new(BunnyIpGate::new( - bunny_ip_gate::build_refresh_client()?, - state.cfg.bunny_ip_gate_trusted_proxies.clone(), - )); - let count = gate - .refresh_once() - .await - .context("initial bunny ip allowlist fetch failed")?; - info!( - count, - trusted_proxies = state.cfg.bunny_ip_gate_trusted_proxies.len(), - refresh_secs = state.cfg.bunny_ip_gate_refresh_secs, - "bunny ip gate enabled" - ); - Arc::clone(&gate) - .spawn_background_refresher(Duration::from_secs(state.cfg.bunny_ip_gate_refresh_secs)); - Ok(Some(gate)) -} - async fn drain_and_shutdown(state: &Arc, drain: &HttpRequestDrain) -> anyhow::Result<()> { let grace_ms = state.cfg.shutdown_grace_ms; let deadline = tokio::time::Instant::now() + Duration::from_millis(grace_ms); @@ -177,11 +136,7 @@ mod tests { fn test_router() -> Router { let cfg = Config::load_from_iter([("FLUXER_MEDIA_PROXY_SECRET_KEY", "secret")]) .expect("test config"); - build_router( - Arc::new(AppState::for_tests(cfg)), - None, - HttpRequestDrain::new(), - ) + build_router(Arc::new(AppState::for_tests(cfg)), HttpRequestDrain::new()) } async fn probe(path: &str, method: Method) -> axum::response::Response { diff --git a/fluxer_media_proxy/src/storage/tests/mod.rs b/fluxer_media_proxy/src/storage/tests/mod.rs index 7a05b5221..b9b1e9aa2 100644 --- a/fluxer_media_proxy/src/storage/tests/mod.rs +++ b/fluxer_media_proxy/src/storage/tests/mod.rs @@ -75,9 +75,6 @@ fn test_config(root: &Path) -> Config { spool_chunk_bytes: 64 * 1024, spool_max_total_bytes: 1 << 30, }, - bunny_ip_gate_enabled: false, - bunny_ip_gate_trusted_proxies: Vec::new(), - bunny_ip_gate_refresh_secs: 3_600, } } diff --git a/packages/config/src/ConfigLoader.ts b/packages/config/src/ConfigLoader.ts index 91e1bc41d..4f145bb89 100644 --- a/packages/config/src/ConfigLoader.ts +++ b/packages/config/src/ConfigLoader.ts @@ -11,7 +11,7 @@ import { parsePublicOrigin, parseWebOrigin, } from '@fluxer/config/src/EndpointDerivation'; -import type {MasterConfig} from '@fluxer/config/src/MasterConfig'; +import {CACHE_PURGE_ADAPTER_NAMES, type MasterConfig} from '@fluxer/config/src/MasterConfig'; let cachedConfig: MasterConfig | null = null; @@ -235,10 +235,13 @@ function defaultConfig(): MasterConfig { youtube: { api_key: '', }, - bunny: { - purge_enabled: false, - api_key: '', - pull_zone_id: 0, + cache_purge: { + adapter: 'none', + http: { + endpoint: '', + token: '', + timeout_ms: 10_000, + }, }, blocklist_feeds: {}, risk_integration: { @@ -449,6 +452,27 @@ function validateApiWorkerConfig(config: MasterConfig): void { } } +function validateCachePurgeConfig(config: MasterConfig): void { + const cachePurge = config.integrations.cache_purge; + if (cachePurge.adapter !== 'http') { + return; + } + requireString(cachePurge.http.endpoint, 'FLUXER_CACHE_PURGE_HTTP_ENDPOINT'); + const endpoint = URL.parse(cachePurge.http.endpoint); + if ( + endpoint === null || + (endpoint.protocol !== 'http:' && endpoint.protocol !== 'https:') || + endpoint.username !== '' || + endpoint.password !== '' + ) { + 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)) { + 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); +} + function validateDomain(value: string, envName: string): void { if (value === '') return; const parsed = URL.parse(`http://${value}/`); @@ -507,6 +531,7 @@ function normalizeConfig(config: MasterConfig): MasterConfig { assertOneOf(config.integrations.email.provider, ['smtp', 'none'], 'FLUXER_EMAIL_PROVIDER'); assertOneOf(config.integrations.captcha.provider, ['hcaptcha', 'turnstile', 'none'], 'FLUXER_CAPTCHA_PROVIDER'); assertOneOf(config.integrations.search.engine, ['elasticsearch', 'meilisearch'], 'FLUXER_SEARCH_ENGINE'); + assertOneOf(config.integrations.cache_purge.adapter, CACHE_PURGE_ADAPTER_NAMES, 'FLUXER_CACHE_PURGE_ADAPTER'); assertOneOf( config.instance.abuse_policy.direct_contact_spam.action, ['flag_spammer', 'suppress_delivery'], @@ -515,6 +540,7 @@ function normalizeConfig(config: MasterConfig): MasterConfig { validatePostgresConfig(config); validateCaptchaConfig(config); validateApiWorkerConfig(config); + validateCachePurgeConfig(config); assertIntegerInRange(config.services.api.max_inflight_requests, 'FLUXER_API_MAX_INFLIGHT_REQUESTS', 1, 100_000); assertIntegerInRange(config.services.api.headers_timeout_ms, 'FLUXER_API_HEADERS_TIMEOUT_MS', 1_000, 3_600_000); assertIntegerInRange(config.services.api.request_timeout_ms, 'FLUXER_API_REQUEST_TIMEOUT_MS', 1_000, 3_600_000); diff --git a/packages/config/src/MasterConfig.ts b/packages/config/src/MasterConfig.ts index 7a6d75731..7f70da0df 100644 --- a/packages/config/src/MasterConfig.ts +++ b/packages/config/src/MasterConfig.ts @@ -5,6 +5,8 @@ import type {DerivedEndpoints} from '@fluxer/config/src/EndpointDerivation'; export type RuntimeEnv = 'development' | 'production' | 'test'; export type DatabaseBackend = 'postgres' | 'cassandra'; export type PublicScheme = 'http' | 'https'; +export const CACHE_PURGE_ADAPTER_NAMES = ['none', 'http'] as const; +export type CachePurgeAdapterName = (typeof CACHE_PURGE_ADAPTER_NAMES)[number]; export interface InstanceBrandingConfig { product_name: string; @@ -274,10 +276,13 @@ export interface MasterConfig { youtube: { api_key: string; }; - bunny: { - purge_enabled: boolean; - api_key: string; - pull_zone_id: number; + cache_purge: { + adapter: CachePurgeAdapterName; + http: { + endpoint: string; + token: string; + timeout_ms: number; + }; }; blocklist_feeds: { enabled?: boolean; diff --git a/packages/config/src/__tests__/ConfigLoader.test.ts b/packages/config/src/__tests__/ConfigLoader.test.ts index 93632f919..cd17b4060 100644 --- a/packages/config/src/__tests__/ConfigLoader.test.ts +++ b/packages/config/src/__tests__/ConfigLoader.test.ts @@ -513,6 +513,90 @@ describe('ConfigLoader', () => { expect((await loadConfig()).integrations.captcha.enabled).toBe(false); }); + test('defaults the cache purge adapter to none', async () => { + stubMinimalEnv(); + expect((await loadConfig()).integrations.cache_purge).toEqual({ + adapter: 'none', + http: {endpoint: '', token: '', timeout_ms: 10_000}, + }); + }); + + test('reads the http cache purge settings from the environment', async () => { + 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: '5000', + }); + expect((await loadConfig()).integrations.cache_purge).toEqual({ + adapter: 'http', + http: {endpoint: 'https://purge.internal/purge', token: 'purge-token', timeout_ms: 5000}, + }); + }); + + test('rejects an unknown cache purge adapter', async () => { + stubMinimalEnv({FLUXER_CACHE_PURGE_ADAPTER: 'varnish'}); + await expect(loadConfig()).rejects.toThrow('Invalid FLUXER_CACHE_PURGE_ADAPTER: varnish'); + }); + + test('rejects the http cache purge adapter without an endpoint', async () => { + stubMinimalEnv({FLUXER_CACHE_PURGE_ADAPTER: 'http'}); + await expect(loadConfig()).rejects.toThrow('FLUXER_CACHE_PURGE_HTTP_ENDPOINT 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}); + await expect(loadConfig()).rejects.toThrow( + 'FLUXER_CACHE_PURGE_HTTP_ENDPOINT must be an absolute http or https URL without credentials', + ); + } + }); + + test('rejects a cache purge endpoint that carries credentials', async () => { + stubMinimalEnv({ + FLUXER_CACHE_PURGE_ADAPTER: 'http', + FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'https://purger:secret@purge.internal/purge', + }); + await expect(loadConfig()).rejects.toThrow( + 'FLUXER_CACHE_PURGE_HTTP_ENDPOINT must be an absolute http or https URL without credentials', + ); + }); + + test('rejects a cache purge timeout outside 1000 to 10000', async () => { + for (const timeout of ['999', '10001']) { + stubMinimalEnv({ + FLUXER_CACHE_PURGE_ADAPTER: 'http', + FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'https://purge.internal/purge', + FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS: timeout, + }); + await expect(loadConfig()).rejects.toThrow( + 'FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS must be an integer between 1000 and 10000', + ); + } + }); + + test('rejects a cache purge token with spaces or control characters without echoing it', async () => { + for (const token of ['secret\nvalue', 'secret\rvalue', 'secret value', 'secret-value\n']) { + stubMinimalEnv({ + FLUXER_CACHE_PURGE_ADAPTER: 'http', + FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'https://purge.internal/purge', + FLUXER_CACHE_PURGE_HTTP_TOKEN: token, + }); + await expect(loadConfig()).rejects.toThrow( + /^FLUXER_CACHE_PURGE_HTTP_TOKEN must contain only visible ASCII characters$/, + ); + } + }); + + test('leaves the cache purge settings unvalidated when the adapter is none', async () => { + stubMinimalEnv({ + FLUXER_CACHE_PURGE_HTTP_ENDPOINT: 'ftp://purge.internal/purge', + FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS: '1', + }); + expect((await loadConfig()).integrations.cache_purge.adapter).toBe('none'); + }); + test('leaves Bluesky login off with no legal URLs by default', async () => { stubMinimalEnv(); diff --git a/packages/config/src/config_loader/EnvironmentOverrides.ts b/packages/config/src/config_loader/EnvironmentOverrides.ts index 885906e00..320562017 100644 --- a/packages/config/src/config_loader/EnvironmentOverrides.ts +++ b/packages/config/src/config_loader/EnvironmentOverrides.ts @@ -318,10 +318,14 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record = { FLUXER_CLAMAV_FAIL_OPEN: {path: ['integrations', 'clamav', 'fail_open'], parse: parseBoolean}, FLUXER_KLIPY_API_KEY: {path: ['integrations', 'klipy', 'api_key']}, FLUXER_YOUTUBE_API_KEY: {path: ['integrations', 'youtube', 'api_key']}, - FLUXER_BUNNY_PURGE_ENABLED: {path: ['integrations', 'bunny', 'purge_enabled'], parse: parseBoolean}, + FLUXER_CACHE_PURGE_ADAPTER: {path: ['integrations', 'cache_purge', 'adapter']}, + FLUXER_CACHE_PURGE_HTTP_ENDPOINT: {path: ['integrations', 'cache_purge', 'http', 'endpoint']}, + FLUXER_CACHE_PURGE_HTTP_TOKEN: {path: ['integrations', 'cache_purge', 'http', 'token']}, + FLUXER_CACHE_PURGE_HTTP_TIMEOUT_MS: { + path: ['integrations', 'cache_purge', 'http', 'timeout_ms'], + parse: parseInteger, + }, FLUXER_BLOCKLIST_FEEDS_ENABLED: {path: ['integrations', 'blocklist_feeds', 'enabled'], parse: parseBoolean}, - FLUXER_BUNNY_API_KEY: {path: ['integrations', 'bunny', 'api_key']}, - FLUXER_BUNNY_PULL_ZONE_ID: {path: ['integrations', 'bunny', 'pull_zone_id'], parse: parseInteger}, FLUXER_RISK_INTEGRATION_ENABLED: {path: ['integrations', 'risk_integration', 'enabled'], parse: parseBoolean}, FLUXER_RISK_IPINFO_API_KEY: {path: ['integrations', 'risk_integration', 'ipinfo_api_key']}, FLUXER_ACCOUNT_POLICY_DSL: {