From dc32a7c70ecfce148cc311ca2bbf45c40b486bae Mon Sep 17 00:00:00 2001 From: Hampus Date: Wed, 30 Sep 2026 21:29:05 +0200 Subject: [PATCH] feat(admin): allow system DMs to all users (#3073) --- fluxer_admin/openapi-admin.json | 15 ++-- fluxer_admin/src/api/system_dm.rs | 9 ++- fluxer_admin/src/api/types/system_dm.rs | 2 +- fluxer_admin/src/routes/message_actions.rs | 3 +- fluxer_admin/src/templates/pages/system_dm.rs | 2 +- fluxer_api/src/api/admin/AdminService.ts | 14 ++-- .../controllers/SystemDmAdminController.ts | 9 ++- .../src/api/worker/tasks/SendSystemDm.ts | 72 +++++++++++++++---- .../worker/tests/SendSystemDmCancel.test.ts | 62 +++++++++++++++- .../src/content/docs/admin-api/system-dms.mdx | 13 ++-- .../schema/src/domains/admin/AdminSchemas.ts | 30 +++++--- 11 files changed, 183 insertions(+), 48 deletions(-) diff --git a/fluxer_admin/openapi-admin.json b/fluxer_admin/openapi-admin.json index 69c5876ec..e0ecd4ebb 100644 --- a/fluxer_admin/openapi-admin.json +++ b/fluxer_admin/openapi-admin.json @@ -5435,7 +5435,7 @@ "content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}} } }, - "description": "Queue a worker job that delivers the same content to every listed user as a direct message from the system account. Progress is observable through the Jobs admin resource (task_type=sendSystemDm), and an in-flight broadcast is stopped by cancelling that job. Requires SYSTEM_DM_SEND permission.", + "description": "Queue a worker job that delivers the same content to every listed user, or to every user when all_users is set, as a direct message from the system account. Progress is observable through the Jobs admin resource (task_type=sendSystemDm), and an in-flight broadcast is stopped by cancelling that job. Requires SYSTEM_DM_SEND permission.", "security": [{"adminApiKey": []}], "requestBody": { "required": true, @@ -10150,20 +10150,25 @@ "description": "Message content to send to each recipient" }, "user_ids": { + "description": "Recipient user IDs. Each receives the same content as a system DM.", "minItems": 1, "maxItems": 10000, "type": "array", - "items": {"$ref": "#/components/schemas/SnowflakeType"}, - "description": "Recipient user IDs. Each receives the same content as a system DM." + "items": {"$ref": "#/components/schemas/SnowflakeType"} + }, + "all_users": { + "description": "Send to every user account, skipping bots, system accounts, and deleted or disabled accounts", + "type": "boolean" } }, - "required": ["content", "user_ids"] + "required": ["content"] }, "SendSystemDmResponse": { "type": "object", "properties": { "recipient_count": { - "description": "Number of recipients the worker job was queued to deliver to", + "nullable": true, + "description": "Number of recipients the worker job was queued to deliver to, or null when sending to all users", "allOf": [{"$ref": "#/components/schemas/Int32Type"}] } }, diff --git a/fluxer_admin/src/api/system_dm.rs b/fluxer_admin/src/api/system_dm.rs index 7e82d2ba7..1b2d8a296 100644 --- a/fluxer_admin/src/api/system_dm.rs +++ b/fluxer_admin/src/api/system_dm.rs @@ -8,13 +8,18 @@ use super::types::SendSystemDmResponse; impl AdminApiClient { pub async fn send_system_dm( &self, - user_ids: &[String], + user_ids: Option<&[String]>, content: &str, ) -> ApiResult { let body = generated_types::SendSystemDmRequest { content: generated_types::SendSystemDmRequestContent::try_from(content) .map_err(|e| ApiError::Parse(e.to_string()))?, - user_ids: user_ids.iter().map(|id| snowflake(id)).collect(), + user_ids: user_ids + .unwrap_or_default() + .iter() + .map(|id| snowflake(id)) + .collect(), + all_users: user_ids.is_none().then_some(true), }; let response = self .generated() diff --git a/fluxer_admin/src/api/types/system_dm.rs b/fluxer_admin/src/api/types/system_dm.rs index 96fa1dad1..08b90eb08 100644 --- a/fluxer_admin/src/api/types/system_dm.rs +++ b/fluxer_admin/src/api/types/system_dm.rs @@ -4,5 +4,5 @@ use serde::{Deserialize, Serialize}; #[derive(Clone, Debug, Deserialize, Serialize)] pub struct SendSystemDmResponse { - pub recipient_count: i64, + pub recipient_count: Option, } diff --git a/fluxer_admin/src/routes/message_actions.rs b/fluxer_admin/src/routes/message_actions.rs index eb3827238..3fe3303f0 100644 --- a/fluxer_admin/src/routes/message_actions.rs +++ b/fluxer_admin/src/routes/message_actions.rs @@ -171,7 +171,8 @@ pub(crate) async fn system_dms_post( let flash = if let Some(content) = content.as_deref() && !user_ids.is_empty() { - match client.send_system_dm(&user_ids, content).await { + let recipients = (user_ids != ["*"]).then_some(user_ids.as_slice()); + match client.send_system_dm(recipients, content).await { Ok(_) => FlashData::success("System DM sent"), Err(error) => { tracing::warn!(%error, "admin API request failed: send system DM"); diff --git a/fluxer_admin/src/templates/pages/system_dm.rs b/fluxer_admin/src/templates/pages/system_dm.rs index e63f15771..39dec86c1 100644 --- a/fluxer_admin/src/templates/pages/system_dm.rs +++ b/fluxer_admin/src/templates/pages/system_dm.rs @@ -58,7 +58,7 @@ pub fn system_dm_page( (form_field_group( "Recipient user IDs", "system-dm-user-ids", true, None, - Some("One per line. Snowflake IDs only."), + Some("One per line. Snowflake IDs only, or a single * to send to every user."), html! { textarea id="system-dm-user-ids" name="user_ids" required rows="10" diff --git a/fluxer_api/src/api/admin/AdminService.ts b/fluxer_api/src/api/admin/AdminService.ts index 1f339e494..08bb0e992 100644 --- a/fluxer_api/src/api/admin/AdminService.ts +++ b/fluxer_api/src/api/admin/AdminService.ts @@ -178,20 +178,20 @@ export class AdminService { } async sendSystemDm( - data: {content: string; userIds: Array}, + data: {content: string; recipients: {kind: 'all'} | {kind: 'list'; userIds: Array}}, adminUserId: UserID, auditLogReason: string | null, ): Promise { + const recipientCount = data.recipients.kind === 'all' ? null : data.recipients.userIds.length; await this.apiContext.services.worker.addJob( 'sendSystemDm', - { - content: data.content, - user_ids: data.userIds, - }, + data.recipients.kind === 'all' + ? {content: data.content, all_users: true} + : {content: data.content, user_ids: data.recipients.userIds}, {requireLedger: true}, ); const metadata = new Map([ - ['recipient_count', data.userIds.length.toString()], + ['recipient_count', recipientCount === null ? 'all' : recipientCount.toString()], ['content_length', data.content.length.toString()], ]); await this.auditService.createAuditLog({ @@ -202,6 +202,6 @@ export class AdminService { auditLogReason, metadata, }); - return {recipient_count: data.userIds.length}; + return {recipient_count: recipientCount}; } } diff --git a/fluxer_api/src/api/admin/controllers/SystemDmAdminController.ts b/fluxer_api/src/api/admin/controllers/SystemDmAdminController.ts index f3b1aa701..28ac4b936 100644 --- a/fluxer_api/src/api/admin/controllers/SystemDmAdminController.ts +++ b/fluxer_api/src/api/admin/controllers/SystemDmAdminController.ts @@ -23,7 +23,7 @@ export function SystemDmAdminController(app: HonoApp) { security: 'adminApiKey', tags: 'Admin', description: - 'Queue a worker job that delivers the same content to every listed user as a direct message from the system account. Progress is observable through the Jobs admin resource (task_type=sendSystemDm), and an in-flight broadcast is stopped by cancelling that job. Requires SYSTEM_DM_SEND permission.', + 'Queue a worker job that delivers the same content to every listed user, or to every user when all_users is set, as a direct message from the system account. Progress is observable through the Jobs admin resource (task_type=sendSystemDm), and an in-flight broadcast is stopped by cancelling that job. Requires SYSTEM_DM_SEND permission.', }), async (ctx) => { const adminService = ctx.get('adminService'); @@ -31,7 +31,12 @@ export function SystemDmAdminController(app: HonoApp) { const auditLogReason = ctx.get('auditLogReason'); const payload = ctx.req.valid('json'); const result = await adminService.sendSystemDm( - {content: payload.content, userIds: payload.user_ids.map((id) => id.toString())}, + { + content: payload.content, + recipients: payload.all_users + ? {kind: 'all'} + : {kind: 'list', userIds: (payload.user_ids ?? []).map((id) => id.toString())}, + }, adminUserId, auditLogReason, ); diff --git a/fluxer_api/src/api/worker/tasks/SendSystemDm.ts b/fluxer_api/src/api/worker/tasks/SendSystemDm.ts index 85370a295..295dad199 100644 --- a/fluxer_api/src/api/worker/tasks/SendSystemDm.ts +++ b/fluxer_api/src/api/worker/tasks/SendSystemDm.ts @@ -2,19 +2,65 @@ import {createUserID, type UserID} from '@app/api/BrandedTypes'; import {createRequestCache} from '@app/api/middleware/RequestCacheMiddleware'; +import type {User} from '@app/api/models/User'; import {UserChannelService} from '@app/api/user/services/UserChannelService'; import {getWorkerDependencies} from '@app/api/worker/WorkerContext'; +import {UserFlags} from '@fluxer/constants/src/UserConstants'; import {JobCancelledError, type WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask'; import {z} from 'zod'; const SYSTEM_USER_ID: UserID = createUserID(0n); -const PayloadSchema = z.object({ - content: z.string().min(1).max(4000), - user_ids: z.array(z.string().regex(/^\d+$/)).min(1), -}); +const ALL_USERS_PAGE_SIZE = 100; +const CURSOR_TTL_SECONDS = 7 * 24 * 60 * 60; +const INELIGIBLE_FLAGS = UserFlags.DELETED | UserFlags.SELF_DELETED | UserFlags.DISABLED; +const PayloadSchema = z.union([ + z.object({ + content: z.string().min(1).max(4000), + user_ids: z.array(z.string().regex(/^\d+$/)).min(1), + }), + z.object({ + content: z.string().min(1).max(4000), + all_users: z.literal(true), + }), +]); + +function isEligibleRecipient(user: User): boolean { + return user.id !== SYSTEM_USER_ID && !user.isBot && !user.isSystem && (user.flags & INELIGIBLE_FLAGS) === 0n; +} + +async function* allUserRecipients(helpers: WorkerTaskHelpers): AsyncGenerator { + const {userRepository, kvClient} = getWorkerDependencies(); + const cursorKey = `system_dm:all_users_cursor:${helpers.jobId}`; + let pageState = await kvClient.get(cursorKey); + if (pageState !== null) { + helpers.logger.info('Resuming system DM broadcast from saved cursor'); + } + do { + const page = await userRepository.scanAllUsersPage(ALL_USERS_PAGE_SIZE, pageState); + for (const user of page.users) { + if (isEligibleRecipient(user)) { + yield user.id; + } + } + pageState = page.pageState; + if (pageState !== null) { + await kvClient.setex(cursorKey, CURSOR_TTL_SECONDS, pageState); + } + } while (pageState !== null); + await kvClient.del(cursorKey); +} + +async function* listedRecipients(userIds: Array): AsyncGenerator { + for (const raw of userIds) { + yield createUserID(BigInt(raw)); + } +} export async function sendSystemDm(payload: unknown, helpers: WorkerTaskHelpers): Promise { - const {content, user_ids} = PayloadSchema.parse(payload); + const parsed = PayloadSchema.parse(payload); + const {content} = parsed; + const total = 'user_ids' in parsed ? parsed.user_ids.length : null; + const recipients = 'user_ids' in parsed ? listedRecipients(parsed.user_ids) : allUserRecipients(helpers); const deps = getWorkerDependencies(); const systemUser = await deps.userRepository.findUniqueAssert(SYSTEM_USER_ID); const userChannelService = new UserChannelService( @@ -29,16 +75,12 @@ export async function sendSystemDm(payload: unknown, helpers: WorkerTaskHelpers) const requestCache = createRequestCache(); let sent = 0; let failed = 0; - for (const raw of user_ids) { + for await (const recipientId of recipients) { if (await helpers.shouldCancel()) { - helpers.logger.info( - {sent, failed, remaining: user_ids.length - sent - failed}, - 'System DM job cancelled mid-flight', - ); + helpers.logger.info({sent, failed, total}, 'System DM job cancelled mid-flight'); requestCache.clear(); throw new JobCancelledError(); } - const recipientId = createUserID(BigInt(raw)); try { const channel = await userChannelService.ensureDmOpenForBothUsers({ userId: SYSTEM_USER_ID, @@ -55,9 +97,13 @@ export async function sendSystemDm(payload: unknown, helpers: WorkerTaskHelpers) sent += 1; } catch (error) { failed += 1; - helpers.logger.warn({recipientId: raw, error}, 'System DM send failed for recipient'); + helpers.logger.warn({recipientId: recipientId.toString(), error}, 'System DM send failed for recipient'); + } + if ((sent + failed) % ALL_USERS_PAGE_SIZE === 0) { + requestCache.clear(); + await helpers.reportProgress(sent + failed, total, `${sent} sent, ${failed} failed`); } } requestCache.clear(); - helpers.logger.info({sent, failed, total: user_ids.length}, 'System DM job complete'); + helpers.logger.info({sent, failed, total: sent + failed}, 'System DM job complete'); } diff --git a/fluxer_api/src/api/worker/tests/SendSystemDmCancel.test.ts b/fluxer_api/src/api/worker/tests/SendSystemDmCancel.test.ts index b103860e4..c11419d2f 100644 --- a/fluxer_api/src/api/worker/tests/SendSystemDmCancel.test.ts +++ b/fluxer_api/src/api/worker/tests/SendSystemDmCancel.test.ts @@ -8,6 +8,7 @@ import type {UserRepository} from '@app/api/user/repositories/UserRepository'; import {sendSystemDm} from '@app/api/worker/tasks/SendSystemDm'; import {clearWorkerDependencies, setWorkerDependenciesForTest} from '@app/api/worker/WorkerContext'; import {WorkerRunner} from '@app/api/worker/WorkerRunner'; +import {UserFlags} from '@fluxer/constants/src/UserConstants'; import type {JsMsg} from '@nats-io/jetstream'; import {afterEach, beforeAll, describe, expect, it, vi} from 'vitest'; @@ -71,11 +72,11 @@ function createWorkerDependencies() { return {sentChannelIds, sentUserIds}; } -function createJobMessage() { +function createJobMessage(recipients: Record = {user_ids: ['11', '12', '13']}) { const envelope = { payload: { content: 'scheduled maintenance tonight', - user_ids: ['11', '12', '13'], + ...recipients, __jobId: LEDGER_JOB_ID.toString(), }, max_attempts: 5, @@ -164,4 +165,61 @@ describe('System DM cancellation', () => { expect(deps.sentUserIds).toEqual([0n, 0n, 0n]); }); + + it('broadcasts to every eligible user when all_users is set', async () => { + const user = (id: bigint, extra: Record = {}) => ({ + id, + isBot: false, + isSystem: false, + flags: 0n, + ...extra, + }); + const pages = [ + {users: [user(0n, {isSystem: true}), user(21n), user(22n, {isBot: true})], pageState: 'page-2'}, + { + users: [user(23n, {flags: UserFlags.DELETED}), user(24n), user(25n, {flags: UserFlags.DISABLED})], + pageState: null, + }, + ]; + const kv = new Map(); + const kvClient = { + get: async (key: string) => kv.get(key) ?? null, + setex: async (key: string, _ttl: number, value: string) => { + kv.set(key, value); + }, + del: async (key: string) => (kv.delete(key) ? 1 : 0), + }; + const recipientIds: Array = []; + const systemUser = {id: 0n, username: 'Fluxer', bot: true, system: true}; + const userRepository = { + findUnique: async () => systemUser, + findUniqueAssert: async () => systemUser, + findExistingDmState: async (_userId: bigint, recipientId: bigint) => { + recipientIds.push(recipientId); + return {id: 500n}; + }, + isDmChannelOpen: async () => true, + scanAllUsersPage: async (_limit: number, pageState: string | null) => + pageState === 'page-2' ? pages[1] : pages[0], + } as unknown as UserRepository; + const channelService = { + messages: {send: {sendMessage: async () => {}}}, + } as unknown as ChannelService; + setWorkerDependenciesForTest({userRepository, channelService, kvClient} as never); + const {ledger, markSucceeded} = createLedgerStub(Number.POSITIVE_INFINITY); + const runner = new TestWorkerRunner({ + tasks: {[TASK_TYPE]: sendSystemDm}, + queue: queueStub, + consumerName: 'workers_batch', + laneName: 'batch', + ledger, + concurrency: 1, + }); + + await expect(runner.runJob(TASK_TYPE, createJobMessage({all_users: true}) as unknown as JsMsg)).resolves.toBe(true); + + expect(recipientIds).toEqual([21n, 24n]); + expect(markSucceeded).toHaveBeenCalledTimes(1); + expect(kv.size).toBe(0); + }); }); diff --git a/fluxer_docs/src/content/docs/admin-api/system-dms.mdx b/fluxer_docs/src/content/docs/admin-api/system-dms.mdx index 52109c2e1..fcdd576e3 100644 --- a/fluxer_docs/src/content/docs/admin-api/system-dms.mdx +++ b/fluxer_docs/src/content/docs/admin-api/system-dms.mdx @@ -6,7 +6,7 @@ description: The system direct message broadcast and the delivery job it queues. import RouteHeader from '@/components/RouteHeader.astro'; -A system DM broadcast delivers one identical message, authored by the system account, to each supplied recipient. A recipient receives a normal message in a normal direct message channel. +A system DM broadcast delivers one identical message, authored by the system account, to each supplied recipient or to every user. A recipient receives a normal message in a normal direct message channel. The [Jobs](/admin-api/jobs/) resource reports delivery progress. [Cancelling that job](/admin-api/jobs/#cancel-job) stops an in-flight broadcast. @@ -14,24 +14,27 @@ The [Jobs](/admin-api/jobs/) resource reports delivery progress. [Cancelling tha -Queues one message for delivery to every supplied recipient. Requires `system_dm:send`. Returns the number of recipients queued. +Queues one message for delivery to every supplied recipient, or to every user when `all_users` is set. Requires `system_dm:send`. Returns the number of recipients queued. ### JSON body | Field | Type | Description | | --- | --- | --- | | content1 | string | The message delivered to every recipient (1-4000 characters) | -| user_ids2 | array[snowflake] | The recipients of the broadcast (1-10,000 entries) | +| user_ids?2 | array[snowflake] | The recipients of the broadcast (1-10,000 entries) | +| all_users?3 | boolean | Whether to deliver to every user instead of `user_ids` | 1 Content with no visible character is accepted and then fails for every recipient 2 Fluxer keeps duplicate IDs, so a repeated recipient receives the message once per occurrence. An ID naming no account is accepted here and skipped at delivery time +3 Supply exactly one of `user_ids` or `all_users: true`. A broadcast to every user skips bots, the system account, and accounts that are deleted, self-deleted, or disabled + ### Response body | Field | Type | Description | | --- | --- | --- | -| recipient_count | integer | The number of entries in the submitted `user_ids` array | +| recipient_count | ?integer | The number of entries in the submitted `user_ids` array, or null for a broadcast to every user | ### Response @@ -56,7 +59,7 @@ Recipients need not have interacted with the system account before. Fluxer skips Cancellation stops remaining deliveries and leaves sent messages in place. The [job](/admin-api/jobs/#admin-job-object) then reports `cancelled`. -The operation records one [Admin audit entry](/admin-api/#admin-audit-entry-object) with the action `system_dm.send`, the target type `system_dm`, and the target ID `0`. Its metadata has `recipient_count` and `content_length` as decimal strings, and the content itself is not recorded. +The operation records one [Admin audit entry](/admin-api/#admin-audit-entry-object) with the action `system_dm.send`, the target type `system_dm`, and the target ID `0`. Its metadata has `recipient_count` and `content_length` as decimal strings, with `recipient_count` set to `all` for a broadcast to every user. The content itself is not recorded. ### Rate limit diff --git a/packages/schema/src/domains/admin/AdminSchemas.ts b/packages/schema/src/domains/admin/AdminSchemas.ts index ac17f4269..8ac000441 100644 --- a/packages/schema/src/domains/admin/AdminSchemas.ts +++ b/packages/schema/src/domains/admin/AdminSchemas.ts @@ -861,19 +861,31 @@ export const LimitConfigUpdateRequest = z.object({ export type LimitConfigUpdateRequest = z.infer; -export const SendSystemDmRequest = z.object({ - content: z.string().min(1).max(4000).describe('Message content to send to each recipient'), - user_ids: z - .array(SnowflakeType) - .min(1) - .max(10000) - .describe('Recipient user IDs. Each receives the same content as a system DM.'), -}); +export const SendSystemDmRequest = z + .object({ + content: z.string().min(1).max(4000).describe('Message content to send to each recipient'), + user_ids: z + .array(SnowflakeType) + .min(1) + .max(10000) + .optional() + .describe('Recipient user IDs. Each receives the same content as a system DM.'), + all_users: z + .boolean() + .optional() + .describe('Send to every user account, skipping bots, system accounts, and deleted or disabled accounts'), + }) + .refine((value) => (value.all_users === true) !== (value.user_ids !== undefined), { + error: 'Provide either user_ids or all_users, not both', + path: ['user_ids'], + }); export type SendSystemDmRequest = z.infer; export const SendSystemDmResponse = z.object({ - recipient_count: Int32Type.describe('Number of recipients the worker job was queued to deliver to'), + recipient_count: Int32Type.nullable().describe( + 'Number of recipients the worker job was queued to deliver to, or null when sending to all users', + ), }); export type SendSystemDmResponse = z.infer;