fix(worker): stop skipped accounts starving the deletion queue (#2294)

This commit is contained in:
Hampus
2026-09-01 00:16:08 +02:00
committed by GitHub
parent a2480c6a02
commit 662f4ac93b
6 changed files with 191 additions and 10 deletions
@@ -6,7 +6,7 @@ import {ms, seconds} from 'itty-time';
import type {UserID} from '../BrandedTypes';
import {Logger} from '../Logger';
import type {UserRepository} from '../user/repositories/UserRepository';
import {resolvePendingDeletionReasonCode} from '../user/services/PendingDeletionCoordinator';
import {isPendingDeletionBlocked, resolvePendingDeletionReasonCode} from '../user/services/PendingDeletionCoordinator';
interface QueuedDeletion {
userId: bigint;
@@ -76,7 +76,7 @@ export class KVAccountDeletionQueueService {
}
let batchQueued = 0;
for (const user of users) {
if (user.pendingDeletionAt) {
if (user.pendingDeletionAt && !isPendingDeletionBlocked(user)) {
const queueItem: QueuedDeletion = {
userId: user.id,
deletionReasonCode: resolvePendingDeletionReasonCode(user, 0),
@@ -0,0 +1,49 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import {describe, expect, it} 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';
function createUser(id: bigint, overrides: Partial<Pick<User, 'isBot' | 'flags'>> = {}): User {
return {
id: createUserID(id),
pendingDeletionAt: new Date('2026-06-01T00:00:00.000Z'),
deletionReasonCode: 0,
isBot: overrides.isBot ?? false,
flags: overrides.flags ?? 0n,
} as unknown as User;
}
function createUserRepository(users: Array<User>): UserRepository {
let served = false;
return {
async scanAllUsersPage() {
if (served) {
return {users: [], pageState: null};
}
served = true;
return {users, pageState: null};
},
} as unknown as UserRepository;
}
describe('KVAccountDeletionQueueService rebuild', () => {
it('does not requeue accounts the deletion worker refuses to process', async () => {
const kvClient = new MockKVProvider();
const users = [
createUser(1n, {isBot: true}),
createUser(2n, {flags: UserFlags.APP_STORE_REVIEWER}),
createUser(3n),
];
const service = new KVAccountDeletionQueueService(kvClient, createUserRepository(users));
await service.rebuildState();
expect(await service.getQueueSize()).toBe(1);
expect(await service.getReadyDeletions(Date.now(), 100)).toEqual([{userId: 3n, deletionReasonCode: 0}]);
});
});
@@ -16,6 +16,8 @@ function createPendingUser(index: number): User {
id: createUserID(BigInt(7000 + index)),
pendingDeletionAt: new Date('2026-06-01T00:00:00.000Z'),
deletionReasonCode: 0,
isBot: false,
flags: 0n,
} as unknown as User;
}
@@ -35,6 +35,11 @@ interface PendingDeletionReasonUserLike {
flags: bigint;
}
interface PendingDeletionEligibilityUserLike {
isBot: boolean;
flags: bigint;
}
export async function reschedulePendingDeletion({
userId,
currentPendingDeletionAt,
@@ -83,3 +88,10 @@ export function resolvePendingDeletionReasonCode(
}
return 0;
}
export function isPendingDeletionBlocked(user: PendingDeletionEligibilityUserLike): boolean {
if (user.isBot) {
return true;
}
return (user.flags & UserFlags.APP_STORE_REVIEWER) !== 0n;
}
@@ -1,10 +1,12 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
import {createUserID} from '../../BrandedTypes';
import {Logger} from '../../Logger';
import {resolvePendingDeletionReasonCode} from '../../user/services/PendingDeletionCoordinator';
import {
isPendingDeletionBlocked,
resolvePendingDeletionReasonCode,
} from '../../user/services/PendingDeletionCoordinator';
import {getWorkerDependencies} from '../WorkerContext';
const userProcessPendingDeletions: WorkerTaskHandler = async (_payload, helpers) => {
@@ -42,12 +44,9 @@ const userProcessPendingDeletions: WorkerTaskHandler = async (_payload, helpers)
await deletionQueueService.removeFromQueue(userId);
continue;
}
if (user.isBot) {
Logger.info({userId}, 'User is a bot, skipping deletion');
continue;
}
if (user.flags & UserFlags.APP_STORE_REVIEWER) {
Logger.info({userId}, 'User is an app store reviewer, skipping deletion');
if (isPendingDeletionBlocked(user)) {
Logger.info({userId}, 'User is not eligible for automated deletion, removing from KV');
await deletionQueueService.removeFromQueue(userId);
continue;
}
const deletionReasonCode = resolvePendingDeletionReasonCode(user, deletion.deletionReasonCode);
@@ -0,0 +1,119 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService';
import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask';
import {afterEach, describe, expect, test} from 'vitest';
import {createUserID, type UserID} from '../../BrandedTypes';
import {KVAccountDeletionQueueService} from '../../infrastructure/KVAccountDeletionQueueService';
import type {User} from '../../models/User';
import {MockKVProvider} from '../../test/mocks/MockKVProvider';
import {NoopLogger} from '../../test/mocks/NoopLogger';
import type {UserRepository} from '../../user/repositories/UserRepository';
import userProcessPendingDeletions from '../tasks/UserProcessPendingDeletions';
import {clearWorkerDependencies, setWorkerDependenciesForTest} from '../WorkerContext';
interface FakeUserOptions {
isBot?: boolean;
flags?: bigint;
}
function createFakeUser(userId: UserID, pendingDeletionAt: Date, options: FakeUserOptions = {}): User {
return {
id: userId,
pendingDeletionAt,
deletionReasonCode: 1,
isBot: options.isBot ?? false,
flags: options.flags ?? 0n,
} as unknown as User;
}
async function createHarness(users: Array<User>) {
const kvClient = new MockKVProvider();
const removedPendingDeletions: Array<string> = [];
const scheduledJobs: Array<string> = [];
const usersById = new Map<string, User>();
for (const user of users) {
usersById.set(user.id.toString(), user);
}
const userRepository = {
async findUnique(userId: UserID): Promise<User | null> {
return usersById.get(userId.toString()) ?? null;
},
async removePendingDeletion(userId: UserID): Promise<void> {
removedPendingDeletions.push(userId.toString());
},
} as unknown as UserRepository;
const workerService = {
async addJob(name: string, payload: {userId: string}): Promise<bigint> {
scheduledJobs.push(`${name}:${payload.userId}`);
return 0n;
},
} as unknown as IWorkerService;
const deletionQueueService = new KVAccountDeletionQueueService(kvClient, userRepository);
await kvClient.set('deletion_queue:state_version', Date.now().toString());
for (const user of users) {
if (user.pendingDeletionAt) {
await deletionQueueService.scheduleDeletion(user.id, user.pendingDeletionAt, 1);
}
}
setWorkerDependenciesForTest({userRepository, workerService, deletionQueueService});
return {deletionQueueService, scheduledJobs, removedPendingDeletions};
}
function createHelpers(): WorkerTaskHelpers {
return {
logger: new NoopLogger(),
jobId: 1n,
addJob: async () => 0n,
reportProgress: async () => {},
shouldCancel: async () => false,
setContextLink: async () => {},
};
}
describe('userProcessPendingDeletions', () => {
afterEach(() => {
clearWorkerDependencies();
});
test('drains skipped accounts so a genuine deletion behind them is not starved', async () => {
const skippedAt = new Date(Date.now() - 10_000_000);
const genuineAt = new Date(Date.now() - 1_000);
const users: Array<User> = [];
for (let i = 0; i < 1000; i++) {
const userId = createUserID(BigInt(100_000 + i));
users.push(
createFakeUser(userId, skippedAt, i % 2 === 0 ? {isBot: true} : {flags: UserFlags.APP_STORE_REVIEWER}),
);
}
const genuineUserId = createUserID(999_999n);
users.push(createFakeUser(genuineUserId, genuineAt));
const harness = await createHarness(users);
await userProcessPendingDeletions({}, createHelpers());
const queueSizeAfterFirstPass = await harness.deletionQueueService.getQueueSize();
await userProcessPendingDeletions({}, createHelpers());
expect(harness.scheduledJobs).toEqual([`userProcessPendingDeletion:${genuineUserId.toString()}`]);
expect(harness.removedPendingDeletions).toEqual([genuineUserId.toString()]);
expect(queueSizeAfterFirstPass).toBe(1);
expect(await harness.deletionQueueService.getQueueSize()).toBe(0);
});
test('keeps skipped accounts out of the queue while the skip condition holds', async () => {
const pendingAt = new Date(Date.now() - 10_000);
const botId = createUserID(1n);
const reviewerId = createUserID(2n);
const harness = await createHarness([
createFakeUser(botId, pendingAt, {isBot: true}),
createFakeUser(reviewerId, pendingAt, {flags: UserFlags.APP_STORE_REVIEWER}),
]);
await userProcessPendingDeletions({}, createHelpers());
expect(await harness.deletionQueueService.getQueueSize()).toBe(0);
expect(harness.scheduledJobs).toEqual([]);
expect(harness.removedPendingDeletions).toEqual([]);
});
});