From 803fdaf443038c89917bbc48d53acd273a6ccb36 Mon Sep 17 00:00:00 2001 From: Hampus Date: Mon, 31 Aug 2026 22:19:22 +0200 Subject: [PATCH] fix(worker): stop replaying requeued asset deletions in a run (#2280) --- .../worker/tasks/ProcessAssetDeletionQueue.ts | 5 +- .../tests/ProcessAssetDeletionQueue.test.ts | 77 +++++++++++++++++++ 2 files changed, 80 insertions(+), 2 deletions(-) create mode 100644 fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts diff --git a/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts index 71b52e22d..e822ae0ac 100644 --- a/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts +++ b/fluxer_api/src/api/worker/tasks/ProcessAssetDeletionQueue.ts @@ -32,13 +32,14 @@ const processAssetDeletionQueue: WorkerTaskHandler = async (_payload, _helpers) return; } Logger.info({queueSize}, 'Starting asset deletion queue processing'); + const maxItems = Math.min(MAX_ITEMS_PER_RUN, queueSize); let totalProcessed = 0; let totalDeleted = 0; let totalSkipped = 0; let totalFailed = 0; let totalCdnPurged = 0; - while (totalProcessed < MAX_ITEMS_PER_RUN) { - const batch = await assetDeletionQueue.getBatch(BATCH_SIZE); + while (totalProcessed < maxItems) { + const batch = await assetDeletionQueue.getBatch(Math.min(BATCH_SIZE, maxItems - totalProcessed)); if (batch.length === 0) { break; } diff --git a/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts b/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts new file mode 100644 index 000000000..20f972017 --- /dev/null +++ b/fluxer_api/src/api/worker/tests/ProcessAssetDeletionQueue.test.ts @@ -0,0 +1,77 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask'; +import {afterEach, describe, expect, it, vi} from 'vitest'; +import {AssetDeletionQueue} from '../../infrastructure/AssetDeletionQueue'; +import {NoopPurgeQueue} from '../../infrastructure/BunnyPurgeQueue'; +import type {IStorageService} from '../../infrastructure/IStorageService'; +import {MockKVProvider} from '../../test/mocks/MockKVProvider'; +import {NoopLogger} from '../../test/mocks/NoopLogger'; +import processAssetDeletionQueue from '../tasks/ProcessAssetDeletionQueue'; +import {clearWorkerDependencies, setWorkerDependenciesForTest} from '../WorkerContext'; + +const HELPERS = {logger: new NoopLogger()} as unknown as WorkerTaskHelpers; + +function createHarness(deleteObject: IStorageService['deleteObject']) { + const kvClient = new MockKVProvider(); + const assetDeletionQueue = new AssetDeletionQueue(kvClient); + setWorkerDependenciesForTest({ + assetDeletionQueue, + purgeQueue: new NoopPurgeQueue(), + storageService: {deleteObject} as unknown as IStorageService, + userRepository: {findUnique: async () => null} as never, + guildRepository: {findUnique: async () => null, getMember: async () => null} as never, + }); + return {assetDeletionQueue, kvClient}; +} + +async function queueAssets(assetDeletionQueue: AssetDeletionQueue, count: number): Promise { + for (let index = 0; index < count; index++) { + await assetDeletionQueue.queueDeletion({ + s3Key: `attachments/100/200/file-${index}.png`, + cdnUrl: null, + reason: 'test', + }); + } +} + +describe('processAssetDeletionQueue', () => { + afterEach(() => { + clearWorkerDependencies(); + }); + + it('attempts each queued asset once per run when storage is failing', async () => { + const deleteObject = vi.fn().mockRejectedValue(new Error('s3 unavailable')); + const {assetDeletionQueue} = createHarness(deleteObject); + await queueAssets(assetDeletionQueue, 3); + + await expect(processAssetDeletionQueue({}, HELPERS)).rejects.toThrow(/3 failures/); + + expect(deleteObject).toHaveBeenCalledTimes(3); + expect(await assetDeletionQueue.getQueueSize()).toBe(3); + const remaining = await assetDeletionQueue.getBatch(10); + expect(remaining.map((item) => item.retryCount)).toEqual([1, 1, 1]); + }); + + it('attempts each queued asset once per run across batch boundaries', async () => { + const deleteObject = vi.fn().mockRejectedValue(new Error('s3 unavailable')); + const {assetDeletionQueue} = createHarness(deleteObject); + await queueAssets(assetDeletionQueue, 60); + + await expect(processAssetDeletionQueue({}, HELPERS)).rejects.toThrow(/60 failures/); + + expect(deleteObject).toHaveBeenCalledTimes(60); + expect(await assetDeletionQueue.getQueueSize()).toBe(60); + }); + + it('drains the queue when storage succeeds', async () => { + const deleteObject = vi.fn().mockResolvedValue(undefined); + const {assetDeletionQueue} = createHarness(deleteObject); + await queueAssets(assetDeletionQueue, 60); + + await processAssetDeletionQueue({}, HELPERS); + + expect(deleteObject).toHaveBeenCalledTimes(60); + expect(await assetDeletionQueue.getQueueSize()).toBe(0); + }); +});