diff --git a/server/public/admin-newsletter.html b/server/public/admin-newsletter.html index 35d9fe4139..2eb69fa5f5 100644 --- a/server/public/admin-newsletter.html +++ b/server/public/admin-newsletter.html @@ -560,7 +560,7 @@

Add Custom Section

if (status === 'draft') { return `
Draft only
-
This edition has not been queued. Send tests, finish edits, then approve it for the scheduled sender.
+
This edition has not been queued. Finish edits, then approve it for the scheduled sender or approve and send it immediately.
${esc(formatSendWindowChip())} ${recipientCount ? `${esc(String(recipientCount.emailCount))} eligible email recipients` : ''} @@ -701,6 +701,7 @@

Add Custom Section

${esc(status)} ${status === 'draft' || status === 'approved' ? `` : ''} ${status === 'draft' ? `` : ''} + ${status === 'draft' ? `` : ''} ${status === 'approved' ? `` : ''}
${optInBar} @@ -1157,8 +1158,11 @@

CUSTOM ${esc(cs.title || 'Untit async function sendNow(event) { const btn = event?.currentTarget; + const originalButtonText = btn?.textContent || 'Send now'; const count = recipientCount ? ` This will go to ${recipientCount.emailCount} recipients and post to Slack if configured.` : ''; - if (!confirm('Send this approved edition now?' + count)) return; + const isDraft = currentDigest.status === 'draft'; + const prompt = isDraft ? 'Approve this draft and send it now?' : 'Send this approved edition now?'; + if (!confirm(prompt + count)) return; try { if (btn) { btn.disabled = true; btn.textContent = 'Sending...'; } const res = await api('/editions/' + currentDigest.id + '/send-now', { method: 'POST' }); @@ -1173,7 +1177,7 @@

CUSTOM ${esc(cs.title || 'Untit toast(Number.isFinite(sent) ? `Sent ${sent} emails` : 'Sent'); } catch (err) { toast('Send failed: ' + esc(err.message)); - if (btn) { btn.disabled = false; btn.textContent = 'Send now'; } + if (btn) { btn.disabled = false; btn.textContent = originalButtonText; } } } diff --git a/server/src/addie/jobs/weekly-digest.ts b/server/src/addie/jobs/weekly-digest.ts index c4d0853ed1..627d1de8d0 100644 --- a/server/src/addie/jobs/weekly-digest.ts +++ b/server/src/addie/jobs/weekly-digest.ts @@ -19,6 +19,7 @@ import { WorkingGroupDatabase } from '../../db/working-group-db.js'; import { sendChannelMessage } from '../../slack/client.js'; import { sendTrackedBatchMarketingEmails, type TrackedBatchMarketingEmail } from '../../notifications/email.js'; import { renderDigestEmail, renderDigestSlack, renderDigestReview, type DigestSegment } from '../templates/weekly-digest.js'; +import { withNewsletterSendLock } from '../../newsletters/send-lock.js'; import { publishDigestAsPerspective } from '../services/digest-publisher.js'; import { generateCoverForEdition } from '../../newsletters/cover.js'; import { markSuggestionsIncluded } from '../../db/newsletter-suggestions-db.js'; @@ -271,6 +272,25 @@ async function sendApprovedDigest(editionDate: string, etHour: number): Promise< * Called from the scheduled job and from the approval handler (for late approvals). */ export async function sendDigest(digest: DigestRecord): Promise<{ sent: number }> { + const editionDate = new Date(digest.edition_date).toISOString().split('T')[0]; + const locked = await withNewsletterSendLock('the_prompt', digest.id, async () => { + // Re-read under the cross-process lock so a stale approved record cannot + // be delivered after a manual sender has already completed it. + const current = await getDigestByDate(editionDate); + if (!current || current.id !== digest.id || current.status !== 'approved') { + return { sent: 0 }; + } + return deliverDigest(current); + }); + + if (!locked.acquired) { + logger.warn({ digestId: digest.id, editionDate }, 'The Prompt delivery already in progress'); + return { sent: 0 }; + } + return locked.value; +} + +async function deliverDigest(digest: DigestRecord): Promise<{ sent: number }> { if (digest.status !== 'approved') { logger.error({ digestId: digest.id, status: digest.status }, 'sendDigest called on non-approved digest'); return { sent: 0 }; @@ -331,7 +351,10 @@ export async function sendDigest(digest: DigestRecord): Promise<{ sent: number } // Mark as sent if (stats.email_count > 0 || stats.slack_count > 0) { - await markSent(digest.id, stats); + const markedSent = await markSent(digest.id, stats); + if (!markedSent) { + throw new Error(`Failed to finalize The Prompt edition ${digest.id} after delivery`); + } // Publish as perspective for SEO/discoverability (non-blocking) void (async () => { diff --git a/server/src/newsletters/admin-routes.ts b/server/src/newsletters/admin-routes.ts index 1fa26872ac..355039b3fc 100644 --- a/server/src/newsletters/admin-routes.ts +++ b/server/src/newsletters/admin-routes.ts @@ -322,13 +322,29 @@ export function createNewsletterAdminRoutes(config: NewsletterConfig): Router { const id = parseInt(req.params.id, 10); if (isNaN(id)) return res.status(400).json({ error: 'Invalid edition ID' }); - const edition = await config.db.getCurrent(); + let edition = await config.db.getCurrent(); if (!edition || edition.id !== id) return res.status(404).json({ error: 'Edition not found' }); - if (edition.status !== 'approved') { - return res.status(400).json({ error: 'Only approved editions can be sent. Approve the draft first.' }); + if (edition.status === 'draft') { + const approvedBy = req.user?.email || 'admin'; + const approved = await config.db.approve(id, approvedBy); + if (!approved) { + return res.status(409).json({ error: 'Edition status changed. Refresh and try again.' }); + } + edition = approved; + } else if (edition.status !== 'approved') { + return res.status(400).json({ error: 'Only draft or approved editions can be sent.' }); } const result = await sendNewsletter(config, edition); + if (result.outcome === 'busy') { + return res.status(409).json({ error: 'This edition is already being sent.' }); + } + if (result.outcome === 'not_sendable') { + return res.status(409).json({ error: 'Edition status changed. Refresh and try again.' }); + } + if (result.outcome === 'failed') { + return res.status(502).json({ error: 'No newsletter deliveries succeeded. The edition remains approved for review.' }); + } const updated = await config.db.getCurrent(); const digest = updated && updated.id === id ? updated : edition; const subject = config.generateSubject(digest.content); diff --git a/server/src/newsletters/send-lock.ts b/server/src/newsletters/send-lock.ts new file mode 100644 index 0000000000..95d5f59bdf --- /dev/null +++ b/server/src/newsletters/send-lock.ts @@ -0,0 +1,59 @@ +import type { PoolClient } from 'pg'; +import { getPool } from '../db/client.js'; +import { createLogger } from '../logger.js'; + +const logger = createLogger('newsletter-send-lock'); + +export type NewsletterSendLockResult = + | { acquired: true; value: T } + | { acquired: false }; + +/** + * Serialize delivery for one newsletter edition across web and worker + * processes. The caller must re-read the edition while holding the lock so a + * stale approved snapshot cannot be delivered after another sender finishes. + */ +export async function withNewsletterSendLock( + newsletterId: string, + editionId: number, + work: () => Promise, +): Promise> { + const client: PoolClient = await getPool().connect(); + const lockKey = `newsletter-send:${newsletterId}:${editionId}`; + let acquired = false; + let destroyClient = false; + + try { + const lockResult = await client.query<{ acquired: boolean }>( + 'SELECT pg_try_advisory_lock(hashtextextended($1, 0)) AS acquired', + [lockKey], + ); + acquired = lockResult.rows[0]?.acquired === true; + if (!acquired) return { acquired: false }; + + return { acquired: true, value: await work() }; + } finally { + if (acquired) { + try { + const unlockResult = await client.query<{ unlocked: boolean }>( + 'SELECT pg_advisory_unlock(hashtextextended($1, 0)) AS unlocked', + [lockKey], + ); + if (unlockResult.rows[0]?.unlocked !== true) { + destroyClient = true; + logger.warn( + { newsletterId, editionId }, + 'Newsletter send lock was not released; discarding pooled connection', + ); + } + } catch (error) { + destroyClient = true; + logger.warn( + { error, newsletterId, editionId }, + 'Failed to release newsletter send lock; discarding pooled connection', + ); + } + } + client.release(destroyClient); + } +} diff --git a/server/src/newsletters/send-pipeline.ts b/server/src/newsletters/send-pipeline.ts index f5e03173e4..462e066f68 100644 --- a/server/src/newsletters/send-pipeline.ts +++ b/server/src/newsletters/send-pipeline.ts @@ -13,6 +13,7 @@ import { sendTrackedBatchMarketingEmails, type TrackedBatchMarketingEmail } from import { proposeContentForUser, type ContentUser } from '../routes/content.js'; import { generateIllustration } from '../services/illustration-generator.js'; import { createIllustration, approveIllustration } from '../db/illustration-db.js'; +import { withNewsletterSendLock } from './send-lock.js'; const logger = createLogger('newsletter-send'); @@ -23,7 +24,31 @@ const logger = createLogger('newsletter-send'); export async function sendNewsletter( config: NewsletterConfig, edition: EditionRecord, -): Promise<{ sent: number }> { +): Promise<{ sent: number; outcome: 'delivered' | 'busy' | 'not_sendable' | 'failed' }> { + const locked = await withNewsletterSendLock(config.id, edition.id, async () => { + // The edition passed by the caller may be stale by the time the lock is + // acquired. Re-read it under the lock and only deliver an approved row. + const current = await config.db.getCurrent(); + if (!current || current.id !== edition.id || current.status !== 'approved') { + return { sent: 0, outcome: 'not_sendable' as const }; + } + return deliverNewsletter(config, current); + }); + + if (!locked.acquired) { + logger.warn( + { newsletterId: config.id, editionId: edition.id }, + 'Newsletter delivery already in progress', + ); + return { sent: 0, outcome: 'busy' }; + } + return locked.value; +} + +async function deliverNewsletter( + config: NewsletterConfig, + edition: EditionRecord, +): Promise<{ sent: number; outcome: 'delivered' | 'failed' }> { const content = edition.content; const editionDate = edition.edition_date.toISOString().split('T')[0]; const subject = config.generateSubject(content); @@ -73,7 +98,10 @@ export async function sendNewsletter( // Mark as sent if (stats.email_count > 0 || stats.slack_count > 0) { - await config.db.markSent(edition.id, stats); + const markedSent = await config.db.markSent(edition.id, stats); + if (!markedSent) { + throw new Error(`Failed to finalize ${config.id} edition ${edition.id} after delivery`); + } // Publish as perspective (non-blocking) publishAsPerspective(config, edition.id, content, editionDate, subject).catch((err) => { @@ -96,7 +124,10 @@ export async function sendNewsletter( }); } - return { sent: stats.email_count }; + return { + sent: stats.email_count, + outcome: stats.email_count > 0 || stats.slack_count > 0 ? 'delivered' : 'failed', + }; } // ─── Perspective Publishing ──────────────────────────────────────────── diff --git a/server/tests/unit/newsletter-admin-send-now.test.ts b/server/tests/unit/newsletter-admin-send-now.test.ts new file mode 100644 index 0000000000..1a4fed247a --- /dev/null +++ b/server/tests/unit/newsletter-admin-send-now.test.ts @@ -0,0 +1,151 @@ +import { readFile } from 'node:fs/promises'; +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import express from 'express'; +import request from 'supertest'; +import type { EditionRecord, NewsletterConfig } from '../../src/newsletters/config.js'; + +const mocks = vi.hoisted(() => ({ + sendNewsletter: vi.fn(), +})); + +vi.mock('../../src/middleware/auth.js', () => ({ + requireAuth: (req: express.Request, _res: express.Response, next: express.NextFunction) => { + req.user = { id: 'admin_01', email: 'admin@example.test', is_admin: true } as typeof req.user; + next(); + }, + requireAdmin: (_req: express.Request, _res: express.Response, next: express.NextFunction) => next(), +})); + +vi.mock('../../src/newsletters/send-pipeline.js', () => ({ + sendNewsletter: mocks.sendNewsletter, +})); + +vi.mock('../../src/db/client.js', () => ({ + query: vi.fn(), +})); + +function edition(status: EditionRecord['status']): EditionRecord { + return { + id: 42, + edition_date: new Date('2026-08-18T00:00:00.000Z'), + status, + content: { emailSubject: 'Test edition' }, + approved_by: status === 'draft' ? null : 'admin@example.test', + approved_at: status === 'draft' ? null : new Date('2026-08-18T12:00:00.000Z'), + review_channel_id: null, + review_message_ts: null, + perspective_id: null, + created_at: new Date('2026-08-18T10:00:00.000Z'), + sent_at: status === 'sent' ? new Date('2026-08-18T12:05:00.000Z') : null, + send_stats: null, + }; +} + +function makeConfig(current: EditionRecord) { + const approved = edition('approved'); + const sent = edition('sent'); + const db = { + getCurrent: vi.fn() + .mockResolvedValueOnce(current) + .mockResolvedValueOnce(sent), + approve: vi.fn().mockResolvedValue(approved), + }; + const config = { + id: 'the_prompt', + name: 'The Prompt', + author: 'Addie', + palette: { primary: '#000', light: '#fff', dark: '#111' }, + editableFields: [], + cadence: { generateHourET: 8, sendHourET: 10, shouldRunToday: () => true }, + sections: [], + db, + generateSubject: () => 'Test edition', + } as unknown as NewsletterConfig; + return { config, db, approved }; +} + +async function makeApp(config: NewsletterConfig) { + const { createNewsletterAdminRoutes } = await import('../../src/newsletters/admin-routes.js'); + const app = express(); + app.use(express.json()); + app.use('/api/admin/newsletters/the_prompt', createNewsletterAdminRoutes(config)); + return app; +} + +describe('newsletter admin send now', () => { + beforeEach(() => { + vi.clearAllMocks(); + mocks.sendNewsletter.mockResolvedValue({ sent: 12, outcome: 'delivered' }); + }); + + it('approves a draft before sending it immediately', async () => { + const { config, db, approved } = makeConfig(edition('draft')); + const app = await makeApp(config); + + const response = await request(app) + .post('/api/admin/newsletters/the_prompt/editions/42/send-now'); + + expect(response.status).toBe(200); + expect(db.approve).toHaveBeenCalledWith(42, 'admin@example.test'); + expect(mocks.sendNewsletter).toHaveBeenCalledWith(config, approved); + expect(response.body.result).toEqual({ sent: 12, outcome: 'delivered' }); + }); + + it('sends an already approved edition without approving it again', async () => { + const current = edition('approved'); + const { config, db } = makeConfig(current); + const app = await makeApp(config); + + const response = await request(app) + .post('/api/admin/newsletters/the_prompt/editions/42/send-now'); + + expect(response.status).toBe(200); + expect(db.approve).not.toHaveBeenCalled(); + expect(mocks.sendNewsletter).toHaveBeenCalledWith(config, current); + }); + + it('does not send when draft approval loses a status race', async () => { + const { config, db } = makeConfig(edition('draft')); + db.approve.mockResolvedValueOnce(null); + const app = await makeApp(config); + + const response = await request(app) + .post('/api/admin/newsletters/the_prompt/editions/42/send-now'); + + expect(response.status).toBe(409); + expect(response.body.error).toContain('status changed'); + expect(mocks.sendNewsletter).not.toHaveBeenCalled(); + }); + + it('does not resend a completed edition', async () => { + const { config, db } = makeConfig(edition('sent')); + const app = await makeApp(config); + + const response = await request(app) + .post('/api/admin/newsletters/the_prompt/editions/42/send-now'); + + expect(response.status).toBe(400); + expect(db.approve).not.toHaveBeenCalled(); + expect(mocks.sendNewsletter).not.toHaveBeenCalled(); + }); + + it('reports a conflict when another sender holds the edition lock', async () => { + const { config } = makeConfig(edition('approved')); + mocks.sendNewsletter.mockResolvedValueOnce({ sent: 0, outcome: 'busy' }); + const app = await makeApp(config); + + const response = await request(app) + .post('/api/admin/newsletters/the_prompt/editions/42/send-now'); + + expect(response.status).toBe(409); + expect(response.body.error).toContain('already being sent'); + }); + + it('shows the immediate-send action while the edition is still a draft', async () => { + const html = await readFile(new URL('../../public/admin-newsletter.html', import.meta.url), 'utf8'); + + expect(html).toContain("status === 'draft' ? `` : ''"); + expect(html).toContain('Approve this draft and send it now?'); + expect(html).toContain('btn.textContent = originalButtonText'); + }); +}); diff --git a/server/tests/unit/newsletter-send-lock.test.ts b/server/tests/unit/newsletter-send-lock.test.ts new file mode 100644 index 0000000000..2cf35d9abc --- /dev/null +++ b/server/tests/unit/newsletter-send-lock.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it, vi } from 'vitest'; + +const mocks = vi.hoisted(() => ({ + connect: vi.fn(), +})); + +vi.mock('../../src/db/client.js', () => ({ + getPool: () => ({ connect: mocks.connect }), +})); + +function fakeClient(acquired: boolean) { + return { + query: vi.fn() + .mockResolvedValueOnce({ rows: [{ acquired }] }) + .mockResolvedValueOnce({ rows: [{ unlocked: true }] }), + release: vi.fn(), + }; +} + +describe('newsletter send lock', () => { + it('allows only one sender to enter delivery for an edition', async () => { + const firstClient = fakeClient(true); + const secondClient = fakeClient(false); + mocks.connect + .mockResolvedValueOnce(firstClient) + .mockResolvedValueOnce(secondClient); + + let releaseFirst!: () => void; + const firstCanFinish = new Promise((resolve) => { releaseFirst = resolve; }); + let firstStarted!: () => void; + const firstDidStart = new Promise((resolve) => { firstStarted = resolve; }); + const firstWork = vi.fn(async () => { + firstStarted(); + await firstCanFinish; + return 'sent'; + }); + const secondWork = vi.fn(async () => 'duplicate'); + + const { withNewsletterSendLock } = await import('../../src/newsletters/send-lock.js'); + const first = withNewsletterSendLock('the_prompt', 42, firstWork); + await firstDidStart; + const second = await withNewsletterSendLock('the_prompt', 42, secondWork); + + expect(second).toEqual({ acquired: false }); + expect(secondWork).not.toHaveBeenCalled(); + + releaseFirst(); + await expect(first).resolves.toEqual({ acquired: true, value: 'sent' }); + expect(firstClient.query).toHaveBeenLastCalledWith( + 'SELECT pg_advisory_unlock(hashtextextended($1, 0)) AS unlocked', + ['newsletter-send:the_prompt:42'], + ); + expect(firstClient.release).toHaveBeenCalledWith(false); + expect(secondClient.release).toHaveBeenCalledWith(false); + }); +}); diff --git a/server/tests/unit/weekly-digest-job.test.ts b/server/tests/unit/weekly-digest-job.test.ts index 418c36207c..07bb182396 100644 --- a/server/tests/unit/weekly-digest-job.test.ts +++ b/server/tests/unit/weekly-digest-job.test.ts @@ -9,6 +9,10 @@ const getUserWorkingGroupMap = vi.fn(); const sendTrackedBatchMarketingEmails = vi.fn(); const buildPromptMarkdown = vi.fn(() => 'Rendered Prompt markdown'); const publishDigestAsPerspective = vi.fn(async () => undefined); +const withNewsletterSendLock = vi.fn(async (_newsletterId: string, _editionId: number, work: () => Promise) => ({ + acquired: true as const, + value: await work(), +})); vi.mock('../../src/db/digest-db.js', async (importOriginal) => { const actual = await importOriginal(); @@ -73,6 +77,10 @@ vi.mock('../../src/newsletters/cover.js', () => ({ generateCoverForEdition: vi.fn(), })); +vi.mock('../../src/newsletters/send-lock.js', () => ({ + withNewsletterSendLock, +})); + vi.mock('../../src/db/newsletter-suggestions-db.js', () => ({ markSuggestionsIncluded: vi.fn(), })); @@ -107,6 +115,10 @@ function approvedDigest(overrides: Partial = {}): DigestRecord { describe('runWeeklyDigestJob', () => { beforeEach(() => { vi.clearAllMocks(); + withNewsletterSendLock.mockImplementation(async (_newsletterId: string, _editionId: number, work: () => Promise) => ({ + acquired: true as const, + value: await work(), + })); getDigestEmailRecipients.mockResolvedValue([ { workos_user_id: 'user_123', @@ -126,6 +138,7 @@ describe('runWeeklyDigestJob', () => { getUserWorkingGroupMap.mockResolvedValue(new Map()); sendTrackedBatchMarketingEmails.mockResolvedValue({ sent: 1, skipped: 0, failed: 0 }); markSent.mockResolvedValue(true); + getDigestByDate.mockResolvedValue(approvedDigest()); }); it('sends an older approved digest even when today is not a cadence day', async () => { @@ -147,6 +160,19 @@ describe('runWeeklyDigestJob', () => { 'Rendered Prompt markdown', ); }); - expect(getDigestByDate).not.toHaveBeenCalled(); + expect(getDigestByDate).toHaveBeenCalledWith('2026-06-05'); + expect(withNewsletterSendLock).toHaveBeenCalledWith('the_prompt', 123, expect.any(Function)); + }); + + it('does not deliver when a manual sender holds the edition lock', async () => { + getLatestApprovedDigest.mockResolvedValue(approvedDigest()); + withNewsletterSendLock.mockResolvedValueOnce({ acquired: false }); + + const { runWeeklyDigestJob } = await import('../../src/addie/jobs/weekly-digest.js'); + const result = await runWeeklyDigestJob(); + + expect(result.sent).toBe(0); + expect(sendTrackedBatchMarketingEmails).not.toHaveBeenCalled(); + expect(markSent).not.toHaveBeenCalled(); }); });