Merge pull request #1359 from PicPeak/fix/qa-shutdown-and-revocation

fix: complete graceful shutdown and revoke tokens without expiry
This commit is contained in:
Paul Nothaft
2026-09-08 19:22:50 +02:00
committed by GitHub
13 changed files with 466 additions and 31 deletions
@@ -0,0 +1,96 @@
const knex = require('knex');
const jwt = require('jsonwebtoken');
const { randomUUID } = require('crypto');
const migration = require('../../migrations/core/211_revocations_without_expiry');
for (const client of ['sqlite3', 'pg']) {
const enabled = client !== 'pg' || process.env.PICPEAK_PG_TEST_URL;
(enabled ? describe : describe.skip)(`token revocation expiry (${client})`, () => {
let db, owner, schema, revocation;
const sign = claims => jwt.sign({ id: 1, type: 'admin', jti: randomUUID(), ...claims }, process.env.JWT_SECRET);
beforeAll(async () => {
if (client === 'pg') {
schema = `revocation_${randomUUID().replace(/-/g, '')}`;
owner = knex({ client, connection: process.env.PICPEAK_PG_TEST_URL });
await owner.schema.createSchema(schema);
db = knex({ client, connection: process.env.PICPEAK_PG_TEST_URL, searchPath: [schema] });
} else {
db = knex({ client, connection: { filename: ':memory:' }, useNullAsDefault: true });
}
// Exercise the upgrade from the real legacy NOT NULL schema as well as
// repeated migration runs, without sharing another test's database.
await require('../../migrations/legacy/017_add_token_revocation_tables').up(db);
await db('revoked_tokens').insert({ token_id: 'existing', expires_at: '2099-01-01T00:00:00.000Z' });
await migration.up(db);
await migration.up(db);
jest.resetModules();
jest.doMock('../../src/database/db', () => ({ db }));
revocation = require('../../src/utils/tokenRevocation');
});
afterAll(async () => {
await db?.destroy();
if (owner) { await owner.schema.dropSchema(schema, true); await owner.destroy(); }
jest.dontMock('../../src/database/db');
});
it('preserves existing revocations and their unique key during upgrade', async () => {
expect(await db('revoked_tokens').where({ token_id: 'existing' }).first()).toBeTruthy();
await expect(db('revoked_tokens').insert({ token_id: 'existing', expires_at: null })).rejects.toThrow();
});
it.each([true, false])('permanently revokes a token without exp (jti: %s)', async withJti => {
const token = sign(withJti ? {} : { jti: undefined });
const payload = jwt.verify(token, process.env.JWT_SECRET);
expect(await revocation.isTokenRevoked(payload)).toBe(false);
expect(await revocation.revokeToken(token, 'logout')).toBe(true);
expect(await revocation.revokeToken(token, 'logout')).toBe(true);
await revocation.cleanupExpiredRevocations();
expect(await revocation.isTokenRevoked(payload)).toBe(true);
const rows = await db('revoked_tokens').where({ token_id: revocation.buildTokenId(payload) });
expect(rows).toHaveLength(1);
expect(rows[0].expires_at).toBeNull();
});
it('cleans up expired revocations and retains future ones', async () => {
const expired = sign({ exp: Math.floor(Date.now() / 1000) - 60 });
const future = sign({ exp: Math.floor(Date.now() / 1000) + 3600 });
expect(await revocation.revokeToken(expired, 'logout')).toBe(true);
expect(await revocation.revokeToken(future, 'logout')).toBe(true);
await revocation.cleanupExpiredRevocations();
expect(await revocation.isTokenRevoked(jwt.decode(expired))).toBe(false);
expect(await revocation.isTokenRevoked(jwt.decode(future))).toBe(true);
});
it.each([true, false])('upgrades an expiring entry with the same key permanently (jti: %s)', async withJti => {
const claims = { id: 99, iat: Math.floor(Date.now() / 1000), jti: withJti ? randomUUID() : undefined };
const expiring = sign({ ...claims, exp: claims.iat - 60 });
const permanent = sign(claims);
expect(await revocation.revokeToken(expiring, 'logout')).toBe(true);
expect(await revocation.revokeToken(permanent, 'logout')).toBe(true);
expect(await revocation.revokeToken(expiring, 'logout')).toBe(true);
await revocation.cleanupExpiredRevocations();
expect(await revocation.isTokenRevoked(jwt.decode(permanent))).toBe(true);
});
it('retains a signed token whose numeric expiry cannot fit a database timestamp', async () => {
const token = sign({ exp: 1e100 });
expect(await revocation.revokeToken(token, 'logout')).toBe(true);
await revocation.cleanupExpiredRevocations();
expect(await revocation.isTokenRevoked(jwt.decode(token))).toBe(true);
});
it('refuses a rollback that would remove permanent revocations', async () => {
await expect(migration.down(db)).rejects.toThrow('permanent token revocations');
expect((await db('revoked_tokens').columnInfo('expires_at')).nullable).toBe(true);
// A rollback with only expiring records remains supported and reversible.
await db('revoked_tokens').whereNull('expires_at').delete();
await migration.down(db);
await migration.down(db);
expect((await db('revoked_tokens').columnInfo('expires_at')).nullable).toBe(false);
await migration.up(db);
expect(await db('revoked_tokens').where({ token_id: 'existing' }).first()).toBeTruthy();
});
});
}
@@ -0,0 +1,75 @@
const request = require('supertest');
const jwt = require('jsonwebtoken');
const { randomUUID } = require('crypto');
const { bootCrmDb, seedMinimal, assignAdminRole, buildRouteApp } = require('../integration/helpers/crmDb');
let db, cleanup, adminId, customerId, eventId, apps, revocation;
const slug = 'logout-revocation';
const cases = [
['auth', '/logout', 'admin', 'admin_token'],
['auth', '/gallery/logout', 'gallery', `gallery_token_${slug}`],
['customerAuth', '/logout', 'customer', 'customer_token'],
['adminAuth', '/logout', 'admin', 'admin_token'],
];
const sign = type => jwt.sign({
type, ...(type === 'admin' ? { id: adminId } : type === 'customer' ? { customerId } : { eventId, eventSlug: slug }),
jti: randomUUID(),
}, process.env.JWT_SECRET, { issuer: 'picpeak-auth' });
beforeAll(async () => {
({ db, cleanup } = await bootCrmDb());
({ adminId, customerId } = await seedMinimal(db));
await assignAdminRole(db, adminId);
const event = await require('../../src/services/eventCreationService').createEvent({
event_type: 'wedding', event_name: 'Logout revocation', event_date: '2026-10-01',
slug, password: 'Logout-Strong-Password-924!', expiration_days: 30,
customer_email: '[email protected]', admin_email: '[email protected]',
}, { actor: { id: adminId }, source: 'v1' });
eventId = event.id;
await db('events').where({ id: eventId }).update({ slug });
revocation = require('../../src/utils/tokenRevocation');
apps = Object.fromEntries(['auth', 'adminAuth', 'customerAuth'].map(name => [
name, buildRouteApp('/', require(`../../src/routes/${name}`)),
]));
});
afterEach(() => jest.restoreAllMocks());
afterAll(async () => {
await require('../../src/services/serviceShutdown').stopServices();
if (cleanup) await cleanup();
});
it.each(cases)('%s%s revokes a no-expiry %s cookie session', async (route, path, type, cookie) => {
const token = sign(type);
const sessionApp = type === 'customer' ? apps.customerAuth : apps.auth;
const sessionPath = type === 'gallery' ? `/session?slug=${slug}` : '/session';
const session = () => request(sessionApp).get(sessionPath).set('Cookie', `${cookie}=${token}`);
const before = await session();
expect(before.status).toBe(200);
if (type !== 'customer') expect(before.body.valid).toBe(true);
const res = await request(apps[route]).post(path).set('Cookie', `${cookie}=${token}`).send({ slug });
expect(res.status).toBe(200);
expect(res.headers['set-cookie'].some(value => value.startsWith(`${cookie}=;`))).toBe(true);
await revocation.cleanupExpiredRevocations();
expect(await revocation.isTokenRevoked(jwt.decode(token))).toBe(true);
const after = await session();
if (type === 'customer') expect(after.status).toBe(401);
else expect(after.body.valid).toBe(false);
});
it.each(cases)('%s%s reports failed persistence for %s logout and clears its cookie', async (route, path, type, cookie) => {
const token = sign(type);
const realQuery = db.client.query;
jest.spyOn(db.client, 'query').mockImplementation(function (connection, query) {
if (/^insert into [`"]revoked_tokens[`"]/.test(query.sql)) {
return Promise.reject(new Error('simulated revocation write failure'));
}
return realQuery.call(this, connection, query);
});
const res = await request(apps[route]).post(path).set('Cookie', `${cookie}=${token}`).send({ slug });
expect(res.status).toBe(500);
expect(res.body.error).toBeTruthy();
expect(res.body.message).not.toBe('Logged out successfully');
expect(res.headers['set-cookie'].some(value => value.startsWith(`${cookie}=;`))).toBe(true);
expect(await revocation.isTokenRevoked(jwt.decode(token))).toBe(false);
});
@@ -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();
});
}
});
@@ -20,7 +20,7 @@ jest.mock('../../src/database/db', () => {
const dbFn = () => ({ const dbFn = () => ({
insert(row) { insert(row) {
inserted.push(row); inserted.push(row);
return { onConflict: () => ({ ignore: async () => undefined }) }; return { onConflict: () => ({ ignore: async () => undefined, merge: async () => undefined }) };
}, },
}); });
return { db: dbFn }; return { db: dbFn };
@@ -0,0 +1,26 @@
/** Non-expiring JWTs need revocation records that cleanup never removes. */
exports.up = async function (knex) {
if (!await knex.schema.hasTable('revoked_tokens')) return;
if (!await knex.schema.hasColumn('revoked_tokens', 'expires_at')) return;
const column = await knex('revoked_tokens').columnInfo('expires_at');
if (!column.nullable) {
await knex.schema.alterTable('revoked_tokens', table => {
table.timestamp('expires_at').nullable().alter();
});
}
};
exports.down = async function (knex) {
if (!await knex.schema.hasTable('revoked_tokens')) return;
if (!await knex.schema.hasColumn('revoked_tokens', 'expires_at')) return;
// Refuse to discard permanent revocations or silently give them a TTL.
if (await knex('revoked_tokens').whereNull('expires_at').first()) {
throw new Error('Cannot roll back while permanent token revocations exist');
}
const column = await knex('revoked_tokens').columnInfo('expires_at');
if (column.nullable) {
await knex.schema.alterTable('revoked_tokens', table => {
table.timestamp('expires_at').notNullable().alter();
});
}
};
+1 -1
View File
@@ -421,7 +421,7 @@ async function initializeDatabase() {
table.integer('user_id').nullable(); // User who owned the token table.integer('user_id').nullable(); // User who owned the token
table.string('token_type', 20); // admin, gallery, etc. table.string('token_type', 20); // admin, gallery, etc.
table.timestamp('revoked_at').defaultTo(db.fn.now()); table.timestamp('revoked_at').defaultTo(db.fn.now());
table.timestamp('expires_at').notNullable(); // When token would have expired table.timestamp('expires_at').nullable(); // NULL retains tokens without a known expiry
table.string('reason', 100); // password_change, logout, compromised, etc. table.string('reason', 100); // password_change, logout, compromised, etc.
table.text('metadata'); // Additional JSON data table.text('metadata'); // Additional JSON data
+4 -3
View File
@@ -182,16 +182,17 @@ router.post('/logout', adminAuth, handleAsync(async (req, res) => {
// old header-only read skipped revocation entirely for cookie-based logout, // old header-only read skipped revocation entirely for cookie-based logout,
// leaving the JWT valid until expiry while reporting a successful logout. // leaving the JWT valid until expiry while reporting a successful logout.
const token = req.token; const token = req.token;
clearAdminAuthCookie(res);
if (token) { if (token) {
// End the in-memory session AND revoke the JWT (GHSA-cjqh) — the token // End the in-memory session AND revoke the JWT (GHSA-cjqh) — the token
// is otherwise valid until expiry, so photoAuth/adminAuth would keep // is otherwise valid until expiry, so photoAuth/adminAuth would keep
// honouring it after logout. isTokenRevoked() checks this store. // honouring it after logout. isTokenRevoked() checks this store.
endSession(token); endSession(token);
const { revokeToken } = require('../utils/tokenRevocation'); const { revokeToken } = require('../utils/tokenRevocation');
await revokeToken(token, 'logout'); if (!await revokeToken(token, 'logout')) {
throw new Error('Token revocation failed');
}
} }
// Clear the auth cookie so the browser stops sending the (now revoked) JWT.
clearAdminAuthCookie(res);
// Log activity // Log activity
await logActivity('admin_logout', await logActivity('admin_logout',
+9 -2
View File
@@ -356,7 +356,9 @@ router.post('/logout', async (req, res) => {
if (token) { if (token) {
// Revoke the token so it can't be reused, then end the session // Revoke the token so it can't be reused, then end the session
await revokeToken(token, 'user_logout'); if (!await revokeToken(token, 'user_logout')) {
throw new Error('Token revocation failed');
}
endSession(token); endSession(token);
try { try {
@@ -400,6 +402,8 @@ router.post('/logout', async (req, res) => {
res.json({ message: 'Logged out successfully', ...(ssoLogoutUrl ? { ssoLogoutUrl } : {}) }); res.json({ message: 'Logged out successfully', ...(ssoLogoutUrl ? { ssoLogoutUrl } : {}) });
} catch (error) { } catch (error) {
clearAdminAuthCookie(res);
clearGalleryAuthCookies(res);
errorResponse(res, error, 500, 'Logout failed'); errorResponse(res, error, 500, 'Logout failed');
} }
}); });
@@ -714,11 +718,14 @@ router.post('/gallery/logout', async (req, res) => {
const { slug } = req.body || {}; const { slug } = req.body || {};
const token = getGalleryTokenFromRequest(req, slug); const token = getGalleryTokenFromRequest(req, slug);
if (token) { if (token) {
await revokeToken(token, 'gallery_logout'); if (!await revokeToken(token, 'gallery_logout')) {
throw new Error('Token revocation failed');
}
} }
clearGalleryAuthCookies(res, slug); clearGalleryAuthCookies(res, slug);
res.json({ message: 'Logged out successfully' }); res.json({ message: 'Logged out successfully' });
} catch (error) { } catch (error) {
clearGalleryAuthCookies(res, req.body?.slug);
errorResponse(res, error, 500, 'Logout failed'); errorResponse(res, error, 500, 'Logout failed');
} }
}); });
+3 -1
View File
@@ -181,7 +181,9 @@ router.post('/logout', async (req, res) => {
try { try {
const token = getCustomerTokenFromRequest(req); const token = getCustomerTokenFromRequest(req);
if (token) { if (token) {
await revokeToken(token, 'user_logout'); if (!await revokeToken(token, 'user_logout')) {
throw new Error('Token revocation failed');
}
} }
clearCustomerAuthCookie(res); clearCustomerAuthCookie(res);
res.json({ message: 'Logged out successfully' }); res.json({ message: 'Logged out successfully' });
+17 -5
View File
@@ -29,6 +29,7 @@
const os = require('os'); const os = require('os');
const { db } = require('../database/db'); const { db } = require('../database/db');
const logger = require('../utils/logger'); const logger = require('../utils/logger');
const { createInterruptibleSleep } = require('../utils/interruptibleSleep');
const { processPhoto } = require('./photoProcessor'); const { processPhoto } = require('./photoProcessor');
const POLL_INTERVAL_MS = parseInt(process.env.UPLOAD_PROCESSOR_POLL_MS || '1000', 10); const POLL_INTERVAL_MS = parseInt(process.env.UPLOAD_PROCESSOR_POLL_MS || '1000', 10);
@@ -69,7 +70,9 @@ let running = false;
let workerHandles = []; let workerHandles = [];
let janitorHandle = null; 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() { function isPostgres() {
const c = db.client.config.client; const c = db.client.config.client;
@@ -175,12 +178,13 @@ async function janitorLoop() {
} }
function start() { function start() {
if (running) return; if (running || stopping) return;
if (process.env.UPLOAD_PROCESSOR_DISABLED === 'true') { if (process.env.UPLOAD_PROCESSOR_DISABLED === 'true') {
logger.info('backgroundProcessor: disabled via UPLOAD_PROCESSOR_DISABLED'); logger.info('backgroundProcessor: disabled via UPLOAD_PROCESSOR_DISABLED');
return; return;
} }
waits = createInterruptibleSleep();
running = true; running = true;
workerHandles = []; workerHandles = [];
for (let i = 0; i < CONCURRENCY; i++) { for (let i = 0; i < CONCURRENCY; i++) {
@@ -199,12 +203,20 @@ function start() {
); );
} }
async function stop() { function stop() {
if (!running) return; if (stopping) return stopping;
if (!running) return Promise.resolve();
running = false; running = false;
await Promise.all([...workerHandles, janitorHandle].filter(Boolean)); // 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 = []; workerHandles = [];
janitorHandle = null; janitorHandle = null;
waits = null;
stopping = null;
});
return stopping;
} }
module.exports = { start, stop, claimNextPhoto }; module.exports = { start, stop, claimNextPhoto };
+19 -5
View File
@@ -28,6 +28,7 @@
const { db } = require('../database/db'); const { db } = require('../database/db');
const logger = require('../utils/logger'); const logger = require('../utils/logger');
const { createInterruptibleSleep } = require('../utils/interruptibleSleep');
const { processPhotoFaces } = require('./faceProcessor'); const { processPhotoFaces } = require('./faceProcessor');
const { SidecarUnavailableError } = require('./faceClient'); const { SidecarUnavailableError } = require('./faceClient');
const { TransientSourceError } = require('./faceProcessor'); const { TransientSourceError } = require('./faceProcessor');
@@ -220,7 +221,9 @@ let running = false;
let workerHandles = []; let workerHandles = [];
let janitorHandle = null; 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() { function isPostgres() {
const c = db.client.config.client; const c = db.client.config.client;
@@ -289,6 +292,8 @@ async function workerLoop(workerIdx) {
continue; continue;
} }
if (!running) break;
let claimed; let claimed;
try { try {
claimed = await claimNextPhoto(currentlyDeferredEventIds()); claimed = await claimNextPhoto(currentlyDeferredEventIds());
@@ -407,12 +412,13 @@ async function janitorLoop() {
} }
function start() { function start() {
if (running) return; if (running || stopping) return;
if (process.env.FACE_PROCESSOR_DISABLED === 'true') { if (process.env.FACE_PROCESSOR_DISABLED === 'true') {
logger.info('faceQueue: disabled via FACE_PROCESSOR_DISABLED'); logger.info('faceQueue: disabled via FACE_PROCESSOR_DISABLED');
return; return;
} }
waits = createInterruptibleSleep();
running = true; running = true;
workerHandles = []; workerHandles = [];
for (let i = 0; i < CONCURRENCY; i++) { for (let i = 0; i < CONCURRENCY; i++) {
@@ -432,12 +438,20 @@ function start() {
); );
} }
async function stop() { function stop() {
if (!running) return; if (stopping) return stopping;
if (!running) return Promise.resolve();
running = false; running = false;
await Promise.all([...workerHandles, janitorHandle].filter(Boolean)); // 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 = []; workerHandles = [];
janitorHandle = null; janitorHandle = null;
waits = null;
stopping = null;
});
return stopping;
} }
// drainConsolidation, touchedEvents and consolidationRetryAt are exported for // drainConsolidation, touchedEvents and consolidationRetryAt are exported for
+25
View File
@@ -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 };
+18 -8
View File
@@ -6,6 +6,7 @@
const jwt = require('jsonwebtoken'); const jwt = require('jsonwebtoken');
const { db } = require('../database/db'); const { db } = require('../database/db');
const logger = require('./logger'); const logger = require('./logger');
const MAX_SQL_EXPIRY_SECONDS = Date.parse('9999-12-31T23:59:59Z') / 1000;
/** /**
* Add a token to the revocation list * Add a token to the revocation list
@@ -52,20 +53,28 @@ async function revokeToken(token, reason, metadata = {}) {
// undefined/string slip through and cause an INSERT type error. // undefined/string slip through and cause an INSERT type error.
const userIdNumeric = Number.isInteger(payload.id) ? payload.id : null; const userIdNumeric = Number.isInteger(payload.id) ? payload.id : null;
// onConflict.ignore: revoking an already-revoked token is a no-op, // JWT permits a missing exp. Keep that revocation permanently: an
// not an error. Hits the unique (token_id) index when the same JWT // arbitrary fallback TTL would make the token usable again after cleanup.
// is logged out twice (e.g. duplicate /logout from two tabs, or a // Retain unrepresentable expiries too, using the common SQL/ISO date range.
// session-expiry path that races with an explicit logout). The const expiresAt = Number.isFinite(payload.exp)
// previous insert was authoritative; nothing to do. && payload.exp >= 0 && payload.exp <= MAX_SQL_EXPIRY_SECONDS
await db('revoked_tokens').insert({ ? new Date(Math.ceil(payload.exp * 1000)).toISOString()
: null;
// Duplicate logouts are idempotent. A permanent revocation must also
// upgrade an existing expiring entry with the same legacy key or jti;
// logging out an expiring token must never shorten that retention again.
const insert = db('revoked_tokens').insert({
token_id: buildTokenId(payload), token_id: buildTokenId(payload),
user_id: userIdNumeric, user_id: userIdNumeric,
token_type: payload.type, token_type: payload.type,
revoked_at: new Date().toISOString(), revoked_at: new Date().toISOString(),
expires_at: new Date(payload.exp * 1000).toISOString(), expires_at: expiresAt,
reason, reason,
metadata: JSON.stringify(metadata) metadata: JSON.stringify(metadata)
}).onConflict('token_id').ignore(); }).onConflict('token_id');
if (expiresAt === null) await insert.merge({ expires_at: null });
else await insert.ignore();
logger.info('Token revoked', { logger.info('Token revoked', {
userId: payload.id ?? payload.customerId ?? null, userId: payload.id ?? payload.customerId ?? null,
@@ -131,6 +140,7 @@ async function revokeAllUserTokens(userId, reason) {
async function cleanupExpiredRevocations() { async function cleanupExpiredRevocations() {
try { try {
const deleted = await db('revoked_tokens') const deleted = await db('revoked_tokens')
.whereNotNull('expires_at')
.where('expires_at', '<', new Date().toISOString()) .where('expires_at', '<', new Date().toISOString())
.delete(); .delete();