mirror of
https://github.com/fluxerapp/fluxer.git
synced 2026-09-02 21:04:06 +03:00
fix(worker): renew the deletion queue lock during a rebuild (#2287)
This commit is contained in:
@@ -59,7 +59,7 @@ export class KVAccountDeletionQueueService {
|
||||
}
|
||||
}
|
||||
|
||||
async rebuildState(): Promise<void> {
|
||||
async rebuildState(lockToken: string | null = null): Promise<void> {
|
||||
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<void> {
|
||||
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<boolean> {
|
||||
try {
|
||||
const released = await this.kvClient.releaseLock(REBUILD_LOCK_KEY, token);
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user