fix(worker): drain in-flight jobs before the runner stops (#2274)

This commit is contained in:
Hampus
2026-08-31 22:06:12 +02:00
committed by GitHub
parent bb7e8cc6f1
commit 44277e6aa2
2 changed files with 133 additions and 1 deletions
+6 -1
View File
@@ -66,6 +66,7 @@ export class WorkerRunner {
private readonly ledger: IJobLedgerRepository;
private running = false;
private consumerMessages: ConsumerMessages | null = null;
private processingLoop: Promise<void> | null = null;
private readonly inFlightJobs = new Set<Promise<void>>();
constructor(options: WorkerRunnerOptions) {
@@ -95,7 +96,7 @@ export class WorkerRunner {
max_messages: prefetch,
idle_heartbeat: 5000,
});
this.processMessages().catch((error) => {
this.processingLoop = this.processMessages().catch((error) => {
Logger.error({workerId: this.workerId, err: error}, 'Worker message processing failed unexpectedly');
});
}
@@ -109,6 +110,10 @@ export class WorkerRunner {
await this.consumerMessages.close();
this.consumerMessages = null;
}
if (this.processingLoop !== null) {
await this.processingLoop;
this.processingLoop = null;
}
Logger.info({workerId: this.workerId}, 'Worker stopped');
}
@@ -0,0 +1,127 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {ConsumerMessages, JsMsg} from 'nats';
import {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 TASK_TYPE = 'processInactivityDeletions';
class FakeConsumerMessages {
private readonly pending: Array<JsMsg> = [];
private notify: (() => void) | null = null;
private closed = false;
push(msg: JsMsg): void {
this.pending.push(msg);
this.wake();
}
async close(): Promise<void> {
this.closed = true;
this.wake();
}
async *[Symbol.asyncIterator](): AsyncGenerator<JsMsg> {
while (true) {
while (this.pending.length > 0) {
yield this.pending.shift()!;
}
if (this.closed) {
return;
}
await new Promise<void>((resolve) => {
this.notify = resolve;
});
}
}
private wake(): void {
const notify = this.notify;
this.notify = null;
notify?.();
}
}
function createQueueStub(messages: FakeConsumerMessages) {
return {
getConnectionManager: () => ({
getJetStreamClient: () => ({
consumers: {
get: async () => ({
consume: async () => messages as unknown as ConsumerMessages,
}),
},
}),
}),
getStreamName: () => 'JOBS',
publishToDlq: vi.fn(),
};
}
function createRunner(task: () => Promise<void>, messages: FakeConsumerMessages): WorkerRunner {
return new WorkerRunner({
tasks: {[TASK_TYPE]: task},
queue: createQueueStub(messages),
consumerName: 'workers_batch',
laneName: 'batch',
ledger: {} as IJobLedgerRepository,
concurrency: 12,
maxDeliver: 25,
ackWaitMs: 120000,
});
}
function createJobMessage() {
const envelope = {
payload: {},
max_attempts: 5,
priority: 0,
created_at: new Date().toISOString(),
};
return {
seq: 1,
subject: `jobs.${TASK_TYPE}`,
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('Worker runner shutdown', () => {
beforeAll(() => {
setInjectedWorkerService(new NoopWorkerService());
});
it('drains in-flight jobs before stop resolves', async () => {
let started = false;
let finished = false;
let release: () => void = () => {};
const pending = new Promise<void>((resolve) => {
release = resolve;
});
const messages = new FakeConsumerMessages();
const runner = createRunner(async () => {
started = true;
await pending;
finished = true;
}, messages);
const msg = createJobMessage();
await runner.start();
messages.push(msg as unknown as JsMsg);
await vi.waitFor(() => expect(started).toBe(true));
setTimeout(release, 0);
await runner.stop();
expect(finished).toBe(true);
expect(msg.ack).toHaveBeenCalledTimes(1);
});
});