Skip to content

Commit 8489f84

Browse files
committed
Handle repeated SIGTERM shutdown
1 parent 4a576a3 commit 8489f84

2 files changed

Lines changed: 200 additions & 12 deletions

File tree

src/worker.ts

Lines changed: 51 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,12 @@ export class Worker {
7474
#generator?: AsyncGenerator<WorkerCycle, void, unknown>
7575
#pool?: JobPool
7676
#lastStalledCheck = 0
77-
#shutdownHandler?: () => Promise<void>
77+
#shutdownHandlers?: {
78+
SIGINT: () => void
79+
SIGTERM: () => void
80+
}
81+
#shutdownInProgress = false
82+
#sigtermReceived = false
7883

7984
/** Unique identifier for this worker instance */
8085
get id() {
@@ -507,25 +512,59 @@ export class Worker {
507512
return
508513
}
509514

510-
this.#shutdownHandler = async () => {
511-
debug('received shutdown signal, stopping worker...')
515+
this.#shutdownInProgress = false
516+
this.#sigtermReceived = false
517+
this.#shutdownHandlers = {
518+
SIGINT: () => void this.#handleShutdownSignal('SIGINT'),
519+
SIGTERM: () => void this.#handleShutdownSignal('SIGTERM'),
520+
}
521+
522+
process.on('SIGINT', this.#shutdownHandlers.SIGINT)
523+
process.on('SIGTERM', this.#shutdownHandlers.SIGTERM)
524+
}
525+
526+
async #handleShutdownSignal(signal: NodeJS.Signals): Promise<void> {
527+
const logger = QueueManager.getLogger()
512528

513-
if (this.#onShutdownSignal) {
514-
await this.#onShutdownSignal()
529+
if (signal === 'SIGTERM') {
530+
if (this.#sigtermReceived) {
531+
logger.warn(
532+
{ workerId: this.#id, signal },
533+
'Received SIGTERM while shutdown is already in progress. Sending SIGKILL.'
534+
)
535+
process.kill(process.pid, 'SIGKILL')
536+
return
515537
}
516538

517-
await this.stop()
539+
this.#sigtermReceived = true
540+
541+
logger.info(
542+
{ workerId: this.#id, signal, runningJobs: this.#pool?.size ?? 0 },
543+
'Received SIGTERM. Waiting for running jobs to drain before shutdown. ' +
544+
'If shutdown never completes, a job may be blocked.'
545+
)
546+
}
547+
548+
if (this.#shutdownInProgress) {
549+
return
550+
}
551+
552+
this.#shutdownInProgress = true
553+
554+
debug('received shutdown signal, stopping worker...')
555+
556+
if (this.#onShutdownSignal) {
557+
await this.#onShutdownSignal()
518558
}
519559

520-
process.on('SIGINT', this.#shutdownHandler)
521-
process.on('SIGTERM', this.#shutdownHandler)
560+
await this.stop()
522561
}
523562

524563
#removeShutdownHandlers() {
525-
if (this.#shutdownHandler) {
526-
process.off('SIGINT', this.#shutdownHandler)
527-
process.off('SIGTERM', this.#shutdownHandler)
528-
this.#shutdownHandler = undefined
564+
if (this.#shutdownHandlers) {
565+
process.off('SIGINT', this.#shutdownHandlers.SIGINT)
566+
process.off('SIGTERM', this.#shutdownHandlers.SIGTERM)
567+
this.#shutdownHandlers = undefined
529568
}
530569
}
531570

tests/worker.spec.ts

Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { test } from '@japa/runner'
22
import { setTimeout } from 'node:timers/promises'
33
import { Worker } from '../src/worker.js'
44
import { MemoryAdapter, memory } from './_mocks/memory_adapter.js'
5+
import { MemoryLogger } from './_mocks/memory_logger.js'
56
import { ChaosAdapter } from './_mocks/chaos_adapter.js'
67
import type { QueueManagerConfig } from '../src/types/main.js'
78
import { Locator } from '../src/locator.js'
@@ -1268,6 +1269,154 @@ test.group('Worker', () => {
12681269

12691270
assert.isTrue(callbackInvoked, 'onShutdownSignal should be called on SIGINT')
12701271
})
1272+
1273+
test('logs when SIGTERM starts draining jobs', async ({ assert, cleanup }) => {
1274+
let markStarted: () => void = () => {}
1275+
let releaseJob: () => void = () => {}
1276+
const started = new Promise<void>((resolve) => {
1277+
markStarted = resolve
1278+
})
1279+
const release = new Promise<void>((resolve) => {
1280+
releaseJob = resolve
1281+
})
1282+
1283+
class DrainingJob extends Job {
1284+
async execute() {
1285+
markStarted()
1286+
await release
1287+
}
1288+
}
1289+
1290+
const logger = new MemoryLogger()
1291+
const sharedAdapter = memory()()
1292+
1293+
const localConfig = {
1294+
default: 'memory',
1295+
adapters: { memory: () => sharedAdapter },
1296+
logger,
1297+
worker: {
1298+
gracefulShutdown: true,
1299+
},
1300+
}
1301+
1302+
Locator.register('DrainingJob', DrainingJob)
1303+
1304+
const worker = new Worker(localConfig)
1305+
let startPromise: Promise<void> = Promise.resolve()
1306+
1307+
cleanup(async () => {
1308+
releaseJob()
1309+
Locator.clear()
1310+
await Promise.race([startPromise, setTimeout(500)])
1311+
})
1312+
1313+
await sharedAdapter.push({
1314+
id: 'draining-job-1',
1315+
name: 'DrainingJob',
1316+
payload: {},
1317+
attempts: 0,
1318+
priority: 0,
1319+
})
1320+
1321+
startPromise = worker.start(['default'])
1322+
await started
1323+
1324+
process.emit('SIGTERM')
1325+
1326+
const stoppedBeforeDrain = await Promise.race([
1327+
startPromise.then(() => true),
1328+
setTimeout(30).then(() => false),
1329+
])
1330+
1331+
const log = logger.logs.find((entry) => entry.level === 'info')
1332+
assert.exists(log)
1333+
assert.include(log!.message, 'Received SIGTERM')
1334+
assert.include(log!.message, 'Waiting for running jobs to drain before shutdown')
1335+
assert.include(log!.message, 'a job may be blocked')
1336+
assert.equal(log!.obj?.signal, 'SIGTERM')
1337+
assert.equal(log!.obj?.runningJobs, 1)
1338+
assert.isFalse(stoppedBeforeDrain)
1339+
1340+
releaseJob()
1341+
await Promise.race([startPromise, setTimeout(500)])
1342+
})
1343+
1344+
test('forces process exit when SIGTERM is received during shutdown', async ({
1345+
assert,
1346+
cleanup,
1347+
}) => {
1348+
let markStarted: () => void = () => {}
1349+
let releaseJob: () => void = () => {}
1350+
const started = new Promise<void>((resolve) => {
1351+
markStarted = resolve
1352+
})
1353+
const release = new Promise<void>((resolve) => {
1354+
releaseJob = resolve
1355+
})
1356+
1357+
class BlockingJob extends Job {
1358+
async execute() {
1359+
markStarted()
1360+
await release
1361+
}
1362+
}
1363+
1364+
const logger = new MemoryLogger()
1365+
const sharedAdapter = memory()()
1366+
const localConfig = {
1367+
default: 'memory',
1368+
adapters: { memory: () => sharedAdapter },
1369+
logger,
1370+
worker: {
1371+
gracefulShutdown: true,
1372+
},
1373+
}
1374+
1375+
Locator.register('BlockingJob', BlockingJob)
1376+
1377+
const worker = new Worker(localConfig)
1378+
const originalKill = process.kill
1379+
let killed: { pid: number; signal?: string | number } | undefined
1380+
let startPromise: Promise<void> = Promise.resolve()
1381+
1382+
process.kill = ((pid: number, signal?: string | number) => {
1383+
killed = { pid, signal }
1384+
return true
1385+
}) as typeof process.kill
1386+
1387+
cleanup(async () => {
1388+
process.kill = originalKill
1389+
releaseJob()
1390+
Locator.clear()
1391+
await Promise.race([startPromise, setTimeout(500)])
1392+
})
1393+
1394+
await sharedAdapter.push({
1395+
id: 'blocking-job-1',
1396+
name: 'BlockingJob',
1397+
payload: {},
1398+
attempts: 0,
1399+
priority: 0,
1400+
})
1401+
1402+
startPromise = worker.start(['default'])
1403+
await started
1404+
1405+
process.emit('SIGTERM')
1406+
await setTimeout(10)
1407+
1408+
process.emit('SIGTERM')
1409+
1410+
assert.deepEqual(killed, { pid: process.pid, signal: 'SIGKILL' })
1411+
1412+
const warning = logger.logs.find((entry) => entry.level === 'warn')
1413+
assert.exists(warning)
1414+
assert.include(warning!.message, 'Sending SIGKILL')
1415+
1416+
releaseJob()
1417+
process.kill = originalKill
1418+
await Promise.race([startPromise, setTimeout(500)])
1419+
})
12711420
})
12721421

12731422
test.group('Worker | jobFactory', () => {

0 commit comments

Comments
 (0)