fix(email): show a queue nobody is working instead of reporting all-clear
Closes #1262. "Gallery email queued" reads as a delivery confirmation, and System Health agreed with it: "No stuck or failed emails -- all clear", while not one email had gone out. Both statements were true and neither was the one the admin needed. Queueing writes an email_queue row at status='pending', retry_count 0 -- nothing more. /failures matched only status='failed' or pending-with-retry_count>=3, so it matched none of those rows, and there are two ordinary ways they never leave that state: - startEmailQueueProcessor() was never reached, so nothing polls the queue. - Every pass returns early. processEmailQueue bails when the transporter will not initialise, before it touches a single row, so retry_count stays 0 and no error_message is ever written. A working SMTP test button does not contradict this: that path builds its own transport. adminSystem.js made it worse by reporting `emailProcessor: { status: 'active' }` as a literal, so the one place that named the worker always said it was fine. - emailProcessor records what each pass did -- started, lastRunAt, lastResult, lastError -- and exports getQueueProcessorStatus(). The transporter bail and the queue-query failure, the two silent early returns, both write lastError. - /failures gains `waitingEmails`: pending, under the retry cap, past any scheduled_at, and queued more than 10 minutes ago. The predicate mirrors the processor's own pickup query, so a row listed there is one it should already have taken; rows over the cap stay in `stuckEmails` and are not counted twice. A future scheduled_at is left alone -- split-payment invoices and the business-hours floor park rows deliberately. - System Health leads with the processor's state (running / stopped / degraded) and lists waiting emails in their own table. The all-clear now needs both buckets empty. - adminSystem reports the real processor state instead of the literal. - The two "queued" toasts say the queue processor is what sends it and where to look if it doesn't arrive. 8 route tests, all 8 failing before the change.
This commit is contained in:
@@ -0,0 +1,158 @@
|
||||
/**
|
||||
* System Health must not report "all clear" over a queue nobody is working (#1262).
|
||||
*
|
||||
* "Gallery email queued" reads as a delivery confirmation, and the two ways the
|
||||
* queue silently stops — the processor never started, or every pass returns
|
||||
* early because the transport will not initialise — leave every row at
|
||||
* status='pending' with retry_count 0. The old /failures query matched only
|
||||
* status='failed' or pending-with-retry_count>=3, so it matched none of them
|
||||
* and the page said everything was fine while nothing had been sent.
|
||||
*/
|
||||
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
const os = require('os');
|
||||
|
||||
process.env.NODE_ENV = 'test';
|
||||
process.env.TEST_DATABASE_PATH = path.join(
|
||||
fs.mkdtempSync(path.join(os.tmpdir(), 'picpeak-mailhealth-')), 'db.sqlite',
|
||||
);
|
||||
process.env.JWT_SECRET = process.env.JWT_SECRET || 'mailhealth-test-secret';
|
||||
|
||||
const request = require('supertest');
|
||||
const bcrypt = require('bcrypt');
|
||||
const jwt = require('jsonwebtoken');
|
||||
|
||||
const { bootCrmDb, seedMinimal, buildRouteApp } = require('../integration/helpers/crmDb');
|
||||
|
||||
const MINUTE = 60 * 1000;
|
||||
const ago = (ms) => new Date(Date.now() - ms).toISOString();
|
||||
const ahead = (ms) => new Date(Date.now() + ms).toISOString();
|
||||
|
||||
describe('GET /admin/system-health/failures — waiting emails (#1262)', () => {
|
||||
let db; let cleanup; let app; let token;
|
||||
|
||||
const queue = (row) => db('email_queue').insert({
|
||||
recipient_email: '[email protected]',
|
||||
email_type: 'gallery_created',
|
||||
email_data: '{}',
|
||||
status: 'pending',
|
||||
retry_count: 0,
|
||||
created_at: ago(60 * MINUTE),
|
||||
...row,
|
||||
});
|
||||
|
||||
const failures = async () => {
|
||||
const res = await request(app)
|
||||
.get('/admin/system-health/failures')
|
||||
.set('Authorization', `Bearer ${token}`);
|
||||
expect(res.status).toBe(200);
|
||||
return res.body.data || res.body;
|
||||
};
|
||||
|
||||
const typesOf = (rows) => rows.map((r) => r.emailType).sort();
|
||||
|
||||
beforeAll(async () => {
|
||||
({ db, cleanup } = await bootCrmDb());
|
||||
await seedMinimal(db);
|
||||
|
||||
const role = await db('roles').where({ name: 'super_admin' }).first();
|
||||
const inserted = await db('admin_users').insert({
|
||||
username: 'mailhealth-admin',
|
||||
email: '[email protected]',
|
||||
password_hash: await bcrypt.hash('Passw0rd!', 4),
|
||||
role_id: role.id,
|
||||
is_active: 1,
|
||||
created_at: new Date().toISOString(),
|
||||
updated_at: new Date().toISOString(),
|
||||
}).returning('id');
|
||||
const adminId = inserted[0]?.id ?? inserted[0];
|
||||
token = jwt.sign(
|
||||
{ id: adminId, username: 'mailhealth-admin', type: 'admin', role: 'super_admin', loginTime: Date.now() },
|
||||
process.env.JWT_SECRET,
|
||||
{ expiresIn: '1h', issuer: 'picpeak-auth' },
|
||||
);
|
||||
|
||||
app = buildRouteApp('/admin/system-health', require('../../src/routes/adminSystemHealth'));
|
||||
});
|
||||
|
||||
afterAll(async () => { await cleanup(); });
|
||||
afterEach(async () => { await db('email_queue').del(); });
|
||||
|
||||
it('reports a due pending email the processor never picked up', async () => {
|
||||
// Exactly the shape "Gallery email queued" leaves behind when the worker
|
||||
// is not running: pending, no retries, no error, no scheduled_at.
|
||||
await queue({ email_type: 'gallery_created' });
|
||||
|
||||
const body = await failures();
|
||||
expect(typesOf(body.waitingEmails)).toEqual(['gallery_created']);
|
||||
expect(body.counts.waitingEmails).toBe(1);
|
||||
// ...and it is NOT a failure, so the two buckets stay distinct.
|
||||
expect(body.stuckEmails).toEqual([]);
|
||||
});
|
||||
|
||||
it('leaves a freshly queued email alone — the processor wakes every 60s', async () => {
|
||||
await queue({ email_type: 'customer_invitation', created_at: ago(30 * 1000) });
|
||||
|
||||
const body = await failures();
|
||||
expect(body.waitingEmails).toEqual([]);
|
||||
expect(body.counts.waitingEmails).toBe(0);
|
||||
});
|
||||
|
||||
it('leaves an email scheduled for later alone', async () => {
|
||||
// Split-payment invoices and the business-hours floor both park rows in
|
||||
// the future on purpose. Not being sent yet is the point of those.
|
||||
await queue({ email_type: 'invoice_due', scheduled_at: ahead(3 * 24 * 60 * MINUTE) });
|
||||
|
||||
const body = await failures();
|
||||
expect(body.waitingEmails).toEqual([]);
|
||||
});
|
||||
|
||||
it('counts a past-due scheduled email once its moment has come', async () => {
|
||||
await queue({ email_type: 'invoice_due', scheduled_at: ago(30 * MINUTE) });
|
||||
|
||||
const body = await failures();
|
||||
expect(typesOf(body.waitingEmails)).toEqual(['invoice_due']);
|
||||
});
|
||||
|
||||
it('does not double-count a retry-exhausted email as waiting', async () => {
|
||||
// retry_count >= 3 is already the `stuckEmails` bucket; listing it in both
|
||||
// would inflate the badge and make the two tables disagree.
|
||||
await queue({ email_type: 'quote_sent', retry_count: 3, error_message: 'template missing' });
|
||||
|
||||
const body = await failures();
|
||||
expect(body.waitingEmails).toEqual([]);
|
||||
expect(typesOf(body.stuckEmails)).toEqual(['quote_sent']);
|
||||
});
|
||||
|
||||
it('ignores emails that were sent', async () => {
|
||||
await queue({ email_type: 'gallery_created', status: 'sent', sent_at: ago(20 * MINUTE) });
|
||||
|
||||
const body = await failures();
|
||||
expect(body.waitingEmails).toEqual([]);
|
||||
expect(body.stuckEmails).toEqual([]);
|
||||
});
|
||||
|
||||
it('reports what the queue processor last did', async () => {
|
||||
const body = await failures();
|
||||
// Never started in this process — which is the condition that makes a
|
||||
// pending row invisible, so the page has to be able to say it.
|
||||
expect(body.processor).toEqual(expect.objectContaining({ started: false }));
|
||||
expect(body.processor).toHaveProperty('lastRunAt');
|
||||
expect(body.processor).toHaveProperty('lastError');
|
||||
});
|
||||
|
||||
it('surfaces the transport failure that makes every pass a no-op', async () => {
|
||||
const { processEmailQueue, getQueueProcessorStatus } = require('../../src/services/emailProcessor');
|
||||
await queue({ email_type: 'gallery_created' });
|
||||
|
||||
// No SMTP configured, so initializeTransporter() yields nothing and the
|
||||
// pass returns early. Before #1262 that left no trace anywhere.
|
||||
await processEmailQueue();
|
||||
|
||||
expect(getQueueProcessorStatus().lastError).toMatch(/transporter could not be initialised/i);
|
||||
const body = await failures();
|
||||
expect(body.processor.lastError).toMatch(/transporter could not be initialised/i);
|
||||
expect(body.counts.waitingEmails).toBe(1);
|
||||
});
|
||||
});
|
||||
@@ -9,6 +9,7 @@ const { formatBoolean } = require('../utils/dbCompat');
|
||||
const logger = require('../utils/logger');
|
||||
const { resolveSqlitePath } = require('../utils/databaseEngine');
|
||||
const { checkForUpdates, getCurrentChannel, getCurrentVersion, getReleasesSince, compareVersions } = require('../services/updateCheckService');
|
||||
const { getQueueProcessorStatus } = require('../services/emailProcessor');
|
||||
const { getAppSetting, upsertAppSetting } = require('../utils/appSettings');
|
||||
const { parseWhatsNew } = require('../utils/whatsNew');
|
||||
const { detectEnvironment, generateUpdateInstructions } = require('../services/environmentService');
|
||||
@@ -346,7 +347,17 @@ router.get('/status', adminAuth, requirePermission(['settings.view', 'system.vie
|
||||
services: {
|
||||
fileWatcher: { status: 'active' }, // These would ideally check actual service status
|
||||
expirationChecker: { status: 'active' },
|
||||
emailProcessor: { status: 'active' }
|
||||
// #1262 — this used to be hardcoded 'active', which reported a healthy
|
||||
// worker on a deployment whose queue had never been touched. Report
|
||||
// what the processor itself recorded instead.
|
||||
emailProcessor: (() => {
|
||||
const p = getQueueProcessorStatus();
|
||||
return {
|
||||
status: p.started && !p.lastError ? 'active' : (p.started ? 'degraded' : 'stopped'),
|
||||
lastRunAt: p.lastRunAt,
|
||||
lastError: p.lastError,
|
||||
};
|
||||
})()
|
||||
},
|
||||
timestamp: new Date()
|
||||
};
|
||||
|
||||
@@ -23,6 +23,7 @@ const { requirePermission } = require('../middleware/permissions');
|
||||
const { handleAsync, validateRequest, successResponse } = require('../utils/routeHelpers');
|
||||
const { verifyDocumentArtefacts } = require('../services/backupIntegrityService');
|
||||
const { getCoverageReport } = require('../services/backupCoverageService');
|
||||
const { getQueueProcessorStatus } = require('../services/emailProcessor');
|
||||
const { db } = require('../database/db');
|
||||
|
||||
const router = express.Router();
|
||||
@@ -85,6 +86,24 @@ router.get(
|
||||
}),
|
||||
);
|
||||
|
||||
/**
|
||||
* How long an email may sit due-but-unsent before it counts as waiting rather
|
||||
* than merely in flight. The processor wakes every 60s and takes 10 rows a
|
||||
* pass, so a genuine backlog of ~6000 clears inside this window — anything
|
||||
* still here has not been worked.
|
||||
*/
|
||||
const WAITING_EMAIL_GRACE_MS = 10 * 60 * 1000;
|
||||
|
||||
const mapEmailRow = (r) => ({
|
||||
id: r.id,
|
||||
recipientEmail: r.recipient_email,
|
||||
emailType: r.email_type,
|
||||
status: r.status,
|
||||
retryCount: r.retry_count,
|
||||
errorMessage: r.error_message,
|
||||
createdAt: r.created_at,
|
||||
});
|
||||
|
||||
/**
|
||||
* GET /api/admin/system-health/failures
|
||||
*
|
||||
@@ -94,6 +113,14 @@ router.get(
|
||||
* (status='pending' AND retry_count >= 3 — the processor only picks up
|
||||
* retry_count < 3). Trigger: a 14h window where 'quote_sent' template
|
||||
* errors left invoices unsent with no admin-visible signal.
|
||||
*
|
||||
* #1262 added the other half. A queue nobody is working produces no failures
|
||||
* at all: the rows sit at status='pending' with retry_count 0, matching
|
||||
* neither branch above, and the page reported "all clear" while not one email
|
||||
* had gone out. That happens whenever the processor never started, or every
|
||||
* pass returns early because the transport will not initialise. So the
|
||||
* response also carries emails that are DUE and still unsent
|
||||
* (`waitingEmails`), plus what the processor itself last did (`processor`).
|
||||
*/
|
||||
router.get(
|
||||
'/failures',
|
||||
@@ -110,17 +137,31 @@ router.get(
|
||||
.limit(200)
|
||||
.select('id', 'recipient_email', 'email_type', 'status', 'retry_count', 'error_message', 'created_at');
|
||||
|
||||
// Deliberately mirrors the processor's own pickup predicate — pending,
|
||||
// under the retry cap, and past any scheduled_at — so a row listed here is
|
||||
// one it should already have taken. Rows over the cap are the `stuckEmails`
|
||||
// set above and must not be counted twice.
|
||||
const now = new Date();
|
||||
const dueBefore = new Date(now.getTime() - WAITING_EMAIL_GRACE_MS);
|
||||
const waitingEmails = await db('email_queue')
|
||||
.where('status', 'pending')
|
||||
.where('retry_count', '<', 3)
|
||||
.where('created_at', '<=', dueBefore.toISOString())
|
||||
.andWhere(function () {
|
||||
this.whereNull('scheduled_at').orWhere('scheduled_at', '<=', now.toISOString());
|
||||
})
|
||||
.orderBy('created_at', 'asc')
|
||||
.limit(200)
|
||||
.select('id', 'recipient_email', 'email_type', 'status', 'retry_count', 'error_message', 'created_at');
|
||||
|
||||
return successResponse(res, {
|
||||
stuckEmails: stuckEmails.map((r) => ({
|
||||
id: r.id,
|
||||
recipientEmail: r.recipient_email,
|
||||
emailType: r.email_type,
|
||||
status: r.status,
|
||||
retryCount: r.retry_count,
|
||||
errorMessage: r.error_message,
|
||||
createdAt: r.created_at,
|
||||
})),
|
||||
counts: { stuckEmails: stuckEmails.length },
|
||||
stuckEmails: stuckEmails.map(mapEmailRow),
|
||||
waitingEmails: waitingEmails.map(mapEmailRow),
|
||||
processor: getQueueProcessorStatus(),
|
||||
counts: {
|
||||
stuckEmails: stuckEmails.length,
|
||||
waitingEmails: waitingEmails.length,
|
||||
},
|
||||
});
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -922,9 +922,32 @@ async function renderQueuedEmail(templateKey, variables = {}, to = '') {
|
||||
// because ignoreSchedule also bypasses that cap.
|
||||
//
|
||||
// Returns { processed, sent, failed }.
|
||||
// What the last pass actually did, so System Health can say whether the queue
|
||||
// is being worked at all (#1262). "Queued" is not "delivered", and the two
|
||||
// ways a queue silently stops -- the processor never started, or every pass
|
||||
// returns early because the transport will not initialise -- both leave rows
|
||||
// at status='pending' with retry_count 0, which no failure query matches.
|
||||
const processorStatus = {
|
||||
started: false,
|
||||
lastRunAt: null,
|
||||
lastResult: null,
|
||||
lastError: null,
|
||||
};
|
||||
|
||||
function getQueueProcessorStatus() {
|
||||
return {
|
||||
started: processorStatus.started,
|
||||
lastRunAt: processorStatus.lastRunAt,
|
||||
lastResult: processorStatus.lastResult,
|
||||
lastError: processorStatus.lastError,
|
||||
};
|
||||
}
|
||||
|
||||
async function processEmailQueue({ ignoreSchedule = false, limit = 10, onlyId = null } = {}) {
|
||||
logger.info('Email queue processor: Checking for pending emails...');
|
||||
const result = { processed: 0, sent: 0, failed: 0 };
|
||||
processorStatus.lastRunAt = new Date().toISOString();
|
||||
processorStatus.lastError = null;
|
||||
|
||||
try {
|
||||
// Try to initialize transporter if it's null (in case it failed at startup).
|
||||
@@ -937,6 +960,10 @@ async function processEmailQueue({ ignoreSchedule = false, limit = 10, onlyId =
|
||||
transporter = await initializeTransporter();
|
||||
if (!transporter) {
|
||||
logger.warn('Email transporter could not be initialized, skipping queue processing');
|
||||
// #1262 — the row stays pending with retry_count 0, so nothing in the
|
||||
// queue itself records that this pass did nothing. Say so here.
|
||||
processorStatus.lastError = 'Email transporter could not be initialised — check the SMTP settings';
|
||||
processorStatus.lastResult = result;
|
||||
return result;
|
||||
}
|
||||
}
|
||||
@@ -969,6 +996,8 @@ async function processEmailQueue({ ignoreSchedule = false, limit = 10, onlyId =
|
||||
.limit(limit);
|
||||
} catch (dbError) {
|
||||
logger.error('Failed to query email queue:', dbError);
|
||||
processorStatus.lastError = dbError.message;
|
||||
processorStatus.lastResult = result;
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -1043,8 +1072,10 @@ async function processEmailQueue({ ignoreSchedule = false, limit = 10, onlyId =
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error('Error processing email queue:', error);
|
||||
processorStatus.lastError = error.message;
|
||||
}
|
||||
|
||||
processorStatus.lastResult = result;
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -1203,6 +1234,7 @@ function startEmailQueueProcessor() {
|
||||
});
|
||||
}, 60000);
|
||||
|
||||
processorStatus.started = true;
|
||||
logger.info('Email queue processor started successfully');
|
||||
} else {
|
||||
logger.info('Email queue processor: Already running');
|
||||
@@ -1213,6 +1245,7 @@ function stopEmailQueueProcessor() {
|
||||
if (emailQueueInterval) {
|
||||
clearInterval(emailQueueInterval);
|
||||
emailQueueInterval = null;
|
||||
processorStatus.started = false;
|
||||
logger.info('Email queue processor stopped');
|
||||
}
|
||||
}
|
||||
@@ -1231,6 +1264,7 @@ module.exports = {
|
||||
sendRawEmail,
|
||||
renderQueuedEmail,
|
||||
processEmailQueue,
|
||||
getQueueProcessorStatus,
|
||||
queueEmail,
|
||||
stopEmailQueueProcessor,
|
||||
testEmailConnection,
|
||||
|
||||
Reference in New Issue
Block a user