diff --git a/fluxer_api/src/api/channel/services/message/MessageMentionService.ts b/fluxer_api/src/api/channel/services/message/MessageMentionService.ts index e31d5638b..c5a077c57 100644 --- a/fluxer_api/src/api/channel/services/message/MessageMentionService.ts +++ b/fluxer_api/src/api/channel/services/message/MessageMentionService.ts @@ -22,10 +22,12 @@ import { } from '../../../BrandedTypes'; import type {IGuildRepositoryAggregate} from '../../../guild/repositories/IGuildRepositoryAggregate'; import type {GatewayChannelMention, IGatewayService} from '../../../infrastructure/IGatewayService'; +import {Logger} from '../../../Logger'; import type {Channel} from '../../../models/Channel'; import type {Message} from '../../../models/Message'; import type {IUserRepository} from '../../../user/IUserRepository'; import type {WorkerTaskName} from '../../../worker/WorkerLaneConfig'; +import {WorkerQueueOverflowError} from '../../../worker/WorkerQueueOverflowError'; import {isOperationDisabled, isPersonalNotesChannel} from './MessageHelpers'; import type {MessageResponseDataService} from './MessageResponseDataService'; @@ -498,8 +500,19 @@ export class MessageMentionService { mentionHere || mentionEveryone || (message.reference && message.type === MessageTypes.REPLY); - if (hasMentions) { + if (!hasMentions) { + return; + } + try { await this.workerService.addJob('handleMentions', taskData, {skipLedger: true}); + } catch (error) { + if (!(error instanceof WorkerQueueOverflowError)) { + throw error; + } + Logger.warn( + {channelId: message.channelId.toString(), messageId: message.id.toString()}, + 'Dropped mention fanout, jobs stream is at its limit', + ); } } } diff --git a/fluxer_api/src/api/infrastructure/EmbedService.ts b/fluxer_api/src/api/infrastructure/EmbedService.ts index 659aa4946..dd6e2c9f2 100644 --- a/fluxer_api/src/api/infrastructure/EmbedService.ts +++ b/fluxer_api/src/api/infrastructure/EmbedService.ts @@ -25,6 +25,7 @@ import {EmbedFooter} from '../models/EmbedFooter'; import {EmbedMedia} from '../models/EmbedMedia'; import * as UnfurlerUtils from '../utils/UnfurlerUtils'; import type {WorkerTaskName} from '../worker/WorkerLaneConfig'; +import {WorkerQueueOverflowError} from '../worker/WorkerQueueOverflowError'; import { type IMediaService, type MediaProxyMetadataResponse, @@ -498,20 +499,30 @@ export class EmbedService { ): Promise { const expectedContentHash = options.content !== undefined ? UnfurlerUtils.hashUnfurlContent(options.content) : undefined; - await this.workerService.addJob( - 'extractEmbeds', - { - guildId: guildId ? guildId.toString() : null, - channelId: channelId.toString(), - messageId: messageId.toString(), - nsfwMode, - ...(expectedContentHash ? {expectedContentHash} : {}), - }, - { - jobKey: expectedContentHash ? `${messageId.toString()}:${expectedContentHash}` : messageId.toString(), - skipLedger: true, - }, - ); + try { + await this.workerService.addJob( + 'extractEmbeds', + { + guildId: guildId ? guildId.toString() : null, + channelId: channelId.toString(), + messageId: messageId.toString(), + nsfwMode, + ...(expectedContentHash ? {expectedContentHash} : {}), + }, + { + jobKey: expectedContentHash ? `${messageId.toString()}:${expectedContentHash}` : messageId.toString(), + skipLedger: true, + }, + ); + } catch (error) { + if (!(error instanceof WorkerQueueOverflowError)) { + throw error; + } + Logger.warn( + {channelId: channelId.toString(), messageId: messageId.toString()}, + 'Dropped url embed extraction, jobs stream is at its limit', + ); + } } private async updateMessageEmbeds(channelId: ChannelID, messageId: MessageID, embeds: Array): Promise { diff --git a/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts b/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts index 9cd57b4d1..dca0eca02 100644 --- a/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts +++ b/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts @@ -3,9 +3,19 @@ import {randomUUID} from 'node:crypto'; import type {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectionManager'; import type {WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes'; -import {AckPolicy, nanos, RetentionPolicy, StorageType} from 'nats'; +import { + AckPolicy, + DiscardPolicy, + type JetStreamManager, + NatsError, + nanos, + RetentionPolicy, + StorageType, + type StreamConfig, +} from 'nats'; import {Logger} from '../Logger'; import type {WorkerLaneDefinition} from './WorkerLaneConfig'; +import {WorkerQueueOverflowError} from './WorkerQueueOverflowError'; const STREAM_NAME = 'JOBS'; const SUBJECT_PREFIX = 'jobs.'; @@ -14,6 +24,29 @@ const LEGACY_CONSUMER_NAME = 'workers'; const DLQ_STREAM_NAME = 'JOBS_DLQ'; const DLQ_SUBJECT_PREFIX = 'dlq.'; const DLQ_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000; +const STREAM_MAX_MSGS = 2_000_000; +const STREAM_MAX_BYTES = 8 * 1024 * 1024 * 1024; +const STREAM_MAX_MSGS_PER_SUBJECT = 250_000; +const STREAM_STORE_ERR_CODE = 10077; + +const STREAM_LIMITS = { + max_msgs: STREAM_MAX_MSGS, + max_bytes: STREAM_MAX_BYTES, + max_msgs_per_subject: STREAM_MAX_MSGS_PER_SUBJECT, + discard: DiscardPolicy.New, + discard_new_per_subject: true, +} satisfies Partial; + +function describeStreamRejection(error: unknown): string | null { + if (!(error instanceof NatsError)) { + return null; + } + const apiError = error.jsError(); + if (apiError?.err_code !== STREAM_STORE_ERR_CODE) { + return null; + } + return apiError.description ?? 'stream rejected the publish'; +} export class JetStreamWorkerQueue { private readonly connectionManager: JetStreamConnectionManager; @@ -30,9 +63,13 @@ export class JetStreamWorkerQueue { return; } const jsm = await this.connectionManager.getJetStreamManager(); + let existingConfig: StreamConfig | null = null; try { - await jsm.streams.info(STREAM_NAME); + existingConfig = (await jsm.streams.info(STREAM_NAME)).config; } catch { + existingConfig = null; + } + if (existingConfig === null) { await jsm.streams.add({ name: STREAM_NAME, subjects: [`${SUBJECT_PREFIX}>`], @@ -41,11 +78,33 @@ export class JetStreamWorkerQueue { max_age: nanos(MAX_AGE_MS), duplicate_window: nanos(2 * 60 * 1000), num_replicas: 1, + ...STREAM_LIMITS, }); + } else if (!this.hasStreamLimits(existingConfig)) { + await this.applyStreamLimits(jsm); } this.streamReady = true; } + private hasStreamLimits(config: StreamConfig): boolean { + return ( + config.max_msgs === STREAM_MAX_MSGS && + config.max_bytes === STREAM_MAX_BYTES && + config.max_msgs_per_subject === STREAM_MAX_MSGS_PER_SUBJECT && + config.discard === DiscardPolicy.New && + config.discard_new_per_subject + ); + } + + private async applyStreamLimits(jsm: JetStreamManager): Promise { + try { + await jsm.streams.update(STREAM_NAME, {...STREAM_LIMITS}); + Logger.info({stream: STREAM_NAME, ...STREAM_LIMITS}, 'Applied jobs stream limits'); + } catch (error) { + Logger.error({err: error, stream: STREAM_NAME}, 'Failed to apply jobs stream limits, stream stays unbounded'); + } + } + async ensureDlqStream(): Promise { if (this.dlqStreamReady) { return; @@ -148,11 +207,19 @@ export class JetStreamWorkerQueue { created_at: new Date().toISOString(), }); const msgID = options?.jobKey ? `${taskType}:${options.jobKey}` : randomUUID(); - const ack = await js.publish(subject, body, { - msgID, - }); - const jobId = `${ack.seq}`; - return jobId; + try { + const ack = await js.publish(subject, body, { + msgID, + }); + const jobId = `${ack.seq}`; + return jobId; + } catch (error) { + const rejection = describeStreamRejection(error); + if (rejection === null) { + throw error; + } + throw new WorkerQueueOverflowError(taskType, rejection); + } } async publishToDlq( diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index 1db068bee..0842e2bfb 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -30,6 +30,7 @@ import { validateLaneCompleteness, type WorkerLaneDefinition, } from './WorkerLaneConfig'; +import {WorkerQueueOverflowError} from './WorkerQueueOverflowError'; import {WorkerRunner} from './WorkerRunner'; import {WorkerService} from './WorkerService'; import {workerTasks} from './WorkerTaskRegistry'; @@ -219,7 +220,14 @@ export async function startWorkerMain(): Promise { ); if (didClaimEmailSync) { Logger.info('Triggering initial disposable email domain sync'); - await workerService.addJob('syncDisposableEmailDomains', {}); + try { + await workerService.addJob('syncDisposableEmailDomains', {}); + } catch (error) { + if (!(error instanceof WorkerQueueOverflowError)) { + throw error; + } + Logger.warn('Dropped initial disposable email domain sync, jobs stream is at its limit'); + } } } cron = new CronScheduler(workerService, Logger, dependencies.kvClient); diff --git a/fluxer_api/src/api/worker/WorkerQueueOverflowError.ts b/fluxer_api/src/api/worker/WorkerQueueOverflowError.ts new file mode 100644 index 000000000..d4436664b --- /dev/null +++ b/fluxer_api/src/api/worker/WorkerQueueOverflowError.ts @@ -0,0 +1,11 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +export class WorkerQueueOverflowError extends Error { + readonly taskType: string; + + constructor(taskType: string, reason: string) { + super(`Jobs stream rejected task "${taskType}": ${reason}`); + this.name = 'WorkerQueueOverflowError'; + this.taskType = taskType; + } +} diff --git a/fluxer_api/src/api/worker/WorkerService.ts b/fluxer_api/src/api/worker/WorkerService.ts index 3f2d84f83..24f1bdce8 100644 --- a/fluxer_api/src/api/worker/WorkerService.ts +++ b/fluxer_api/src/api/worker/WorkerService.ts @@ -7,6 +7,7 @@ import type {IJobLedgerRepository} from '../jobs/IJobLedgerRepository'; import {Logger} from '../Logger'; import type {JetStreamWorkerQueue} from './JetStreamWorkerQueue'; import {findLaneForTask, type WorkerTaskName} from './WorkerLaneConfig'; +import {WorkerQueueOverflowError} from './WorkerQueueOverflowError'; export class WorkerService implements IWorkerService { private readonly queue: JetStreamWorkerQueue; @@ -56,6 +57,10 @@ export class WorkerService implements IWorkerService { Logger.debug({taskType, jobId: jobId.toString(), seq}, 'Job queued successfully'); return jobId; } catch (error) { + if (error instanceof WorkerQueueOverflowError) { + Logger.warn({taskType, jobId: jobId.toString()}, 'Jobs stream is at its limit, shedding job'); + throw error; + } Logger.error({error, taskType, payload}, 'Failed to queue job'); throw error; } diff --git a/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts b/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts new file mode 100644 index 000000000..41c3be07e --- /dev/null +++ b/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts @@ -0,0 +1,145 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectionManager'; +import {DiscardPolicy, NatsError, RetentionPolicy, StorageType, type StreamConfig} from 'nats'; +import {describe, expect, it} from 'vitest'; +import {JetStreamWorkerQueue} from '../JetStreamWorkerQueue'; +import {WorkerQueueOverflowError} from '../WorkerQueueOverflowError'; + +const EXPECTED_LIMITS = { + max_msgs: 2_000_000, + max_bytes: 8 * 1024 * 1024 * 1024, + max_msgs_per_subject: 250_000, + discard: DiscardPolicy.New, + discard_new_per_subject: true, +}; + +const LEGACY_CONFIG = { + name: 'JOBS', + subjects: ['jobs.>'], + retention: RetentionPolicy.Workqueue, + storage: StorageType.File, + max_msgs: -1, + max_bytes: -1, + max_msgs_per_subject: -1, + discard: DiscardPolicy.Old, + discard_new_per_subject: false, +} as unknown as StreamConfig; + +function streamLimitError(description: string): NatsError { + const error = new NatsError('503', '503'); + error.api_error = {code: 503, err_code: 10077, description}; + return error; +} + +function boundedPublisher(maxPerSubject: number): (subject: string) => {seq: number} { + const counts = new Map(); + let seq = 0; + return (subject) => { + const stored = counts.get(subject) ?? 0; + if (stored >= maxPerSubject) { + throw streamLimitError('maximum messages per subject exceeded'); + } + counts.set(subject, stored + 1); + seq += 1; + return {seq}; + }; +} + +function createQueue(params: { + existing?: StreamConfig | null; + updateError?: Error; + publish?: (subject: string) => {seq: number}; +}): {queue: JetStreamWorkerQueue; added: Array>; updated: Array>} { + const added: Array> = []; + const updated: Array> = []; + const connectionManager = { + getJetStreamManager: () => + Promise.resolve({ + streams: { + info: () => { + if (!params.existing) { + return Promise.reject(new Error('stream not found')); + } + return Promise.resolve({config: params.existing}); + }, + add: (config: Partial) => { + added.push(config); + return Promise.resolve({config}); + }, + update: (_name: string, config: Partial) => { + if (params.updateError) { + return Promise.reject(params.updateError); + } + updated.push(config); + return Promise.resolve({config}); + }, + }, + }), + getJetStreamClient: () => ({ + publish: (subject: string) => { + const publish = params.publish ?? (() => ({seq: 1})); + return Promise.resolve(publish(subject)); + }, + }), + } as unknown as JetStreamConnectionManager; + return {queue: new JetStreamWorkerQueue(connectionManager), added, updated}; +} + +describe('jobs stream limits', () => { + it('creates the stream bounded and discarding new messages', async () => { + const {queue, added, updated} = createQueue({existing: null}); + await queue.ensureStream(); + expect(updated).toHaveLength(0); + expect(added).toHaveLength(1); + expect(added[0]).toMatchObject(EXPECTED_LIMITS); + expect(added[0]).toMatchObject({retention: RetentionPolicy.Workqueue, storage: StorageType.File}); + }); + + it('applies the limits to an existing unbounded stream', async () => { + const {queue, added, updated} = createQueue({existing: LEGACY_CONFIG}); + await queue.ensureStream(); + expect(added).toHaveLength(0); + expect(updated).toHaveLength(1); + expect(updated[0]).toEqual(EXPECTED_LIMITS); + }); + + it('leaves an already bounded stream alone', async () => { + const {queue, added, updated} = createQueue({existing: {...LEGACY_CONFIG, ...EXPECTED_LIMITS}}); + await queue.ensureStream(); + expect(added).toHaveLength(0); + expect(updated).toHaveLength(0); + }); + + it('keeps startup alive when the limit update is rejected', async () => { + const {queue} = createQueue({existing: LEGACY_CONFIG, updateError: new Error('stream update rejected')}); + await expect(queue.ensureStream()).resolves.toBeUndefined(); + }); +}); + +describe('jobs stream enqueue shedding', () => { + it('rejects enqueues once the stream is at its cap', async () => { + const {queue} = createQueue({existing: LEGACY_CONFIG, publish: boundedPublisher(2)}); + await expect(queue.enqueue('extractEmbeds', {})).resolves.toBe('1'); + await expect(queue.enqueue('extractEmbeds', {})).resolves.toBe('2'); + await expect(queue.enqueue('extractEmbeds', {})).rejects.toBeInstanceOf(WorkerQueueOverflowError); + }); + + it('caps each task type independently', async () => { + const {queue} = createQueue({existing: LEGACY_CONFIG, publish: boundedPublisher(1)}); + await expect(queue.enqueue('extractEmbeds', {})).resolves.toBe('1'); + await expect(queue.enqueue('extractEmbeds', {})).rejects.toBeInstanceOf(WorkerQueueOverflowError); + await expect(queue.enqueue('handleMentions', {})).resolves.toBe('2'); + }); + + it('rethrows publish failures that are not stream limits', async () => { + const failure = new Error('no responders'); + const {queue} = createQueue({ + existing: LEGACY_CONFIG, + publish: () => { + throw failure; + }, + }); + await expect(queue.enqueue('extractEmbeds', {})).rejects.toBe(failure); + }); +}); diff --git a/fluxer_api/src/api/worker/tests/PerMessageJobLedger.test.ts b/fluxer_api/src/api/worker/tests/PerMessageJobLedger.test.ts index 037ca498b..f59e81ca4 100644 --- a/fluxer_api/src/api/worker/tests/PerMessageJobLedger.test.ts +++ b/fluxer_api/src/api/worker/tests/PerMessageJobLedger.test.ts @@ -9,6 +9,7 @@ import {MessageMentionService} from '../../channel/services/message/MessageMenti import {EmbedService} from '../../infrastructure/EmbedService'; import type {Message} from '../../models/Message'; import type {WorkerTaskName} from '../WorkerLaneConfig'; +import {WorkerQueueOverflowError} from '../WorkerQueueOverflowError'; class RecordingWorkerService implements IWorkerService { readonly jobs: Array<{taskType: WorkerTaskName; options: WorkerJobOptions | undefined}> = []; @@ -31,6 +32,24 @@ class RecordingWorkerService implements IWorkerService { } } +class OverflowingWorkerService implements IWorkerService { + async addJob( + taskType: WorkerTaskName, + _payload: TPayload, + _options?: WorkerJobOptions, + ): Promise { + throw new WorkerQueueOverflowError(taskType, 'maximum messages per subject exceeded'); + } + + async cancelJob(): Promise { + return false; + } + + async retryDeadLetterJob(): Promise { + return false; + } +} + function makeMentionMessage(): Message { return { id: createMessageID(2n), @@ -71,4 +90,28 @@ describe('per-message worker jobs', () => { expect(workerService.jobs[0]!.taskType).toBe('extractEmbeds'); expect(workerService.jobs[0]!.options?.skipLedger).toBe(true); }); + + it('drops mention fanout instead of failing the send when the jobs stream is full', async () => { + const mentionService = new MessageMentionService( + null as never, + null as never, + null as never, + new OverflowingWorkerService(), + null as never, + ); + await expect( + mentionService.handleMentionTasks({ + guildId: null, + message: makeMentionMessage(), + authorId: createUserID(1n), + }), + ).resolves.toBeUndefined(); + }); + + it('drops embed extraction instead of failing the send when the jobs stream is full', async () => { + const embedService = new EmbedService(null as never, null as never, null as never, new OverflowingWorkerService()); + await expect( + embedService.enqueueUrlEmbedExtraction(createChannelID(3n), createMessageID(2n), null, 'block'), + ).resolves.toBeUndefined(); + }); });