fix(api): dead-letter retired worker task types (#2323)

This commit is contained in:
Hampus
2026-09-01 20:47:19 +02:00
committed by GitHub
parent bc40073a02
commit aa267b54ec
8 changed files with 337 additions and 5 deletions
@@ -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,
@@ -5,6 +5,7 @@ import type {APIWorkerLaneName, APIWorkerMode} from '../config/APIConfig';
interface LaneSettings {
readonly consumerName: string;
readonly tasks: ReadonlyArray<string>;
readonly retiredTasks: ReadonlyArray<string>;
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<WorkerTaskName>;
retiredTaskTypes: ReadonlyArray<string>;
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<string, unknown>): voi
if (missingFromRegistry.length > 0) {
errors.push(`Lane tasks not found in registry: ${missingFromRegistry.join(', ')}`);
}
const retiredCollisions = WORKER_LANES.flatMap<string>((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')}`);
}
+1
View File
@@ -240,6 +240,7 @@ export async function startWorkerMain(): Promise<void> {
}
const runner = new WorkerRunner({
tasks: laneTasks,
retiredTaskTypes: lane.retiredTaskTypes,
queue,
consumerName: lane.consumerName,
laneName: lane.name,
+49
View File
@@ -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<string, WorkerTaskHandler>;
retiredTaskTypes?: ReadonlyArray<string>;
queue: WorkerRunnerQueue;
consumerName: string;
laneName: string;
@@ -56,6 +58,7 @@ interface WorkerRunnerOptions {
export class WorkerRunner {
private readonly tasks: Record<string, WorkerTaskHandler>;
private readonly retiredTaskTypes: Set<string>;
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<boolean> {
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<void> {
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<typeof setInterval> {
const heartbeat = setInterval(
() => {
@@ -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<string>;
}
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<Partial<StreamConfig>>; updated: Array<Partial<StreamConfig>>} {
}): {
queue: JetStreamWorkerQueue;
added: Array<Partial<StreamConfig>>;
updated: Array<Partial<StreamConfig>>;
consumerAdds: Array<ConsumerAddConfig>;
} {
const added: Array<Partial<StreamConfig>> = [];
const updated: Array<Partial<StreamConfig>> = [];
const consumerAdds: Array<ConsumerAddConfig> = [];
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);
});
});
@@ -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({
@@ -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<boolean> {
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<string, unknown>) {
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/));
});
});
+61 -1
View File
@@ -3101,7 +3101,7 @@ impl From<MessageDbRow> 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();