From 44277e6aa2acd597880efafb596e2e6b1e91eaf9 Mon Sep 17 00:00:00 2001 From: Hampus Date: Mon, 31 Aug 2026 22:06:12 +0200 Subject: [PATCH] fix(worker): drain in-flight jobs before the runner stops (#2274) --- fluxer_api/src/api/worker/WorkerRunner.ts | 7 +- .../worker/tests/WorkerRunnerShutdown.test.ts | 127 ++++++++++++++++++ 2 files changed, 133 insertions(+), 1 deletion(-) create mode 100644 fluxer_api/src/api/worker/tests/WorkerRunnerShutdown.test.ts diff --git a/fluxer_api/src/api/worker/WorkerRunner.ts b/fluxer_api/src/api/worker/WorkerRunner.ts index 33ee95964..806e05c75 100644 --- a/fluxer_api/src/api/worker/WorkerRunner.ts +++ b/fluxer_api/src/api/worker/WorkerRunner.ts @@ -66,6 +66,7 @@ export class WorkerRunner { private readonly ledger: IJobLedgerRepository; private running = false; private consumerMessages: ConsumerMessages | null = null; + private processingLoop: Promise | null = null; private readonly inFlightJobs = new Set>(); 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'); } diff --git a/fluxer_api/src/api/worker/tests/WorkerRunnerShutdown.test.ts b/fluxer_api/src/api/worker/tests/WorkerRunnerShutdown.test.ts new file mode 100644 index 000000000..c93ef34d0 --- /dev/null +++ b/fluxer_api/src/api/worker/tests/WorkerRunnerShutdown.test.ts @@ -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 = []; + private notify: (() => void) | null = null; + private closed = false; + + push(msg: JsMsg): void { + this.pending.push(msg); + this.wake(); + } + + async close(): Promise { + this.closed = true; + this.wake(); + } + + async *[Symbol.asyncIterator](): AsyncGenerator { + while (true) { + while (this.pending.length > 0) { + yield this.pending.shift()!; + } + if (this.closed) { + return; + } + await new Promise((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, 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((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); + }); +});