diff --git a/fluxer_api/src/api/Config.ts b/fluxer_api/src/api/Config.ts index 8a0d906a3..c02a6bf67 100644 --- a/fluxer_api/src/api/Config.ts +++ b/fluxer_api/src/api/Config.ts @@ -187,6 +187,7 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig { headersTimeoutMs: master.services.api.headers_timeout_ms, requestTimeoutMs: master.services.api.request_timeout_ms, maxInflightRequests: master.services.api.max_inflight_requests, + automatedMessageDeletionDelayDays: master.services.api.automated_message_deletion_delay_days, ipBanExemptIps: normalizeIpBanExemptIps(master.services.api.ip_ban_exempt_ips), cassandra: { hosts: cassandraSource?.hosts.join(',') ?? '', diff --git a/fluxer_api/src/api/config/APIConfig.ts b/fluxer_api/src/api/config/APIConfig.ts index f71332380..e87e2ed5f 100644 --- a/fluxer_api/src/api/config/APIConfig.ts +++ b/fluxer_api/src/api/config/APIConfig.ts @@ -54,6 +54,7 @@ export interface APIConfig { headersTimeoutMs: number; requestTimeoutMs: number; maxInflightRequests: number; + automatedMessageDeletionDelayDays: number; ipBanExemptIps: Array; cassandra: { hosts: string; diff --git a/fluxer_api/src/api/user/services/AccountStateApplier.ts b/fluxer_api/src/api/user/services/AccountStateApplier.ts index e973b949e..47b312807 100644 --- a/fluxer_api/src/api/user/services/AccountStateApplier.ts +++ b/fluxer_api/src/api/user/services/AccountStateApplier.ts @@ -2,7 +2,7 @@ import type {ApiContext} from '@app/api/ApiContext'; import type {AdminRepository} from '@app/api/admin/AdminRepository'; -import type {AdminMessageDeletionService} from '@app/api/admin/services/AdminMessageDeletionService'; +import type {AdminArchiveService} from '@app/api/admin/services/AdminArchiveService'; import {type ChannelID, createUserID, type MessageID, type UserID} from '@app/api/BrandedTypes'; import {isIpBanExempt} from '@app/api/ban/IpBanExemptions'; import type {IChannelRepository} from '@app/api/channel/IChannelRepository'; @@ -17,6 +17,7 @@ import type { OutcomeStatus, } from '@app/api/infrastructure/activity/Contract.generated'; import type {IGatewayService} from '@app/api/infrastructure/IGatewayService'; +import type {KVBulkMessageDeletionQueueService} from '@app/api/infrastructure/KVBulkMessageDeletionQueueService'; import type {User} from '@app/api/models/User'; import {isAccountLimitExempt} from '@app/api/user/AccountLimit'; import type {IUserRepository} from '@app/api/user/IUserRepository'; @@ -34,6 +35,7 @@ import {UserFlags} from '@fluxer/constants/src/UserConstants'; import {getSameIpDecisionKey, isPublicIpAddress, parseIpAddress} from '@fluxer/ip_utils/src/IpAddress'; import {snowflakeToDate} from '@fluxer/snowflake/src/Snowflake'; import type {ICacheService} from '@pkgs/cache/src/ICacheService'; +import {ms} from 'itty-time'; export type ActionOf = Extract; @@ -49,7 +51,9 @@ export interface AccountStateDeps { dispatch: AccountUpdateDispatch; ipBans: Pick; cache: Pick; - messages: Pick; + archives: Pick; + messageDeletionQueue: Pick; + messageDeletionDelayMs: number; authored: Pick; now?: () => number; } @@ -57,6 +61,8 @@ export interface AccountStateDeps { const FLAGS_WRITE_ATTEMPTS = 3; const AUTHORED_PAGE_SIZE = 200; const MIN_TEMP_BAN_SECONDS = 60; +const RECENT_ARCHIVE_MS = ms('1 day'); +const RECENT_ARCHIVE_SCAN = 20; type ProfilePropagation = Omit; @@ -94,7 +100,7 @@ export function accountStateDepsFromContext( ctx: ApiContext, ipBans: AccountStateDeps['ipBans'], profile: Omit, - messages: AccountStateDeps['messages'], + messageDeletion: Pick, channels: Pick, ): AccountStateDeps { return { @@ -102,7 +108,7 @@ export function accountStateDepsFromContext( dispatch: gatewayDispatch(ctx.services.gateway, {...profile, userRepository: ctx.services.users}, channels), ipBans, cache: ctx.services.cache, - messages, + ...messageDeletion, authored: channels, }; } @@ -135,6 +141,10 @@ function isIneligible(user: User): boolean { return user.isBot || (user.flags & UserFlags.DELETED) !== 0n; } +function isDeletionComplete(user: User): boolean { + return (user.flags & UserFlags.DELETED) !== 0n && user.pendingDeletionAt === null; +} + export async function applySetAccountLimit( deps: AccountStateDeps, env: ActionOf<'set_account_limit'>, @@ -190,7 +200,7 @@ export async function applyHideRecentMessages( return withAccountChangeSource('action', async () => { const user = await deps.users.findUnique(createUserID(BigInt(env.user_id))); if (!user) return outcomeOf(env, 'ineligible'); - if (isIneligible(user)) return outcomeOf(env, 'ineligible', user); + if (user.isBot || isDeletionComplete(user)) return outcomeOf(env, 'ineligible', user); const current = user.contentHiddenSince?.getTime() ?? null; if (!env.on) { if (current === null) return outcomeOf(env, 'noop', user); @@ -235,18 +245,64 @@ export async function applyDeleteUserMessages( deps: AccountStateDeps, env: ActionOf<'delete_user_messages'>, ): Promise { - if (!env.on) return outcomeOf(env, 'unsupported', null, 'a message purge cannot be reversed'); const user = await deps.users.findUnique(createUserID(BigInt(env.user_id))); if (!user) return outcomeOf(env, 'ineligible'); - if (user.isBot) return outcomeOf(env, 'ineligible', user); + if (!env.on) return cancelScheduledMessageDeletion(deps, env, user); + if (user.isBot || isDeletionComplete(user)) return outcomeOf(env, 'ineligible', user); if (isAccountLimitExempt(user)) return outcomeOf(env, 'exempt', user); - const purge = await deps.messages.deleteAllUserMessages( - {user_id: user.id, dry_run: false}, - SYSTEM_USER_ID, - `Automated action ${env.id}`, + const now = deps.now?.() ?? Date.now(); + const archive = await archiveBeforeDeletion(deps, user.id, now); + const target = now + deps.messageDeletionDelayMs; + const current = user.pendingBulkMessageDeletionAt?.getTime() ?? null; + const scheduledAt = current === null ? target : Math.min(current, target); + await deps.messageDeletionQueue.scheduleDeletion(user.id, new Date(scheduledAt)); + const detail = `${archive} scheduled_at=${new Date(scheduledAt).toISOString()}`; + if (scheduledAt === current) return outcomeOf(env, 'noop', user, detail); + const scheduled = await deps.users.patchUpsert( + user.id, + { + pending_bulk_message_deletion_at: new Date(scheduledAt), + pending_bulk_message_deletion_channel_count: null, + pending_bulk_message_deletion_message_count: null, + }, + user.toRow(), ); - if (purge.message_count === 0) return outcomeOf(env, 'noop', user); - return outcomeOf(env, 'applied', user, `messages=${purge.message_count} job=${purge.job_id ?? ''}`); + await deps.dispatch.userUpdated(scheduled); + return outcomeOf(env, 'applied', scheduled, detail); +} + +async function archiveBeforeDeletion(deps: AccountStateDeps, userId: UserID, now: number): Promise { + const archives = await deps.archives.listArchives({ + subjectType: 'user', + subjectId: userId, + limit: RECENT_ARCHIVE_SCAN, + }); + const recent = archives.find( + (archive) => archive.failed_at === null && now - Date.parse(archive.requested_at) < RECENT_ARCHIVE_MS, + ); + if (recent) return `archive=${recent.archive_id} reused=true`; + const created = await deps.archives.triggerUserArchive(userId, SYSTEM_USER_ID, true); + return `archive=${created.archive_id} reused=false`; +} + +async function cancelScheduledMessageDeletion( + deps: AccountStateDeps, + env: ActionOf<'delete_user_messages'>, + user: User, +): Promise { + if (user.pendingBulkMessageDeletionAt === null) return outcomeOf(env, 'noop', user); + const cancelled = await deps.users.patchUpsert( + user.id, + { + pending_bulk_message_deletion_at: null, + pending_bulk_message_deletion_channel_count: null, + pending_bulk_message_deletion_message_count: null, + }, + user.toRow(), + ); + await deps.messageDeletionQueue.removeFromQueue(user.id); + await deps.dispatch.userUpdated(cancelled); + return outcomeOf(env, 'applied', cancelled); } function parseBanTarget(value: string): ReturnType { diff --git a/fluxer_api/src/api/user/services/UserContentService.ts b/fluxer_api/src/api/user/services/UserContentService.ts index 1f63c63ee..216f41693 100644 --- a/fluxer_api/src/api/user/services/UserContentService.ts +++ b/fluxer_api/src/api/user/services/UserContentService.ts @@ -42,6 +42,7 @@ import {UserHarvestRepository} from '@app/api/user/UserHarvestRepository'; import {serializeSelfMessageFilter} from '@app/api/worker/utils/SelfMessageFilterPayload'; import type {WorkerTaskName} from '@app/api/worker/WorkerLaneConfig'; import {MAX_BOOKMARKS_NON_PREMIUM} from '@fluxer/constants/src/LimitConstants'; +import {UserFlags} from '@fluxer/constants/src/UserConstants'; import {ValidationErrorCodes} from '@fluxer/constants/src/ValidationErrorCodes'; import {UnknownChannelError} from '@fluxer/errors/src/domains/channel/UnknownChannelError'; import {UnknownMessageError} from '@fluxer/errors/src/domains/channel/UnknownMessageError'; @@ -748,6 +749,7 @@ export class UserContentService { async cancelBulkMessageDeletion(userId: UserID): Promise { Logger.debug({userId: userId.toString()}, 'Canceling pending bulk message deletion'); const user = await this.userRepository.findUniqueAssert(userId); + if ((user.flags & UserFlags.SPAMMER) !== 0n) return; const updatedUser = await this.userRepository.patchUpsert( userId, { diff --git a/fluxer_api/src/api/user/services/UserDeletionService.ts b/fluxer_api/src/api/user/services/UserDeletionService.ts index 864b92b3a..9d8690df2 100644 --- a/fluxer_api/src/api/user/services/UserDeletionService.ts +++ b/fluxer_api/src/api/user/services/UserDeletionService.ts @@ -7,6 +7,7 @@ import {createMessageID, createUserID, type MessageID, type UserID} from '@app/a import {Config} from '@app/api/Config'; import {mapChannelToResponse} from '@app/api/channel/ChannelMappers'; import type {ChannelRepository} from '@app/api/channel/ChannelRepository'; +import {UserMessageDeletionService} from '@app/api/channel/services/message/UserMessageDeletionService'; import type {IConnectionRepository} from '@app/api/connection/IConnectionRepository'; import type {FavoriteMemeRepository} from '@app/api/favorite_meme/FavoriteMemeRepository'; import type {GuildRepository} from '@app/api/guild/repositories/GuildRepository'; @@ -371,6 +372,25 @@ export async function processUserDeletion( Logger.error({error, userId, channelId: channel.id}, 'Failed to leave group DM'); } } + if (user.pendingBulkMessageDeletionAt) { + const deleted = await new UserMessageDeletionService({ + channelRepository, + gatewayService, + storageService, + purgeQueue, + workerService, + }).deleteUserMessagesBulk(userId); + await userRepository.patchUpsert( + userId, + { + pending_bulk_message_deletion_at: null, + pending_bulk_message_deletion_channel_count: null, + pending_bulk_message_deletion_message_count: null, + }, + user.toRow(), + ); + Logger.debug({userId, deleted}, 'Deleted messages scheduled for deletion before anonymizing the rest'); + } Logger.debug({userId}, 'Anonymizing user messages'); let lastMessageId: MessageID | undefined; let processedCount = 0; diff --git a/fluxer_api/src/api/user/tests/AccountDeleteScheduledMessageDeletion.test.ts b/fluxer_api/src/api/user/tests/AccountDeleteScheduledMessageDeletion.test.ts new file mode 100644 index 000000000..784b76ea8 --- /dev/null +++ b/fluxer_api/src/api/user/tests/AccountDeleteScheduledMessageDeletion.test.ts @@ -0,0 +1,66 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createTestAccount} from '@app/api/auth/tests/AuthTestUtils'; +import {acceptInvite, createChannel, createChannelInvite, createGuild} from '@app/api/guild/tests/GuildTestUtils'; +import {ensureSessionStarted, sendMessage} from '@app/api/message/tests/MessageTestUtils'; +import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {HTTP_STATUS} from '@app/api/test/TestConstants'; +import {createBuilder, createBuilderWithoutAuth} from '@app/api/test/TestRequestBuilder'; +import { + deleteAccount, + setPendingDeletionAt, + triggerDeletionWorker, + waitForDeletionCompletion, +} from '@app/api/user/tests/UserTestUtils'; +import type {MessageResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas'; +import {afterEach, beforeEach, describe, expect, test} from 'vitest'; + +describe('Account deletion with a scheduled message deletion', () => { + let harness: ApiTestHarness; + beforeEach(async () => { + harness = await createApiTestHarness(); + }); + afterEach(async () => { + await harness?.shutdown(); + }); + test('deletes the scheduled messages instead of anonymizing them', async () => { + const account = await createTestAccount(harness); + const guild = await createGuild(harness, account.token, 'Scheduled Deletion Guild'); + let channelId = guild.system_channel_id; + if (!channelId) { + const channel = await createChannel(harness, account.token, guild.id, 'general'); + channelId = channel.id; + } + await ensureSessionStarted(harness, account.token); + for (let i = 0; i < 3; i++) { + await sendMessage(harness, account.token, channelId, `Scheduled ${i + 1}`); + } + const newOwner = await createTestAccount(harness); + const invite = await createChannelInvite(harness, account.token, channelId); + await acceptInvite(harness, newOwner.token, invite.code); + await createBuilder(harness, account.token) + .post(`/guilds/${guild.id}/transfer-ownership`) + .body({new_owner_id: newOwner.userId, password: account.password}) + .expect(HTTP_STATUS.OK) + .execute(); + await createBuilder(harness, account.token) + .post('/users/@me/messages/delete') + .body({password: account.password}) + .expect(HTTP_STATUS.NO_CONTENT) + .execute(); + await deleteAccount(harness, account.token, account.password); + await setPendingDeletionAt(harness, account.userId, new Date(Date.now() - 60_000)); + await triggerDeletionWorker(harness); + await waitForDeletionCompletion(harness, account.userId); + const messages = await createBuilder>(harness, newOwner.token) + .get(`/channels/${channelId}/messages?limit=100`) + .expect(HTTP_STATUS.OK) + .execute(); + expect(messages.filter((message) => message.content.startsWith('Scheduled '))).toEqual([]); + const countJson = await createBuilderWithoutAuth<{count: number}>(harness) + .get(`/test/users/${account.userId}/messages/count`) + .expect(HTTP_STATUS.OK) + .execute(); + expect(countJson.count).toBe(0); + }, 60000); +}); diff --git a/fluxer_api/src/api/user/tests/ScheduledMessageDeletionCancel.test.ts b/fluxer_api/src/api/user/tests/ScheduledMessageDeletionCancel.test.ts new file mode 100644 index 000000000..5128edc40 --- /dev/null +++ b/fluxer_api/src/api/user/tests/ScheduledMessageDeletionCancel.test.ts @@ -0,0 +1,52 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createTestAccount} from '@app/api/auth/tests/AuthTestUtils'; +import {type ApiTestHarness, createApiTestHarness} from '@app/api/test/ApiTestHarness'; +import {HTTP_STATUS} from '@app/api/test/TestConstants'; +import {createBuilder} from '@app/api/test/TestRequestBuilder'; +import {UserFlags} from '@fluxer/constants/src/UserConstants'; +import {afterEach, beforeEach, describe, expect, test} from 'vitest'; + +interface Me { + id: string; + flags?: string | number; + pending_bulk_message_deletion: {scheduled_at: string} | null; +} + +describe('Cancelling a scheduled message deletion', () => { + let harness: ApiTestHarness; + beforeEach(async () => { + harness = await createApiTestHarness(); + }); + afterEach(async () => { + await harness?.shutdown(); + }); + test('a flagged account cannot cancel it until the flag is cleared', async () => { + const account = await createTestAccount(harness); + const me = () => createBuilder(harness, account.token).get('/users/@me').expect(HTTP_STATUS.OK).execute(); + const setFlags = (flags: bigint) => + createBuilder(harness, account.token) + .patch(`/test/users/${account.userId}/flags`) + .body({flags: flags.toString()}) + .execute(); + const cancel = () => + createBuilder<{success: boolean}>(harness, account.token) + .delete('/users/@me/messages/delete') + .expect(HTTP_STATUS.OK) + .execute(); + await createBuilder(harness, account.token) + .post('/users/@me/messages/delete') + .body({password: account.password}) + .expect(HTTP_STATUS.NO_CONTENT) + .execute(); + const scheduled = (await me()).pending_bulk_message_deletion; + expect(scheduled).not.toBeNull(); + const flags = BigInt((await me()).flags ?? 0); + await setFlags(flags | UserFlags.SPAMMER); + await cancel(); + expect((await me()).pending_bulk_message_deletion?.scheduled_at).toBe(scheduled?.scheduled_at); + await setFlags(flags); + await cancel(); + expect((await me()).pending_bulk_message_deletion).toBeNull(); + }, 60000); +}); diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index b6921cf7a..b4ee997d2 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -1,8 +1,5 @@ // SPDX-License-Identifier: AGPL-3.0-or-later -import {AdminAuditService} from '@app/api/admin/services/AdminAuditService'; -import {AdminMessageDeletionService} from '@app/api/admin/services/AdminMessageDeletionService'; -import {AdminMessageShredService} from '@app/api/admin/services/AdminMessageShredService'; import {Config} from '@app/api/Config'; import {createApiContext} from '@app/api/CreateApiContext'; import {setDatabaseQueryExecutor} from '@app/api/database/CassandraQueryExecution'; @@ -29,6 +26,7 @@ import { shutdownVoiceResources, } from '@app/api/middleware/ServiceRegistry'; import { + getAdminArchiveService, getAdminRepository, getCacheService, getInstanceConfigRepository, @@ -58,6 +56,7 @@ import {BACKGROUND_READ_TIMEOUT_MS, initCassandra, shutdownCassandra} from '@pkg import {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectionManager'; import {getDefaultPostgresClient, initPostgres, shutdownPostgres} from '@pkgs/postgres/src/Client'; import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask'; +import {ms} from 'itty-time'; function registerCronJobs(cron: CronScheduler, jobsStreamMaxAgeMs: number): void { cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *', {ledger: false}); @@ -286,12 +285,6 @@ export async function startWorkerMain(): Promise { startSharedListWatch(jsConnectionManager.getJetStreamClient()); if (activeWorkerLanes.some((lane) => lane.name === 'lifecycle')) { const apiContext = createApiContext(); - const auditService = new AdminAuditService(getAdminRepository(), apiContext.services.snowflake); - const messagePurge = new AdminMessageDeletionService({ - channelRepository: dependencies.channelRepository, - messageShredService: new AdminMessageShredService({apiContext, auditService}), - auditService, - }); startAccountActionConsumer({ js: jsConnectionManager.getJetStreamClient(), state: accountStateDepsFromContext( @@ -301,7 +294,11 @@ export async function startWorkerMain(): Promise { userCacheService: dependencies.userCacheService, guildRepository: dependencies.guildRepository, }, - messagePurge, + { + archives: getAdminArchiveService(), + messageDeletionQueue: dependencies.bulkMessageDeletionQueueService, + messageDeletionDelayMs: Config.automatedMessageDeletionDelayDays * ms('1 day'), + }, dependencies.channelRepository, ), }); diff --git a/fluxer_api/src/api/worker/tests/AccountActionConsumer.test.ts b/fluxer_api/src/api/worker/tests/AccountActionConsumer.test.ts index 03cba33e5..fe33d0fcc 100644 --- a/fluxer_api/src/api/worker/tests/AccountActionConsumer.test.ts +++ b/fluxer_api/src/api/worker/tests/AccountActionConsumer.test.ts @@ -1,6 +1,5 @@ // SPDX-License-Identifier: AGPL-3.0-or-later -import {AdminMessageDeletionService} from '@app/api/admin/services/AdminMessageDeletionService'; import { type ChannelID, createChannelID, @@ -29,6 +28,7 @@ import { stopAccountActionConsumer, } from '@app/api/worker/AccountActionConsumer'; import {UserFlags} from '@fluxer/constants/src/UserConstants'; +import type {AdminArchiveResponse} from '@fluxer/schema/src/domains/admin/AdminArchiveSchemas'; import {createSnowflakeFromTimestamp} from '@fluxer/snowflake/src/Snowflake'; import {AckPolicy, DeliverPolicy, type JsMsg, jetstream, jetstreamManager} from '@nats-io/jetstream'; import {connect} from '@nats-io/transport-node'; @@ -37,6 +37,33 @@ import {afterAll, beforeAll, beforeEach, describe, expect, it} from 'vitest'; const USER_ID = '1174109840998400001'; const NOW = 1_759_000_000_000; const NATS_URL = process.env.FLUXER_TEST_ACTIVITY_NATS_URL; +const DAY_MS = 86_400_000; +const DELETION_DELAY_MS = 7 * DAY_MS; + +function archiveAt( + requestedAtMs: number, + userId: UserID, + requestedBy: UserID, + archiveId: string, + failed = false, +): AdminArchiveResponse { + return { + archive_id: archiveId, + subject_type: 'user', + subject_id: userId.toString(), + requested_by: requestedBy.toString(), + requested_at: new Date(requestedAtMs).toISOString(), + started_at: null, + completed_at: null, + failed_at: failed ? new Date(requestedAtMs + 1_000).toISOString() : null, + file_size: null, + progress_percent: 0, + progress_step: 'Queued', + error_message: null, + download_url_expires_at: null, + expires_at: null, + }; +} function limitEnvelope(overrides: Partial> = {}) { return envelope<'set_account_limit'>({ @@ -108,8 +135,9 @@ interface Harness { authored: Array<{channelId: ChannelID; messageId: MessageID}>; visibility: Array; removed: Array<{channelId: ChannelID; authorId: UserID; messageIds: Array}>; - audits: Array<{adminUserId: UserID; action: string; targetId: bigint; auditLogReason: string | null}>; - shreds: Array<{userId: bigint; entries: number; adminUserId: UserID; auditLogReason: string | null}>; + archives: Array; + archiveRequests: Array<{userId: UserID; requestedBy: UserID; includeAttachments: boolean}>; + queued: Map; deps: AccountActionDeps; } @@ -125,8 +153,9 @@ function harness(): Harness { authored: [], visibility: [], removed: [], - audits: [], - shreds: [], + archives: [], + archiveRequests: [], + queued: new Map(), deps: null as never, }; const authored = { @@ -136,34 +165,6 @@ function harness(): Harness { .filter((m) => before === undefined || m.messageId < before) .slice(0, limit), }; - const messages = new AdminMessageDeletionService({ - channelRepository: authored, - messageShredService: { - queueMessageShred: async ( - data: {user_id: bigint; entries: Array}, - adminUserId: UserID, - auditLogReason: string | null, - ) => { - h.shreds.push({userId: data.user_id, entries: data.entries.length, adminUserId, auditLogReason}); - return {success: true, job_id: '77', requested: data.entries.length}; - }, - }, - auditService: { - createAuditLog: async (log: { - adminUserId: UserID; - action: string; - targetId: bigint; - auditLogReason: string | null; - }) => { - h.audits.push({ - adminUserId: log.adminUserId, - action: log.action, - targetId: log.targetId, - auditLogReason: log.auditLogReason, - }); - }, - }, - } as unknown as ConstructorParameters[0]); const state: AccountStateDeps = { users: users as unknown as AccountStateDeps['users'], dispatch: { @@ -198,7 +199,25 @@ function harness(): Harness { h.cached.delete(key); }, } as unknown as AccountStateDeps['cache'], - messages, + archives: { + listArchives: async ({subjectId}) => + h.archives.filter((archive) => archive.subject_id === subjectId?.toString()).reverse(), + triggerUserArchive: async (userId, requestedBy, includeAttachments) => { + h.archiveRequests.push({userId, requestedBy, includeAttachments}); + const archive = archiveAt(NOW, userId, requestedBy, String(9000 + h.archives.length)); + h.archives.push(archive); + return archive; + }, + }, + messageDeletionQueue: { + scheduleDeletion: async (userId, scheduledAt) => { + h.queued.set(userId.toString(), scheduledAt.getTime()); + }, + removeFromQueue: async (userId) => { + h.queued.delete(userId.toString()); + }, + }, + messageDeletionDelayMs: DELETION_DELAY_MS, authored: authored as unknown as AccountStateDeps['authored'], now: () => NOW, }; @@ -453,7 +472,7 @@ describe('account action apply', () => { } } expect(h.users.current().flags).toBe(0n); - expect(h.shreds).toHaveLength(0); + expect(h.queued.size).toBe(0); }); it('keeps the wider window when asked again and resends the deletes as a noop', async () => { @@ -529,103 +548,121 @@ describe('account action apply', () => { expect(h.users.current().contentHiddenSince?.getTime()).toBe(since); }); - it('still purges a hidden account through the admin purge', async () => { - const since = NOW - 86_400_000; - authorAt(since - 5_000, 100); - authorAt(since + 5_000, 101); - await applyAction(h.deps, hideEnvelope(true, since)); - const purge = await applyAction( - h.deps, - envelope<'delete_user_messages'>({type: 'delete_user_messages', user_id: USER_ID, on: true}), - ); - expect(purge).toMatchObject({status: 'applied', detail: 'messages=2 job=77'}); - expect(h.shreds).toEqual([expect.objectContaining({entries: 2})]); + it('hides every message the account ever sent when the window starts at zero', async () => { + const all = [authorAt(1_420_070_401_000, 100), authorAt(NOW - 400 * DAY_MS, 101), authorAt(NOW - 1_000, 102)]; + const outcome = await applyAction(h.deps, hideEnvelope(true, 0)); + expect(outcome).toMatchObject({status: 'applied', detail: 'messages=3'}); + expect(h.users.current().contentHiddenSince?.getTime()).toBe(0); + expect(removedIds()).toEqual([...all].sort((a, b) => (a < b ? -1 : 1))); + expect((await applyAction(h.deps, hideEnvelope(true, NOW - DAY_MS))).status).toBe('noop'); + expect(h.users.current().contentHiddenSince?.getTime()).toBe(0); + expect((await applyAction(h.deps, hideEnvelope(false, 0))).status).toBe('applied'); + expect(h.users.current().contentHiddenSince).toBeNull(); + }); + + it('hides the messages of an account whose deletion is only scheduled', async () => { + authorAt(NOW - 1_000, 100); + h.users.put({flags: UserFlags.DELETED | UserFlags.SPAMMER, pending_deletion_at: new Date(NOW + 60 * DAY_MS)}); + expect((await applyAction(h.deps, hideEnvelope(true, 0))).status).toBe('applied'); + expect(h.users.current().contentHiddenSince?.getTime()).toBe(0); }); function purgeEnvelope(on = true) { return envelope<'delete_user_messages'>({type: 'delete_user_messages', user_id: USER_ID, on}); } - function author(count: number): void { - for (let i = 0; i < count; i++) { - h.authored.push({ - channelId: createChannelID(BigInt(100 + (i % 3))), - messageId: createMessageID(BigInt(1000 + i)), - }); - } - } + const userId = createUserID(BigInt(USER_ID)); - it('purges every message through the admin purge and audits it as the system', async () => { - author(450); + it('archives with attachments and schedules the deletion instead of deleting now', async () => { + authorAt(NOW - 1_000, 100); const outcome = await applyAction(h.deps, purgeEnvelope()); + const scheduledAt = NOW + DELETION_DELAY_MS; expect(outcome).toEqual({ action_id: 'a:07:4242:0', action_type: 'delete_user_messages', status: 'applied', - detail: 'messages=450 job=77', + detail: `archive=9000 reused=false scheduled_at=${new Date(scheduledAt).toISOString()}`, observed: {flags: '0', deleted: false}, user_id: USER_ID, }); - expect(h.shreds).toEqual([ - { - userId: BigInt(USER_ID), - entries: 450, - adminUserId: SYSTEM_USER_ID, - auditLogReason: 'Automated action a:07:4242:0', - }, - ]); - expect(h.audits).toEqual([ - { - adminUserId: SYSTEM_USER_ID, - action: 'delete_all_user_messages', - targetId: BigInt(USER_ID), - auditLogReason: 'Automated action a:07:4242:0', - }, - ]); - expect(h.users.current().flags).toBe(0n); - expect(h.presence).toHaveLength(0); + expect(h.archiveRequests).toEqual([{userId, requestedBy: SYSTEM_USER_ID, includeAttachments: true}]); + expect(h.users.current().pendingBulkMessageDeletionAt?.getTime()).toBe(scheduledAt); + expect(h.queued.get(USER_ID)).toBe(scheduledAt); + expect(h.removed).toEqual([]); + expect(h.authored).toHaveLength(1); + expect(h.presence).toEqual([userId]); }); - it('purges accounts already scheduled for deletion and paid accounts', async () => { - author(2); + it('reuses an archive from the last day and makes a new one when it is older or failed', async () => { + h.archives.push(archiveAt(NOW - DAY_MS + 60_000, userId, createUserID(42n), '777')); + expect((await applyAction(h.deps, purgeEnvelope())).detail).toMatch(/^archive=777 reused=true /); + expect(h.archiveRequests).toEqual([]); + h.archives.length = 0; + h.archives.push(archiveAt(NOW - DAY_MS - 60_000, userId, createUserID(42n), '778')); + h.archives.push(archiveAt(NOW - 60_000, userId, createUserID(42n), '779', true)); + h.users.put(); + expect((await applyAction(h.deps, purgeEnvelope())).detail).toMatch(/^archive=9002 reused=false /); + expect(h.archiveRequests).toHaveLength(1); + }); + + it('keeps an earlier schedule and pulls a later one forward', async () => { + const earlier = NOW + DAY_MS; + h.users.put({pending_bulk_message_deletion_at: new Date(earlier), pending_bulk_message_deletion_message_count: 4}); + const kept = await applyAction(h.deps, purgeEnvelope()); + expect(kept).toMatchObject({status: 'noop'}); + expect(kept.detail).toContain(`scheduled_at=${new Date(earlier).toISOString()}`); + expect(h.users.current().pendingBulkMessageDeletionAt?.getTime()).toBe(earlier); + expect(h.users.current().pendingBulkMessageDeletionMessageCount).toBe(4); + expect(h.queued.get(USER_ID)).toBe(earlier); + expect(h.presence).toHaveLength(0); + h.users.put({pending_bulk_message_deletion_at: new Date(NOW + 30 * DAY_MS)}); + expect((await applyAction(h.deps, purgeEnvelope())).status).toBe('applied'); + expect(h.users.current().pendingBulkMessageDeletionAt?.getTime()).toBe(NOW + DELETION_DELAY_MS); + expect(h.queued.get(USER_ID)).toBe(NOW + DELETION_DELAY_MS); + expect(h.archiveRequests).toHaveLength(1); + }); + + it('schedules for accounts already scheduled for deletion and paid accounts', async () => { for (const overrides of [ - {flags: UserFlags.DELETED | UserFlags.SPAMMER, pending_deletion_at: new Date(NOW + 86_400_000)}, + {flags: UserFlags.DELETED | UserFlags.SPAMMER, pending_deletion_at: new Date(NOW + 60 * DAY_MS)}, {has_ever_purchased: true, premium_type: 2}, ] satisfies Array>) { h.users.put(overrides); expect((await applyAction(h.deps, purgeEnvelope())).status).toBe('applied'); + expect(h.users.current().pendingBulkMessageDeletionAt?.getTime()).toBe(NOW + DELETION_DELAY_MS); } - expect(h.shreds).toHaveLength(2); }); - it('reports a noop when nothing is left to purge', async () => { - const outcome = await applyAction(h.deps, purgeEnvelope()); - expect(outcome).toMatchObject({status: 'noop', detail: null}); - expect(h.shreds).toHaveLength(0); + it('cancels the schedule on off and is a noop when nothing is scheduled', async () => { + await applyAction(h.deps, purgeEnvelope()); + const cancelled = await applyAction(h.deps, purgeEnvelope(false)); + expect(cancelled).toMatchObject({status: 'applied', detail: null, observed: {flags: '0', deleted: false}}); + expect(h.users.current().pendingBulkMessageDeletionAt).toBeNull(); + expect(h.queued.has(USER_ID)).toBe(false); + expect(h.presence).toHaveLength(2); + expect((await applyAction(h.deps, purgeEnvelope(false))).status).toBe('noop'); + expect(h.archives).toHaveLength(1); + h.users.put({flags: UserFlags.STAFF, pending_bulk_message_deletion_at: new Date(NOW + DAY_MS)}); + expect((await applyAction(h.deps, purgeEnvelope(false))).status).toBe('applied'); + h.users.rows.clear(); + expect(await applyAction(h.deps, purgeEnvelope(false))).toMatchObject({status: 'ineligible', observed: null}); }); - it('never purges staff, trusted, system, bot or missing accounts', async () => { - author(3); + it('never schedules for staff, trusted, system, bot, deleted or missing accounts', async () => { for (const overrides of [{flags: UserFlags.STAFF}, {flags: UserFlags.LIMIT_EXEMPT}, {system: true}] satisfies Array< Partial >) { h.users.put(overrides); expect((await applyAction(h.deps, purgeEnvelope())).status).toBe('exempt'); } - h.users.put({bot: true}); - expect((await applyAction(h.deps, purgeEnvelope())).status).toBe('ineligible'); + for (const overrides of [{bot: true}, {flags: UserFlags.DELETED}] satisfies Array>) { + h.users.put(overrides); + expect((await applyAction(h.deps, purgeEnvelope())).status).toBe('ineligible'); + } h.users.rows.clear(); expect(await applyAction(h.deps, purgeEnvelope())).toMatchObject({status: 'ineligible', observed: null}); - expect(h.shreds).toHaveLength(0); - expect(h.audits).toHaveLength(0); - }); - - it('answers a purge with on false as unsupported, since it cannot be reversed', async () => { - author(3); - const outcome = await applyAction(h.deps, purgeEnvelope(false)); - expect(outcome).toMatchObject({status: 'unsupported', detail: 'a message purge cannot be reversed'}); - expect(h.shreds).toHaveLength(0); - expect(h.audits).toHaveLength(0); + expect(h.archiveRequests).toEqual([]); + expect(h.queued.size).toBe(0); }); it('answers expired actions and unknown shapes without touching the account', async () => { diff --git a/packages/config/src/ConfigLoader.ts b/packages/config/src/ConfigLoader.ts index 7d618dd89..966c2bbf2 100644 --- a/packages/config/src/ConfigLoader.ts +++ b/packages/config/src/ConfigLoader.ts @@ -102,6 +102,7 @@ function defaultConfig(): MasterConfig { headers_timeout_ms: 30_000, request_timeout_ms: 120_000, max_inflight_requests: 512, + automated_message_deletion_delay_days: 7, ip_ban_exempt_ips: [], donation_proxy_key: '', trusted_callers: [], @@ -604,6 +605,12 @@ function normalizeConfig(config: MasterConfig): MasterConfig { validateReplyToEmail(config.integrations.email.reply_to_email); normalizeAppOriginAliases(config); assertIntegerInRange(config.services.api.max_inflight_requests, 'FLUXER_API_MAX_INFLIGHT_REQUESTS', 1, 100_000); + assertIntegerInRange( + config.services.api.automated_message_deletion_delay_days, + 'FLUXER_API_AUTOMATED_MESSAGE_DELETION_DELAY_DAYS', + 1, + 365, + ); assertIntegerInRange(config.services.api.headers_timeout_ms, 'FLUXER_API_HEADERS_TIMEOUT_MS', 1_000, 3_600_000); assertIntegerInRange(config.services.api.request_timeout_ms, 'FLUXER_API_REQUEST_TIMEOUT_MS', 1_000, 3_600_000); assertIntegerInRange(config.domain.public_port, 'FLUXER_PUBLIC_PORT', 1, 65_535); diff --git a/packages/config/src/MasterConfig.ts b/packages/config/src/MasterConfig.ts index 3832d9a81..3b3c052c6 100644 --- a/packages/config/src/MasterConfig.ts +++ b/packages/config/src/MasterConfig.ts @@ -88,6 +88,7 @@ export interface MasterConfig { headers_timeout_ms: number; request_timeout_ms: number; max_inflight_requests: number; + automated_message_deletion_delay_days: number; ip_ban_exempt_ips: Array; donation_proxy_key: string; trusted_callers: Array<{ diff --git a/packages/config/src/__tests__/SelfHostingComposeEnvironment.test.ts b/packages/config/src/__tests__/SelfHostingComposeEnvironment.test.ts index 8cf501861..81a40e14a 100644 --- a/packages/config/src/__tests__/SelfHostingComposeEnvironment.test.ts +++ b/packages/config/src/__tests__/SelfHostingComposeEnvironment.test.ts @@ -20,6 +20,7 @@ const API_SETTINGS_NOT_FORWARDED: Record = { FLUXER_TEST_MODE_ENABLED: 'test and development only', FLUXER_TEST_HARNESS_TOKEN: 'test and development only', FLUXER_VALIDATE_RESPONSES: 'test and development only', + FLUXER_API_AUTOMATED_MESSAGE_DELETION_DELAY_DAYS: 'automated account actions run only on the hosted service', ...Object.fromEntries( ['MONTHLY', 'YEARLY', 'GIFT_1_MONTH', 'GIFT_1_YEAR'].flatMap((slot) => ['USD', 'EUR', 'BRL', 'DKK', 'INR', 'NOK', 'PLN', 'SEK', 'TRY'].map((currency) => [ diff --git a/packages/config/src/config_loader/EnvironmentOverrides.ts b/packages/config/src/config_loader/EnvironmentOverrides.ts index f20039f06..393874df5 100644 --- a/packages/config/src/config_loader/EnvironmentOverrides.ts +++ b/packages/config/src/config_loader/EnvironmentOverrides.ts @@ -71,6 +71,10 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record = { FLUXER_API_HEADERS_TIMEOUT_MS: {path: ['services', 'api', 'headers_timeout_ms'], parse: parseInteger}, FLUXER_API_REQUEST_TIMEOUT_MS: {path: ['services', 'api', 'request_timeout_ms'], parse: parseInteger}, FLUXER_API_MAX_INFLIGHT_REQUESTS: {path: ['services', 'api', 'max_inflight_requests'], parse: parseInteger}, + FLUXER_API_AUTOMATED_MESSAGE_DELETION_DELAY_DAYS: { + path: ['services', 'api', 'automated_message_deletion_delay_days'], + parse: parseInteger, + }, FLUXER_API_IP_BAN_EXEMPT_IPS: {path: ['services', 'api', 'ip_ban_exempt_ips'], parse: parseCsv}, FLUXER_API_DONATION_PROXY_KEY: {path: ['services', 'api', 'donation_proxy_key']}, FLUXER_API_TRUSTED_CALLERS: {path: ['services', 'api', 'trusted_callers'], parse: parseJsonArray},