feat(api): archive and schedule automated message deletion (#3231)

This commit is contained in:
Hampus
2026-10-05 21:03:29 +02:00
committed by GitHub
parent f4e545e090
commit e26c8c870d
13 changed files with 365 additions and 120 deletions
+1
View File
@@ -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(',') ?? '',
+1
View File
@@ -54,6 +54,7 @@ export interface APIConfig {
headersTimeoutMs: number;
requestTimeoutMs: number;
maxInflightRequests: number;
automatedMessageDeletionDelayDays: number;
ipBanExemptIps: Array<string>;
cassandra: {
hosts: string;
@@ -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<T extends ActionEnvelope['type']> = Extract<ActionEnvelope, {type: T}>;
@@ -49,7 +51,9 @@ export interface AccountStateDeps {
dispatch: AccountUpdateDispatch;
ipBans: Pick<AdminRepository, 'isIpBanned' | 'banIpTemp'>;
cache: Pick<ICacheService, 'publish' | 'get' | 'set' | 'delete'>;
messages: Pick<AdminMessageDeletionService, 'deleteAllUserMessages'>;
archives: Pick<AdminArchiveService, 'triggerUserArchive' | 'listArchives'>;
messageDeletionQueue: Pick<KVBulkMessageDeletionQueueService, 'scheduleDeletion' | 'removeFromQueue'>;
messageDeletionDelayMs: number;
authored: Pick<IChannelRepository, 'listMessagesByAuthor'>;
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<PartialUserChangePropagationDeps, 'gatewayService'>;
@@ -94,7 +100,7 @@ export function accountStateDepsFromContext(
ctx: ApiContext,
ipBans: AccountStateDeps['ipBans'],
profile: Omit<ProfilePropagation, 'userRepository'>,
messages: AccountStateDeps['messages'],
messageDeletion: Pick<AccountStateDeps, 'archives' | 'messageDeletionQueue' | 'messageDeletionDelayMs'>,
channels: Pick<IChannelRepository, 'findUnique' | 'listMessagesByAuthor'>,
): 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<ActionOutcome> {
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<string> {
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<ActionOutcome> {
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<typeof parseIpAddress> {
@@ -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<void> {
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,
{
@@ -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;
@@ -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<void>(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<Array<MessageResponse>>(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);
});
@@ -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<Me>(harness, account.token).get('/users/@me').expect(HTTP_STATUS.OK).execute();
const setFlags = (flags: bigint) =>
createBuilder<unknown>(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<void>(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);
});
+7 -10
View File
@@ -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<void> {
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<void> {
userCacheService: dependencies.userCacheService,
guildRepository: dependencies.guildRepository,
},
messagePurge,
{
archives: getAdminArchiveService(),
messageDeletionQueue: dependencies.bulkMessageDeletionQueueService,
messageDeletionDelayMs: Config.automatedMessageDeletionDelayDays * ms('1 day'),
},
dependencies.channelRepository,
),
});
@@ -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<Extract<ActionEnvelope, {type: 'set_account_limit'}>> = {}) {
return envelope<'set_account_limit'>({
@@ -108,8 +135,9 @@ interface Harness {
authored: Array<{channelId: ChannelID; messageId: MessageID}>;
visibility: Array<UserID>;
removed: Array<{channelId: ChannelID; authorId: UserID; messageIds: Array<MessageID>}>;
audits: Array<{adminUserId: UserID; action: string; targetId: bigint; auditLogReason: string | null}>;
shreds: Array<{userId: bigint; entries: number; adminUserId: UserID; auditLogReason: string | null}>;
archives: Array<AdminArchiveResponse>;
archiveRequests: Array<{userId: UserID; requestedBy: UserID; includeAttachments: boolean}>;
queued: Map<string, number>;
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<unknown>},
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<typeof AdminMessageDeletionService>[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<Partial<UserRow>>) {
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<UserRow>
>) {
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<Partial<UserRow>>) {
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 () => {
+7
View File
@@ -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);
+1
View File
@@ -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<string>;
donation_proxy_key: string;
trusted_callers: Array<{
@@ -20,6 +20,7 @@ const API_SETTINGS_NOT_FORWARDED: Record<string, string> = {
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) => [
@@ -71,6 +71,10 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record<string, NamedEnvOverride> = {
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},