From 9d95a808572246ddfea5fa4cda513698cee5ef02 Mon Sep 17 00:00:00 2001 From: Hampus Date: Mon, 31 Aug 2026 22:38:00 +0200 Subject: [PATCH] fix(worker): renew the deletion queue lock during a rebuild (#2287) --- .../KVAccountDeletionQueueService.ts | 16 +++++- .../KVAccountDeletionQueueRebuildLock.test.ts | 54 +++++++++++++++++++ .../tasks/UserProcessPendingDeletions.ts | 2 +- 3 files changed, 70 insertions(+), 2 deletions(-) create mode 100644 fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts diff --git a/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts b/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts index 661175056..e8d0554fa 100644 --- a/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts +++ b/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts @@ -59,7 +59,7 @@ export class KVAccountDeletionQueueService { } } - async rebuildState(): Promise { + async rebuildState(lockToken: string | null = null): Promise { Logger.info('Starting deletion queue rebuild from primary database'); try { await this.kvClient.del(QUEUE_KEY); @@ -91,6 +91,9 @@ export class KVAccountDeletionQueueService { totalQueued += batchQueued; totalProcessed += users.length; pageState = page.pageState; + if (lockToken !== null) { + await this.renewRebuildLock(lockToken); + } if (totalProcessed % 10000 === 0) { Logger.debug({totalProcessed, totalQueued}, 'Deletion queue rebuild progress'); } @@ -174,6 +177,17 @@ export class KVAccountDeletionQueueService { } } + private async renewRebuildLock(token: string): Promise { + try { + const renewed = await this.kvClient.extendLock(REBUILD_LOCK_KEY, token, REBUILD_LOCK_TTL); + if (!renewed) { + Logger.warn({token}, 'Deletion queue rebuild lock was no longer held on renewal'); + } + } catch (error) { + Logger.error({error, token}, 'Failed to renew deletion queue rebuild lock'); + } + } + async releaseRebuildLock(token: string): Promise { try { const released = await this.kvClient.releaseLock(REBUILD_LOCK_KEY, token); diff --git a/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts new file mode 100644 index 000000000..75f7a35c7 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts @@ -0,0 +1,54 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {ms} from 'itty-time'; +import {afterEach, describe, expect, it, vi} from 'vitest'; +import {createUserID} from '../../BrandedTypes'; +import type {User} from '../../models/User'; +import {MockKVProvider} from '../../test/mocks/MockKVProvider'; +import type {UserRepository} from '../../user/repositories/UserRepository'; +import {KVAccountDeletionQueueService} from '../KVAccountDeletionQueueService'; + +const PAGE_DURATION_MS = ms('3 minutes'); +const PAGE_COUNT = 3; + +function createPendingUser(index: number): User { + return { + id: createUserID(BigInt(7000 + index)), + pendingDeletionAt: new Date('2026-06-01T00:00:00.000Z'), + deletionReasonCode: 0, + } as unknown as User; +} + +function createSlowUserRepository(): UserRepository { + let page = 0; + return { + async scanAllUsersPage() { + vi.advanceTimersByTime(PAGE_DURATION_MS); + page += 1; + return { + users: [createPendingUser(page)], + pageState: page < PAGE_COUNT ? `page-${page}` : null, + }; + }, + } as unknown as UserRepository; +} + +describe('KVAccountDeletionQueueService rebuild lock', () => { + afterEach(() => { + vi.useRealTimers(); + }); + + it('keeps holding the rebuild lock across a scan longer than the lock ttl', async () => { + vi.useFakeTimers(); + vi.setSystemTime(new Date('2026-06-01T00:00:00.000Z')); + const kvClient = new MockKVProvider(); + const service = new KVAccountDeletionQueueService(kvClient, createSlowUserRepository()); + const token = await service.acquireRebuildLock(); + expect(token).not.toBeNull(); + + await service.rebuildState(token); + + expect(await service.acquireRebuildLock()).toBeNull(); + expect(await service.releaseRebuildLock(token!)).toBe(true); + }); +}); diff --git a/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts b/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts index 72b2b6ec8..90f6dd504 100644 --- a/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts +++ b/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts @@ -18,7 +18,7 @@ const userProcessPendingDeletions: WorkerTaskHandler = async (_payload, helpers) const lockToken = await deletionQueueService.acquireRebuildLock(); if (lockToken) { try { - await deletionQueueService.rebuildState(); + await deletionQueueService.rebuildState(lockToken); await deletionQueueService.releaseRebuildLock(lockToken); } catch (error) { await deletionQueueService.releaseRebuildLock(lockToken);