From a31a2e25e2666989ee226c0d7884e23f50fe58da Mon Sep 17 00:00:00 2001 From: Paul Nothaft <53005142+the-luap@users.noreply.github.com> Date: Tue, 8 Sep 2026 15:53:53 +0200 Subject: [PATCH] fix: interrupt idle worker waits during shutdown --- .../__tests__/services/workerShutdown.test.js | 167 ++++++++++++++++++ backend/src/services/backgroundProcessor.js | 26 ++- backend/src/services/faceQueue.js | 28 ++- backend/src/utils/interruptibleSleep.js | 25 +++ 4 files changed, 232 insertions(+), 14 deletions(-) create mode 100644 backend/__tests__/services/workerShutdown.test.js create mode 100644 backend/src/utils/interruptibleSleep.js diff --git a/backend/__tests__/services/workerShutdown.test.js b/backend/__tests__/services/workerShutdown.test.js new file mode 100644 index 00000000..ab7b816d --- /dev/null +++ b/backend/__tests__/services/workerShutdown.test.js @@ -0,0 +1,167 @@ +jest.mock('../../src/utils/logger', () => ({ info: jest.fn(), warn: jest.fn(), error: jest.fn() })); + +describe.each(['backgroundProcessor', 'faceQueue'])('%s shutdown', name => { + let worker, db, processPhoto, featureEnabled, janitorUpdate, releases; + const prefix = name === 'faceQueue' ? 'FACE_PROCESSOR' : 'UPLOAD_PROCESSOR'; + let previousEnv; + class SidecarUnavailableError extends Error {} + + function deferred() { + let resolve; + const promise = new Promise(done => { resolve = done; }); + releases.push(resolve); + return { promise, resolve }; + } + + beforeEach(() => { + jest.resetModules(); + jest.useFakeTimers(); + releases = []; + previousEnv = { ...process.env }; + delete process.env[`${prefix}_DISABLED`]; + process.env[`${prefix}_CONCURRENCY`] = '2'; + process.env[`${prefix}_POLL_MS`] = '2000'; + process.env.FACE_PROCESSOR_BACKOFF_MS = '30000'; + processPhoto = jest.fn().mockResolvedValue({ status: 'skipped' }); + featureEnabled = jest.fn().mockResolvedValue(true); + janitorUpdate = jest.fn().mockResolvedValue(0); + const chain = { where: jest.fn().mockReturnThis(), update: janitorUpdate }; + db = jest.fn(() => chain); + db.client = { config: { client: 'pg' } }; + // Claims and processing are controlled at the I/O boundary; the real + // workers, janitors, idle sleeps and stopServices run in every test. + db.transaction = jest.fn().mockResolvedValue(null); + jest.doMock('../../src/database/db', () => ({ db })); + jest.doMock('../../src/services/photoProcessor', () => ({ processPhoto })); + jest.doMock('../../src/services/faceProcessor', () => ({ + processPhotoFaces: processPhoto, TransientSourceError: class extends Error {}, + })); + jest.doMock('../../src/services/faceClient', () => ({ SidecarUnavailableError })); + jest.doMock('../../src/services/faceSettings', () => ({ + isFeatureEnabled: featureEnabled, isEnabledForEvent: jest.fn().mockResolvedValue(false), + })); + worker = require(`../../src/services/${name}`); + }); + + afterEach(async () => { + releases.forEach(resolve => resolve()); + const stopped = worker.stop(); + // Also cleans up the original, non-interruptible implementation when a + // regression assertion fails; the test never needs to wait a real minute. + await jest.advanceTimersByTimeAsync(60000); + await stopped; + jest.useRealTimers(); + process.env = previousEnv; + }); + + async function expectPromptStop(stop = () => worker.stop()) { + let done = false; + const stopped = stop().then(() => { done = true; }); + await jest.advanceTimersByTimeAsync(0); + expect(done).toBe(true); + expect(jest.getTimerCount()).toBe(0); + await stopped; + } + + it('wakes all idle workers and the minute-long janitor through stopServices', async () => { + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(jest.getTimerCount()).toBe(3); + // Confirm normal polling still runs before shutdown. + await jest.advanceTimersByTimeAsync(2000); + expect(db.transaction).toHaveBeenCalledTimes(4); + await expectPromptStop(() => require('../../src/services/serviceShutdown').stopServices()); + }); + + it('wakes claim-error backoff and can start a fresh run after stopping', async () => { + db.transaction.mockRejectedValue(new Error('database unavailable')); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + await expectPromptStop(); + db.transaction.mockResolvedValue(null); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(jest.getTimerCount()).toBe(3); + await expectPromptStop(); + }); + + it('drains active processing for every stop caller and prevents overlapping restarts', async () => { + const processing = deferred(); + processPhoto.mockReturnValue(processing.promise); + db.transaction.mockResolvedValueOnce({ id: 1 }); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(processPhoto).toHaveBeenCalledWith(1); + let done = false; + const first = worker.stop(); + expect(worker.stop()).toBe(first); + first.then(() => { done = true; }); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(done).toBe(false); + expect(db.transaction).toHaveBeenCalledTimes(2); + processing.resolve({ status: 'skipped' }); + await expectPromptStop(); + expect(done).toBe(true); + expect(db.transaction).toHaveBeenCalledTimes(2); + }); + + it('drains a claim already in flight without starting another poll', async () => { + const claim = deferred(); + db.transaction.mockReturnValueOnce(claim.promise); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + const stopped = worker.stop(); + claim.resolve({ id: 2 }); + await expectPromptStop(() => stopped); + expect(processPhoto).toHaveBeenCalledWith(2); + expect(db.transaction).toHaveBeenCalledTimes(2); + }); + + it('does not schedule new waits when pending database work finishes after stop', async () => { + const claim = deferred(); + const janitor = deferred(); + db.transaction.mockReturnValue(claim.promise); + janitorUpdate.mockReturnValue(janitor.promise); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + let done = false; + const stopped = worker.stop().then(() => { done = true; }); + claim.resolve(null); + await jest.advanceTimersByTimeAsync(0); + expect(done).toBe(false); + janitor.resolve(0); + await expectPromptStop(() => stopped); + }); + + if (name === 'faceQueue') { + it('interrupts the ten-second sleep with faces disabled by default', async () => { + featureEnabled.mockResolvedValue(false); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(db.transaction).not.toHaveBeenCalled(); + expect(jest.getTimerCount()).toBe(3); + await expectPromptStop(); + }); + + it('releases a claimed photo and interrupts the sidecar outage backoff', async () => { + db.transaction.mockResolvedValueOnce({ id: 3, event_id: 5 }); + processPhoto.mockRejectedValue(new SidecarUnavailableError('offline')); + worker.start(); + await jest.advanceTimersByTimeAsync(0); + expect(janitorUpdate).toHaveBeenCalledWith({ face_status: 'pending', face_started_at: null }); + await expectPromptStop(); + expect(worker.inFlightByEvent.size).toBe(0); + }); + + it('does not claim new work after a pending feature check resolves during shutdown', async () => { + const feature = deferred(); + featureEnabled.mockReturnValue(feature.promise); + worker.start(); + const stopped = worker.stop(); + feature.resolve(true); + await expectPromptStop(() => stopped); + expect(db.transaction).not.toHaveBeenCalled(); + }); + } +}); diff --git a/backend/src/services/backgroundProcessor.js b/backend/src/services/backgroundProcessor.js index 004c9c51..61a7d5ca 100644 --- a/backend/src/services/backgroundProcessor.js +++ b/backend/src/services/backgroundProcessor.js @@ -29,6 +29,7 @@ const os = require('os'); const { db } = require('../database/db'); const logger = require('../utils/logger'); +const { createInterruptibleSleep } = require('../utils/interruptibleSleep'); const { processPhoto } = require('./photoProcessor'); const POLL_INTERVAL_MS = parseInt(process.env.UPLOAD_PROCESSOR_POLL_MS || '1000', 10); @@ -69,7 +70,9 @@ let running = false; let workerHandles = []; let janitorHandle = null; -const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); +let waits = null; +let stopping = null; +const sleep = ms => waits.sleep(ms); function isPostgres() { const c = db.client.config.client; @@ -175,12 +178,13 @@ async function janitorLoop() { } function start() { - if (running) return; + if (running || stopping) return; if (process.env.UPLOAD_PROCESSOR_DISABLED === 'true') { logger.info('backgroundProcessor: disabled via UPLOAD_PROCESSOR_DISABLED'); return; } + waits = createInterruptibleSleep(); running = true; workerHandles = []; for (let i = 0; i < CONCURRENCY; i++) { @@ -199,12 +203,20 @@ function start() { ); } -async function stop() { - if (!running) return; +function stop() { + if (stopping) return stopping; + if (!running) return Promise.resolve(); running = false; - await Promise.all([...workerHandles, janitorHandle].filter(Boolean)); - workerHandles = []; - janitorHandle = null; + // Interrupt idle/backoff waits only. Claims, processing and janitor work + // already in flight still drain before the database can be closed. + waits.cancel(); + stopping = Promise.all([...workerHandles, janitorHandle].filter(Boolean)).finally(() => { + workerHandles = []; + janitorHandle = null; + waits = null; + stopping = null; + }); + return stopping; } module.exports = { start, stop, claimNextPhoto }; diff --git a/backend/src/services/faceQueue.js b/backend/src/services/faceQueue.js index 190e1454..72c0a80d 100644 --- a/backend/src/services/faceQueue.js +++ b/backend/src/services/faceQueue.js @@ -28,6 +28,7 @@ const { db } = require('../database/db'); const logger = require('../utils/logger'); +const { createInterruptibleSleep } = require('../utils/interruptibleSleep'); const { processPhotoFaces } = require('./faceProcessor'); const { SidecarUnavailableError } = require('./faceClient'); const { TransientSourceError } = require('./faceProcessor'); @@ -220,7 +221,9 @@ let running = false; let workerHandles = []; let janitorHandle = null; -const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); +let waits = null; +let stopping = null; +const sleep = ms => waits.sleep(ms); function isPostgres() { const c = db.client.config.client; @@ -289,6 +292,8 @@ async function workerLoop(workerIdx) { continue; } + if (!running) break; + let claimed; try { claimed = await claimNextPhoto(currentlyDeferredEventIds()); @@ -407,12 +412,13 @@ async function janitorLoop() { } function start() { - if (running) return; + if (running || stopping) return; if (process.env.FACE_PROCESSOR_DISABLED === 'true') { logger.info('faceQueue: disabled via FACE_PROCESSOR_DISABLED'); return; } + waits = createInterruptibleSleep(); running = true; workerHandles = []; for (let i = 0; i < CONCURRENCY; i++) { @@ -432,12 +438,20 @@ function start() { ); } -async function stop() { - if (!running) return; +function stop() { + if (stopping) return stopping; + if (!running) return Promise.resolve(); running = false; - await Promise.all([...workerHandles, janitorHandle].filter(Boolean)); - workerHandles = []; - janitorHandle = null; + // Interrupt idle/backoff waits only. Claims, processing and janitor work + // already in flight still drain before the database can be closed. + waits.cancel(); + stopping = Promise.all([...workerHandles, janitorHandle].filter(Boolean)).finally(() => { + workerHandles = []; + janitorHandle = null; + waits = null; + stopping = null; + }); + return stopping; } // drainConsolidation, touchedEvents and consolidationRetryAt are exported for diff --git a/backend/src/utils/interruptibleSleep.js b/backend/src/utils/interruptibleSleep.js new file mode 100644 index 00000000..533a7420 --- /dev/null +++ b/backend/src/utils/interruptibleSleep.js @@ -0,0 +1,25 @@ +/** A run owns its waits; cancelling wakes every waiter and clears its timer. */ +function createInterruptibleSleep() { + let cancelled = false; + const waiters = new Set(); + return { + sleep(ms) { + if (cancelled) return Promise.resolve(); + return new Promise(resolve => { + const finish = () => { + clearTimeout(timer); + waiters.delete(finish); + resolve(); + }; + const timer = setTimeout(finish, ms); + waiters.add(finish); + }); + }, + cancel() { + cancelled = true; + for (const finish of waiters) finish(); + }, + }; +} + +module.exports = { createInterruptibleSleep };