mirror of
https://github.com/fluxerapp/fluxer.git
synced 2026-09-02 21:04:06 +03:00
fix(api): abort startup on unverifiable deletion queue state (#2311)
This commit is contained in:
@@ -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');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<string> = [];
|
||||
|
||||
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<number> {
|
||||
throw new Error('kv unavailable');
|
||||
}
|
||||
}
|
||||
|
||||
class UnlockableKVProvider extends MockKVProvider {
|
||||
override async acquireLock(): Promise<boolean> {
|
||||
throw new Error('kv lock unavailable');
|
||||
}
|
||||
}
|
||||
|
||||
class UnreleasableKVProvider extends MockKVProvider {
|
||||
override async releaseLock(): Promise<boolean> {
|
||||
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(
|
||||
|
||||
Reference in New Issue
Block a user