diff --git a/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts b/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts index e8d0554fa..5945c500b 100644 --- a/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts +++ b/fluxer_api/src/api/infrastructure/KVAccountDeletionQueueService.ts @@ -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), diff --git a/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuild.test.ts b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuild.test.ts new file mode 100644 index 000000000..33605e0a5 --- /dev/null +++ b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuild.test.ts @@ -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> = {}): 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): 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}]); + }); +}); diff --git a/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts index 75f7a35c7..257544ee8 100644 --- a/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts +++ b/fluxer_api/src/api/infrastructure/tests/KVAccountDeletionQueueRebuildLock.test.ts @@ -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; } diff --git a/fluxer_api/src/api/user/services/PendingDeletionCoordinator.ts b/fluxer_api/src/api/user/services/PendingDeletionCoordinator.ts index b7608b8bd..aef3e9ee6 100644 --- a/fluxer_api/src/api/user/services/PendingDeletionCoordinator.ts +++ b/fluxer_api/src/api/user/services/PendingDeletionCoordinator.ts @@ -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; +} diff --git a/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts b/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts index 90f6dd504..35f344e37 100644 --- a/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts +++ b/fluxer_api/src/api/worker/tasks/UserProcessPendingDeletions.ts @@ -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); diff --git a/fluxer_api/src/api/worker/tests/UserProcessPendingDeletions.test.ts b/fluxer_api/src/api/worker/tests/UserProcessPendingDeletions.test.ts new file mode 100644 index 000000000..327120e69 --- /dev/null +++ b/fluxer_api/src/api/worker/tests/UserProcessPendingDeletions.test.ts @@ -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) { + const kvClient = new MockKVProvider(); + const removedPendingDeletions: Array = []; + const scheduledJobs: Array = []; + const usersById = new Map(); + for (const user of users) { + usersById.set(user.id.toString(), user); + } + const userRepository = { + async findUnique(userId: UserID): Promise { + return usersById.get(userId.toString()) ?? null; + }, + async removePendingDeletion(userId: UserID): Promise { + removedPendingDeletions.push(userId.toString()); + }, + } as unknown as UserRepository; + const workerService = { + async addJob(name: string, payload: {userId: string}): Promise { + 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 = []; + 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([]); + }); +});