From aa267b54ecbce395dd74cdad1c83ea528bed4e81 Mon Sep 17 00:00:00 2001 From: Hampus Date: Tue, 1 Sep 2026 20:47:19 +0200 Subject: [PATCH] fix(api): dead-letter retired worker task types (#2323) --- .../src/api/worker/JetStreamWorkerQueue.ts | 2 +- fluxer_api/src/api/worker/WorkerLaneConfig.ts | 14 ++ fluxer_api/src/api/worker/WorkerMain.ts | 1 + fluxer_api/src/api/worker/WorkerRunner.ts | 49 +++++++ .../worker/tests/JetStreamWorkerQueue.test.ts | 48 ++++++- .../api/worker/tests/WorkerLaneConfig.test.ts | 34 ++++- .../worker/tests/WorkerRetiredTask.test.ts | 132 ++++++++++++++++++ fluxer_messages/src/shard_impl.rs | 62 +++++++- 8 files changed, 337 insertions(+), 5 deletions(-) create mode 100644 fluxer_api/src/api/worker/tests/WorkerRetiredTask.test.ts diff --git a/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts b/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts index dca0eca02..45816b86f 100644 --- a/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts +++ b/fluxer_api/src/api/worker/JetStreamWorkerQueue.ts @@ -131,7 +131,7 @@ export class JetStreamWorkerQueue { } const jsm = await this.connectionManager.getJetStreamManager(); for (const lane of lanes) { - const filterSubjects = lane.taskTypes.map((t) => `${SUBJECT_PREFIX}${t}`); + const filterSubjects = [...lane.taskTypes, ...lane.retiredTaskTypes].map((t) => `${SUBJECT_PREFIX}${t}`); const config = { durable_name: lane.consumerName, ack_policy: AckPolicy.Explicit, diff --git a/fluxer_api/src/api/worker/WorkerLaneConfig.ts b/fluxer_api/src/api/worker/WorkerLaneConfig.ts index 9fc4f6f0c..ef99ca715 100644 --- a/fluxer_api/src/api/worker/WorkerLaneConfig.ts +++ b/fluxer_api/src/api/worker/WorkerLaneConfig.ts @@ -5,6 +5,7 @@ import type {APIWorkerLaneName, APIWorkerMode} from '../config/APIConfig'; interface LaneSettings { readonly consumerName: string; readonly tasks: ReadonlyArray; + readonly retiredTasks: ReadonlyArray; readonly concurrency: number; readonly maxAckPending: number; readonly ackWaitMs: number; @@ -15,6 +16,7 @@ const LANE_CONFIG = { realtime: { consumerName: 'workers_realtime', tasks: ['handleMentions', 'handleMentionChunk'] as const, + retiredTasks: [], concurrency: 10, maxAckPending: 50, ackWaitMs: 15000, @@ -23,6 +25,7 @@ const LANE_CONFIG = { unfurl: { consumerName: 'workers_unfurl', tasks: ['extractEmbeds'] as const, + retiredTasks: [], concurrency: 20, maxAckPending: 200, ackWaitMs: 30000, @@ -54,6 +57,7 @@ const LANE_CONFIG = { 'bulkAddGuildMembers', 'bulkBanFileShas', ] as const, + retiredTasks: ['sendScheduledMessage'], concurrency: 8, maxAckPending: 50, ackWaitMs: 60000, @@ -79,6 +83,7 @@ const LANE_CONFIG = { 'syncFileShaBlocklists', 'flushUserActivityBuffer', ] as const, + retiredTasks: [], concurrency: 12, maxAckPending: 100, ackWaitMs: 120000, @@ -92,6 +97,7 @@ interface WorkerLaneDefinition { name: APIWorkerLaneName; consumerName: string; taskTypes: ReadonlyArray; + retiredTaskTypes: ReadonlyArray; concurrency: number; maxAckPending: number; ackWaitMs: number; @@ -111,6 +117,7 @@ function makeLane(name: APIWorkerLaneName): WorkerLaneDefinition { name, consumerName: config.consumerName, taskTypes: config.tasks, + retiredTaskTypes: config.retiredTasks, concurrency: config.concurrency, maxAckPending: config.maxAckPending, ackWaitMs: config.ackWaitMs, @@ -160,6 +167,7 @@ function resolveSingleTaskLane(taskName: WorkerTaskName | undefined): WorkerLane name: parentLane.name, consumerName: `worker_${taskName}`, taskTypes: [taskName], + retiredTaskTypes: [], concurrency: parentLane.concurrency, maxAckPending: parentLane.maxAckPending, ackWaitMs: parentLane.ackWaitMs, @@ -231,6 +239,12 @@ function validateLaneCompleteness(registeredTasks: Record): voi if (missingFromRegistry.length > 0) { errors.push(`Lane tasks not found in registry: ${missingFromRegistry.join(', ')}`); } + const retiredCollisions = WORKER_LANES.flatMap((lane) => [...lane.retiredTaskTypes]).filter((task) => + registeredTaskNames.has(task), + ); + if (retiredCollisions.length > 0) { + errors.push(`Retired tasks registered again: ${retiredCollisions.join(', ')}`); + } if (errors.length > 0) { throw new Error(`Worker lane configuration mismatch:\n${errors.join('\n')}`); } diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index ebe015d97..2173623e3 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -240,6 +240,7 @@ export async function startWorkerMain(): Promise { } const runner = new WorkerRunner({ tasks: laneTasks, + retiredTaskTypes: lane.retiredTaskTypes, queue, consumerName: lane.consumerName, laneName: lane.name, diff --git a/fluxer_api/src/api/worker/WorkerRunner.ts b/fluxer_api/src/api/worker/WorkerRunner.ts index c920db248..a0e25addf 100644 --- a/fluxer_api/src/api/worker/WorkerRunner.ts +++ b/fluxer_api/src/api/worker/WorkerRunner.ts @@ -12,6 +12,7 @@ import {isJsonRecord, parseJsonRecord} from '../utils/JsonBoundaryUtils'; const MAX_DLQ_PUBLISH_ATTEMPTS = 3; const MIN_ACK_HEARTBEAT_MS = 1000; const RESUBSCRIBE_DELAY_MS = 5000; +const RETIRED_TASK_REASON = 'task type retired'; interface WorkerRunnerJetStreamClient { consumers: { @@ -44,6 +45,7 @@ interface WorkerRunnerQueue { interface WorkerRunnerOptions { tasks: Record; + retiredTaskTypes?: ReadonlyArray; queue: WorkerRunnerQueue; consumerName: string; laneName: string; @@ -56,6 +58,7 @@ interface WorkerRunnerOptions { export class WorkerRunner { private readonly tasks: Record; + private readonly retiredTaskTypes: Set; private readonly queue: WorkerRunnerQueue; private readonly consumerName: string; private readonly laneName: string; @@ -72,6 +75,7 @@ export class WorkerRunner { constructor(options: WorkerRunnerOptions) { this.tasks = options.tasks; + this.retiredTaskTypes = new Set(options.retiredTaskTypes ?? []); this.queue = options.queue; this.consumerName = options.consumerName; this.laneName = options.laneName; @@ -198,6 +202,10 @@ export class WorkerRunner { } protected async processJob(taskType: string, msg: JsMsg): Promise { + if (this.retiredTaskTypes.has(taskType)) { + await this.retireJob(taskType, msg); + return false; + } const task = this.tasks[taskType]; if (!task) { Logger.error({taskType, seq: msg.seq}, 'Unknown task type, terminating message'); @@ -364,6 +372,47 @@ export class WorkerRunner { } } + private async retireJob(taskType: string, msg: JsMsg): Promise { + const decoded = parseJsonRecord(new TextDecoder().decode(msg.data)); + const jobPayload = decoded && isJsonRecord(decoded.payload) ? decoded.payload : {}; + const runAt = decoded && typeof decoded.run_at === 'string' ? decoded.run_at : undefined; + let ledgerJobId: bigint | null = null; + const embedded = jobPayload['__jobId']; + if (typeof embedded === 'string') { + try { + ledgerJobId = BigInt(embedded); + } catch { + ledgerJobId = null; + } + delete jobPayload['__jobId']; + } + Logger.warn( + {taskType, seq: msg.seq, jobId: ledgerJobId?.toString()}, + 'Retired task type from an older release, moving to dead-letter queue', + ); + try { + await this.queue.publishToDlq(taskType, jobPayload, { + originalSeq: msg.seq, + errorMessage: RETIRED_TASK_REASON, + deliveryCount: msg.info.deliveryCount, + lane: this.laneName, + runAt, + }); + } catch (error) { + Logger.error({taskType, seq: msg.seq, err: error}, 'Failed to dead-letter a retired job'); + msg.nak(5000); + return; + } + if (ledgerJobId !== null) { + try { + await this.ledger.markDeadletter(ledgerJobId, RETIRED_TASK_REASON); + } catch (err) { + Logger.warn({err, jobId: ledgerJobId.toString()}, 'Ledger markDeadletter failed'); + } + } + msg.term(RETIRED_TASK_REASON); + } + private startAckHeartbeat(taskType: string, msg: JsMsg): ReturnType { const heartbeat = setInterval( () => { diff --git a/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts b/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts index 41c3be07e..a010eef3a 100644 --- a/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts +++ b/fluxer_api/src/api/worker/tests/JetStreamWorkerQueue.test.ts @@ -4,6 +4,7 @@ import type {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectio import {DiscardPolicy, NatsError, RetentionPolicy, StorageType, type StreamConfig} from 'nats'; import {describe, expect, it} from 'vitest'; import {JetStreamWorkerQueue} from '../JetStreamWorkerQueue'; +import {WORKER_LANES} from '../WorkerLaneConfig'; import {WorkerQueueOverflowError} from '../WorkerQueueOverflowError'; const EXPECTED_LIMITS = { @@ -26,6 +27,11 @@ const LEGACY_CONFIG = { discard_new_per_subject: false, } as unknown as StreamConfig; +interface ConsumerAddConfig { + durable_name: string; + filter_subjects: Array; +} + function streamLimitError(description: string): NatsError { const error = new NatsError('503', '503'); error.api_error = {code: 503, err_code: 10077, description}; @@ -50,12 +56,26 @@ function createQueue(params: { existing?: StreamConfig | null; updateError?: Error; publish?: (subject: string) => {seq: number}; -}): {queue: JetStreamWorkerQueue; added: Array>; updated: Array>} { +}): { + queue: JetStreamWorkerQueue; + added: Array>; + updated: Array>; + consumerAdds: Array; +} { const added: Array> = []; const updated: Array> = []; + const consumerAdds: Array = []; const connectionManager = { getJetStreamManager: () => Promise.resolve({ + consumers: { + add: (_stream: string, config: ConsumerAddConfig) => { + consumerAdds.push(config); + return Promise.resolve({}); + }, + delete: () => Promise.resolve(true), + info: () => Promise.reject(new Error('consumer not found')), + }, streams: { info: () => { if (!params.existing) { @@ -83,7 +103,7 @@ function createQueue(params: { }, }), } as unknown as JetStreamConnectionManager; - return {queue: new JetStreamWorkerQueue(connectionManager), added, updated}; + return {queue: new JetStreamWorkerQueue(connectionManager), added, updated, consumerAdds}; } describe('jobs stream limits', () => { @@ -143,3 +163,27 @@ describe('jobs stream enqueue shedding', () => { await expect(queue.enqueue('extractEmbeds', {})).rejects.toBe(failure); }); }); + +describe('lane consumer filters', () => { + it('keeps consuming retired subjects so legacy jobs are drained instead of orphaned', async () => { + const {queue, consumerAdds} = createQueue({existing: LEGACY_CONFIG}); + await queue.ensureConsumers(WORKER_LANES); + const lifecycle = consumerAdds.find((config) => config.durable_name === 'workers_lifecycle'); + expect(lifecycle?.filter_subjects).toContain('jobs.sendSystemDm'); + expect(lifecycle?.filter_subjects).toContain('jobs.sendScheduledMessage'); + }); + + it('leaves lanes without retired tasks filtering only their own subjects', async () => { + const {queue, consumerAdds} = createQueue({existing: LEGACY_CONFIG}); + await queue.ensureConsumers(WORKER_LANES); + const unfurl = consumerAdds.find((config) => config.durable_name === 'workers_unfurl'); + expect(unfurl?.filter_subjects).toEqual(['jobs.extractEmbeds']); + }); + + it('never claims the same subject from two lane consumers', async () => { + const {queue, consumerAdds} = createQueue({existing: LEGACY_CONFIG}); + await queue.ensureConsumers(WORKER_LANES); + const allSubjects = consumerAdds.flatMap((config) => config.filter_subjects); + expect(new Set(allSubjects).size).toBe(allSubjects.length); + }); +}); diff --git a/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts b/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts index 661912944..b11d2be4c 100644 --- a/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts +++ b/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts @@ -1,7 +1,12 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import {describe, expect, it} from 'vitest'; -import {resolveCronSchedulerEnabled, resolveWorkerLanes, WORKER_LANES} from '../WorkerLaneConfig'; +import { + resolveCronSchedulerEnabled, + resolveWorkerLanes, + validateLaneCompleteness, + WORKER_LANES, +} from '../WorkerLaneConfig'; describe('WorkerLaneConfig', () => { it('returns all lanes in all_lanes mode', () => { @@ -56,6 +61,33 @@ describe('WorkerLaneConfig', () => { expect(embedLane[0]!.name).toBe('unfurl'); expect(embedLane[0]!.taskTypes).toEqual(['extractEmbeds']); }); + it('keeps the retired scheduled message subject on the lifecycle lane only', () => { + const lanes = resolveWorkerLanes({ + mode: 'all_lanes', + laneConcurrencyOverrides: {}, + }); + const lifecycleLane = lanes.find((lane) => lane.name === 'lifecycle'); + expect(lifecycleLane?.retiredTaskTypes).toEqual(['sendScheduledMessage']); + expect(lifecycleLane?.taskTypes).not.toContain('sendScheduledMessage'); + for (const lane of lanes.filter((lane) => lane.name !== 'lifecycle')) { + expect(lane.retiredTaskTypes).toEqual([]); + } + }); + it('never claims a retired subject from a single_task lane', () => { + const lanes = resolveWorkerLanes({ + mode: 'single_task', + taskName: 'processStripeWebhook', + laneConcurrencyOverrides: {}, + }); + expect(lanes[0]!.retiredTaskTypes).toEqual([]); + }); + it('rejects a registry that brings a retired task name back', () => { + const registry = Object.fromEntries(WORKER_LANES.flatMap((lane) => lane.taskTypes).map((task) => [task, () => {}])); + expect(() => validateLaneCompleteness(registry)).not.toThrow(); + expect(() => validateLaneCompleteness({...registry, sendScheduledMessage: () => {}})).toThrow( + /Retired tasks registered again: sendScheduledMessage/, + ); + }); it('throws when single_lane mode has no lane configured', () => { expect(() => resolveWorkerLanes({ diff --git a/fluxer_api/src/api/worker/tests/WorkerRetiredTask.test.ts b/fluxer_api/src/api/worker/tests/WorkerRetiredTask.test.ts new file mode 100644 index 000000000..80c228ce9 --- /dev/null +++ b/fluxer_api/src/api/worker/tests/WorkerRetiredTask.test.ts @@ -0,0 +1,132 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {JsMsg} from 'nats'; +import {afterEach, beforeAll, describe, expect, it, vi} from 'vitest'; +import type {IJobLedgerRepository} from '../../jobs/IJobLedgerRepository'; +import {setInjectedWorkerService} from '../../middleware/ServiceRegistry'; +import {NoopWorkerService} from '../../test/NoopWorkerService'; +import {WorkerRunner} from '../WorkerRunner'; + +const RETIRED_TASK_TYPE = 'sendScheduledMessage'; +const RETIRED_REASON = 'task type retired'; +const LEDGER_JOB_ID = '1509197195776110592'; + +const queueStub = { + getConnectionManager: () => { + throw new Error('WorkerRunner tests never consume messages'); + }, + getStreamName: () => 'JOBS', + publishToDlq: vi.fn(), +}; + +const ledgerStub = { + markDeadletter: vi.fn(), +}; + +class TestWorkerRunner extends WorkerRunner { + async runJob(taskType: string, msg: JsMsg): Promise { + return await this.processJob(taskType, msg); + } +} + +function createRunner(): TestWorkerRunner { + return new TestWorkerRunner({ + tasks: {}, + retiredTaskTypes: [RETIRED_TASK_TYPE], + queue: queueStub, + consumerName: 'workers_lifecycle', + laneName: 'lifecycle', + ledger: ledgerStub as unknown as IJobLedgerRepository, + concurrency: 8, + maxDeliver: 25, + ackWaitMs: 60000, + }); +} + +function createJobMessage(taskType: string, payload: Record) { + const envelope = { + payload, + run_at: new Date(Date.now() + 30 * 24 * 60 * 60 * 1000).toISOString(), + max_attempts: 5, + priority: 0, + created_at: new Date().toISOString(), + }; + return { + seq: 7, + subject: `jobs.${taskType}`, + redelivered: false, + data: new TextEncoder().encode(JSON.stringify(envelope)), + info: {deliveryCount: 1}, + ack: vi.fn(), + nak: vi.fn(), + term: vi.fn(), + working: vi.fn(), + }; +} + +describe('Retired worker task types', () => { + beforeAll(() => { + setInjectedWorkerService(new NoopWorkerService()); + }); + + afterEach(() => { + queueStub.publishToDlq.mockReset(); + ledgerStub.markDeadletter.mockReset(); + }); + + it('dead-letters a legacy job and closes its ledger row instead of redelivering it', async () => { + const runner = createRunner(); + const msg = createJobMessage(RETIRED_TASK_TYPE, { + userId: '1', + scheduledMessageId: '2', + __jobId: LEDGER_JOB_ID, + }); + + await expect(runner.runJob(RETIRED_TASK_TYPE, msg as unknown as JsMsg)).resolves.toBe(false); + + expect(queueStub.publishToDlq).toHaveBeenCalledTimes(1); + expect(queueStub.publishToDlq).toHaveBeenCalledWith( + RETIRED_TASK_TYPE, + {userId: '1', scheduledMessageId: '2'}, + expect.objectContaining({errorMessage: RETIRED_REASON, lane: 'lifecycle', originalSeq: 7}), + ); + expect(ledgerStub.markDeadletter).toHaveBeenCalledWith(BigInt(LEDGER_JOB_ID), RETIRED_REASON); + expect(msg.term).toHaveBeenCalledWith(RETIRED_REASON); + expect(msg.nak).not.toHaveBeenCalled(); + expect(msg.ack).not.toHaveBeenCalled(); + }); + + it('dead-letters a legacy job that carries no ledger id', async () => { + const runner = createRunner(); + const msg = createJobMessage(RETIRED_TASK_TYPE, {userId: '1', scheduledMessageId: '2'}); + + await expect(runner.runJob(RETIRED_TASK_TYPE, msg as unknown as JsMsg)).resolves.toBe(false); + + expect(queueStub.publishToDlq).toHaveBeenCalledTimes(1); + expect(ledgerStub.markDeadletter).not.toHaveBeenCalled(); + expect(msg.term).toHaveBeenCalledWith(RETIRED_REASON); + expect(msg.nak).not.toHaveBeenCalled(); + }); + + it('redelivers a retired job when the dead-letter publish fails', async () => { + const runner = createRunner(); + queueStub.publishToDlq.mockRejectedValueOnce(new Error('no responders')); + const msg = createJobMessage(RETIRED_TASK_TYPE, {__jobId: LEDGER_JOB_ID}); + + await expect(runner.runJob(RETIRED_TASK_TYPE, msg as unknown as JsMsg)).resolves.toBe(false); + + expect(ledgerStub.markDeadletter).not.toHaveBeenCalled(); + expect(msg.term).not.toHaveBeenCalled(); + expect(msg.nak).toHaveBeenCalledTimes(1); + }); + + it('still terminates a task type that was never registered or retired', async () => { + const runner = createRunner(); + const msg = createJobMessage('neverShippedTask', {}); + + await expect(runner.runJob('neverShippedTask', msg as unknown as JsMsg)).resolves.toBe(false); + + expect(queueStub.publishToDlq).not.toHaveBeenCalled(); + expect(msg.term).toHaveBeenCalledWith(expect.stringMatching(/unknown task type/)); + }); +}); diff --git a/fluxer_messages/src/shard_impl.rs b/fluxer_messages/src/shard_impl.rs index 78b28ff54..1e4704a70 100644 --- a/fluxer_messages/src/shard_impl.rs +++ b/fluxer_messages/src/shard_impl.rs @@ -3101,7 +3101,7 @@ impl From for Message { #[cfg(test)] mod tests { use super::*; - use fluxer_svc::transport::InMemoryTransport; + use fluxer_svc::transport::{InMemoryTransport, TransportSubscriber, reply_message}; use serde_json::json; #[test] @@ -3581,10 +3581,70 @@ mod tests { .unwrap() } + fn authored_message(message_id: i64) -> Message { + decode_postgres_message(json!({ + "channel_id": {"__fluxer_type": "bigint", "value": "10"}, + "bucket": 416, + "message_id": {"__fluxer_type": "bigint", "value": message_id.to_string()}, + "author_id": {"__fluxer_type": "bigint", "value": "1472426752046002208"}, + "content": "kept" + })) + .unwrap() + } + + fn legacy_string_author_message(message_id: i64) -> Message { + decode_postgres_message(json!({ + "channel_id": {"__fluxer_type": "bigint", "value": "10"}, + "bucket": 416, + "message_id": {"__fluxer_type": "bigint", "value": message_id.to_string()}, + "author_id": "1472426752046002208", + "content": "kept" + })) + .unwrap() + } + + async fn stub_user_service(transport: &InMemoryTransport) -> tokio::task::JoinHandle<()> { + let mut subscriber = transport.subscribe("svc.users").await.unwrap(); + let transport = transport.clone(); + tokio::spawn(async move { + while let Some(message) = subscriber.next().await { + let _ = reply_message(&message, &transport, b"\"NotFound\"").await; + } + }) + } + fn recorded_deletions(deleted: &DeletedMessageKeys) -> Vec<(i64, i32, i64)> { deleted.lock().unwrap().clone() } + #[tokio::test] + async fn build_path_keeps_authored_rows_from_older_releases() { + let deleted = DeletedMessageKeys::default(); + let shard = recording_shard(&deleted); + let users = stub_user_service(&shard.transport).await; + let wrapped_id = 1_509_197_195_776_110_592; + let legacy_id = 1_509_197_195_776_110_593; + + let responses = shard + .build_api_responses_from_messages( + vec![ + authored_message(wrapped_id), + legacy_string_author_message(legacy_id), + ], + build_options(), + ) + .await + .unwrap(); + users.abort(); + + assert_eq!(responses.len(), 2); + assert_eq!(responses[0].id, wrapped_id.to_string()); + assert_eq!(responses[1].id, legacy_id.to_string()); + assert_eq!(responses[0].author.id, "1472426752046002208"); + assert_eq!(responses[1].author.id, "1472426752046002208"); + assert!(recorded_deletions(&deleted).is_empty()); + } + #[tokio::test] async fn build_path_reaps_orphaned_messages() { let deleted = DeletedMessageKeys::default();