From 49f76e5b4049ced4c0070ef927116ccffde5f25a Mon Sep 17 00:00:00 2001 From: Hampus Date: Tue, 1 Sep 2026 02:45:29 +0200 Subject: [PATCH] fix(api): abort startup on unverifiable deletion queue state (#2311) --- .../src/api/app/DeletionQueueStartup.ts | 37 ++++++--- .../app/tests/DeletionQueueStartup.test.ts | 81 ++++++++++++++++++- 2 files changed, 105 insertions(+), 13 deletions(-) diff --git a/fluxer_api/src/api/app/DeletionQueueStartup.ts b/fluxer_api/src/api/app/DeletionQueueStartup.ts index 8c7eeb3e1..9bf65fbf0 100644 --- a/fluxer_api/src/api/app/DeletionQueueStartup.ts +++ b/fluxer_api/src/api/app/DeletionQueueStartup.ts @@ -12,18 +12,31 @@ export async function ensureDeletionQueueState( logger.info('KV deletion queue state is healthy'); return; } - const lockToken = await deletionQueue.acquireRebuildLock(); - if (!lockToken) { - logger.info('Another instance is rebuilding the KV deletion queue, skipping'); - return; - } - logger.info('KV deletion queue needs rebuild, rebuilding...'); - try { - await deletionQueue.rebuildState(lockToken); - } finally { - await deletionQueue.releaseRebuildLock(lockToken); - } } catch (error) { - logger.error({error}, 'KV deletion queue rebuild failed, continuing startup'); + logger.error({error}, 'Failed to read KV deletion queue state, aborting startup'); + throw error; + } + let lockToken: string | null; + try { + lockToken = await deletionQueue.acquireRebuildLock(); + } catch (error) { + logger.error({error}, 'Failed to acquire the KV deletion queue rebuild lock, aborting startup'); + throw error; + } + if (!lockToken) { + logger.info('Another instance is rebuilding the KV deletion queue, skipping'); + return; + } + logger.info('KV deletion queue needs rebuild, rebuilding...'); + try { + await deletionQueue.rebuildState(lockToken); + } catch (error) { + logger.error({error}, 'KV deletion queue rebuild failed, leaving the rebuild to the deletion worker'); + } finally { + try { + await deletionQueue.releaseRebuildLock(lockToken); + } catch (error) { + logger.error({error}, 'Failed to release the KV deletion queue rebuild lock'); + } } } diff --git a/fluxer_api/src/api/app/tests/DeletionQueueStartup.test.ts b/fluxer_api/src/api/app/tests/DeletionQueueStartup.test.ts index 3d19bceb6..699eca0be 100644 --- a/fluxer_api/src/api/app/tests/DeletionQueueStartup.test.ts +++ b/fluxer_api/src/api/app/tests/DeletionQueueStartup.test.ts @@ -3,6 +3,7 @@ import {describe, expect, it} from 'vitest'; import {createUserID} from '../../BrandedTypes'; import {EMPTY_USER_ROW} from '../../database/types/UserTypes'; +import type {ILogger} from '../../ILogger'; import {KVAccountDeletionQueueService} from '../../infrastructure/KVAccountDeletionQueueService'; import {User} from '../../models/User'; import {MockKVProvider} from '../../test/mocks/MockKVProvider'; @@ -12,6 +13,50 @@ import {ensureDeletionQueueState} from '../DeletionQueueStartup'; const WORKER_PAGE_COUNT = 2; +class RecordingLogger implements ILogger { + readonly errors: Array = []; + + trace(): void {} + debug(): void {} + info(): void {} + warn(): void {} + fatal(): void {} + + error(obj: object | string, msg?: string): void { + this.errors.push(typeof obj === 'string' ? obj : (msg ?? '')); + } + + child(): ILogger { + return this; + } +} + +class UnreadableStateKVProvider extends MockKVProvider { + override async exists(): Promise { + throw new Error('kv unavailable'); + } +} + +class UnlockableKVProvider extends MockKVProvider { + override async acquireLock(): Promise { + throw new Error('kv lock unavailable'); + } +} + +class UnreleasableKVProvider extends MockKVProvider { + override async releaseLock(): Promise { + throw new Error('kv lock release failed'); + } +} + +function createFailingRepository(): UserRepository { + return { + async scanAllUsersPage() { + throw new Error('paged user scan failed'); + }, + } as unknown as UserRepository; +} + function createPendingUser(index: number): User { return new User({ ...EMPTY_USER_ROW, @@ -82,13 +127,47 @@ describe('ensureDeletionQueueState', () => { }, } as unknown as UserRepository; const queue = new KVAccountDeletionQueueService(kvClient, repository); + const logger = new RecordingLogger(); - await expect(ensureDeletionQueueState(queue, new NoopLogger())).resolves.toBeUndefined(); + await expect(ensureDeletionQueueState(queue, logger)).resolves.toBeUndefined(); expect(scans).toBe(1); + expect(logger.errors).toEqual(['KV deletion queue rebuild failed, leaving the rebuild to the deletion worker']); expect(await queue.acquireRebuildLock()).not.toBeNull(); }); + it('aborts startup when the queue state cannot be read', async () => { + const queue = new KVAccountDeletionQueueService(new UnreadableStateKVProvider(), createFailingRepository()); + const logger = new RecordingLogger(); + + await expect(ensureDeletionQueueState(queue, logger)).rejects.toThrow('kv unavailable'); + + expect(logger.errors).toEqual(['Failed to read KV deletion queue state, aborting startup']); + }); + + it('aborts startup when the rebuild lock cannot be acquired', async () => { + const queue = new KVAccountDeletionQueueService(new UnlockableKVProvider(), createFailingRepository()); + const logger = new RecordingLogger(); + + await expect(ensureDeletionQueueState(queue, logger)).rejects.toThrow('kv lock unavailable'); + + expect(logger.errors).toEqual(['Failed to acquire the KV deletion queue rebuild lock, aborting startup']); + }); + + it('does not abort startup when releasing the rebuild lock fails', async () => { + const queue = new KVAccountDeletionQueueService( + new UnreleasableKVProvider(), + createWorkerRepository(async () => {}), + ); + const logger = new RecordingLogger(); + + await expect(ensureDeletionQueueState(queue, logger)).resolves.toBeUndefined(); + + const queued = await queue.getReadyDeletions(Date.parse('2026-06-02T00:00:00.000Z'), 10); + expect(queued.map((deletion) => deletion.userId).sort()).toEqual([9101n, 9102n]); + expect(logger.errors).toEqual(['Failed to release the KV deletion queue rebuild lock']); + }); + it('rebuilds under the lock when no other instance holds it', async () => { const kvClient = new MockKVProvider(); const queue = new KVAccountDeletionQueueService(