diff --git a/README.md b/README.md index 74fc8cba..7afbd998 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ NestJS PostgreSQL Apache 2.0 License -1578 Tests Passing 12 AI Agents + 1591 Tests Passing 12 AI Agents

@@ -510,7 +510,7 @@ Operator notes for activating existing adapters, metasearch landings on the dire | OTA Channels | Booking.com + Expedia (EQC) + SiteMinder + DerbySoft | Direct + aggregated OTA connectivity (ARI + content) | | XML Processing | fast-xml-parser | Booking.com OTA XML protocol | | Package Manager | pnpm workspaces | Monorepo management | -| Testing | Vitest (1578 tests across 219 test files) | Unit and integration tests || Build | tsup (packages) + Vite (dashboard) + nest build (API) | Fast builds | +| Testing | Vitest (1591 tests across 220 test files) | Unit and integration tests || Build | tsup (packages) + Vite (dashboard) + nest build (API) | Fast builds | | Containers | Docker + docker-compose | Local dev and production deployment | | CI/CD | GitHub Actions | Automated testing, builds, and releases | @@ -642,7 +642,7 @@ Before going live, verify the items in [`docs/deployment.md`](./docs/deployment. ### Run tests ```bash -# All tests (1578 tests across 219 test files) +# All tests (1591 tests across 220 test files) # API tests only pnpm --filter @telivityhaip/api test @@ -1190,7 +1190,7 @@ HAIP is built in public and contributions are welcome. pnpm install # Install dependencies pnpm build # Build all workspace packages pnpm dev # Start API in dev mode (hot reload) -pnpm test # Run all tests (1578 tests, 219 files) +pnpm test # Run all tests (1591 tests, 220 files) pnpm lint # ESLint ``` diff --git a/apps/api/package.json b/apps/api/package.json index a1ed495b..67ea9504 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -42,6 +42,7 @@ "fast-xml-parser": "^5.5.10", "jsonwebtoken": "^9.0.3", "jwks-rsa": "^4.0.1", + "nodemailer": "^9.0.5", "passport": "^0.7.0", "passport-jwt": "^4.0.1", "postgres": "^3.4.5", diff --git a/apps/api/src/modules/agent/guest-comms/email-provider.interface.ts b/apps/api/src/modules/agent/guest-comms/email-provider.interface.ts index dbca9dda..483aa004 100644 --- a/apps/api/src/modules/agent/guest-comms/email-provider.interface.ts +++ b/apps/api/src/modules/agent/guest-comms/email-provider.interface.ts @@ -4,19 +4,51 @@ export interface EmailMessage { html: string; text: string; from?: string; + /** + * Correlation key forwarded to providers as custom metadata + * (e.g. Mailgun `v:haip-idempotency-key`, SendGrid custom_args). Useful for + * log/trace correlation across retries — NOT an exactly-once or deduplication + * guarantee in Mailgun, SendGrid, SES, or SMTP. + */ + idempotencyKey?: string; + /** + * Stable RFC Message-ID reused across retries for correlation when the + * provider supports setting it. Does not prevent duplicate delivery. + */ + messageId?: string; } +/** Provider-confirmed acceptance vs definite failure vs ambiguous response. */ +export type EmailDeliveryStatus = 'sent' | 'notSent' | 'outcomeUnknown'; + export interface EmailResult { + /** `sent` = provider confirmed; `notSent` = safe to auto-retry; `outcomeUnknown` = do not auto-retry. */ + status: EmailDeliveryStatus; + /** Convenience mirror of `status === 'sent'`. */ sent: boolean; messageId?: string; provider?: string; error?: string; } +export interface EmailSendOptions { + /** + * Hard send deadline in milliseconds. HTTP transports return at the deadline + * even when the underlying fetch ignores abort; SMTP hard-closes owned sockets + * and returns via `Promise.race` even if `sendMail` is still settling. + */ + timeoutMs?: number; + /** + * Max send attempts for definitely-not-sent failures (`status: 'notSent'`). + * Does not retry `outcomeUnknown` (ambiguous acceptance). + */ + maxAttempts?: number; +} + export interface EmailProvider { readonly name: string; isConfigured(): boolean; - send(message: EmailMessage): Promise; + send(message: EmailMessage, options?: EmailSendOptions): Promise; } export const EMAIL_PROVIDERS = Symbol('EMAIL_PROVIDERS'); diff --git a/apps/api/src/modules/agent/guest-comms/email.service.spec.ts b/apps/api/src/modules/agent/guest-comms/email.service.spec.ts index f7d5fc80..71f457e0 100644 --- a/apps/api/src/modules/agent/guest-comms/email.service.spec.ts +++ b/apps/api/src/modules/agent/guest-comms/email.service.spec.ts @@ -1,12 +1,17 @@ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; import { EmailService } from './email.service'; -import type { EmailProvider } from './email-provider.interface'; +import type { EmailProvider, EmailResult } from './email-provider.interface'; describe('EmailService', () => { const consoleProvider: EmailProvider = { name: 'console', isConfigured: () => true, - send: vi.fn().mockResolvedValue({ sent: false, provider: 'console', error: 'logged' }), + send: vi.fn().mockResolvedValue({ + status: 'notSent', + sent: false, + provider: 'console', + error: 'logged', + } satisfies EmailResult), }; beforeEach(() => { @@ -17,7 +22,12 @@ describe('EmailService', () => { const sendgrid = { name: 'sendgrid', isConfigured: () => true, - send: vi.fn().mockResolvedValue({ sent: true, provider: 'sendgrid', messageId: 'sg-1' }), + send: vi.fn().mockResolvedValue({ + status: 'sent', + sent: true, + provider: 'sendgrid', + messageId: 'sg-1', + } satisfies EmailResult), }; const smtp = { name: 'smtp', @@ -32,10 +42,37 @@ describe('EmailService', () => { text: 'Hi', }); expect(result.sent).toBe(true); - expect(sendgrid.send).toHaveBeenCalled(); + expect(result.status).toBe('sent'); + expect(sendgrid.send).toHaveBeenCalledTimes(1); expect(smtp.send).not.toHaveBeenCalled(); }); + it('passes stable transport identity through to the selected provider', async () => { + const provider = { + name: 'sendgrid', + isConfigured: () => true, + send: vi.fn().mockResolvedValue({ + status: 'sent', + sent: true, + messageId: 'provider-id', + } satisfies EmailResult), + }; + const service = new EmailService([provider]); + const message = { + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + idempotencyKey: 'booking-request-email:delivery-1', + messageId: '', + }; + await service.send(message, { timeoutMs: 1_234 }); + expect(provider.send).toHaveBeenCalledWith(expect.objectContaining({ + idempotencyKey: 'booking-request-email:delivery-1', + messageId: '', + }), { timeoutMs: 1_234 }); + }); + it('falls back to console when no real provider is configured', async () => { const smtp = { name: 'smtp', isConfigured: () => false, send: vi.fn() }; const sendgrid = { name: 'sendgrid', isConfigured: () => false, send: vi.fn() }; @@ -49,6 +86,87 @@ describe('EmailService', () => { }); expect(consoleProvider.send).toHaveBeenCalled(); }); + + it('retries definitely-not-sent failures but not outcomeUnknown', async () => { + const provider = { + name: 'sendgrid', + isConfigured: () => true, + send: vi.fn() + .mockResolvedValueOnce({ + status: 'notSent', + sent: false, + provider: 'sendgrid', + error: 'SendGrid HTTP 503', + } satisfies EmailResult) + .mockResolvedValueOnce({ + status: 'outcomeUnknown', + sent: false, + provider: 'sendgrid', + error: 'Email transport timed out', + } satisfies EmailResult), + }; + const service = new EmailService([provider]); + const message = { + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + idempotencyKey: 'delivery-1', + messageId: '', + }; + + const result = await service.send(message, { maxAttempts: 3 }); + expect(result.status).toBe('outcomeUnknown'); + expect(provider.send).toHaveBeenCalledTimes(2); + }); + + it('does not retry outcomeUnknown on the first attempt', async () => { + const provider = { + name: 'sendgrid', + isConfigured: () => true, + send: vi.fn().mockResolvedValue({ + status: 'outcomeUnknown', + sent: false, + provider: 'sendgrid', + error: 'Email transport timed out', + } satisfies EmailResult), + }; + const service = new EmailService([provider]); + const result = await service.send({ + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + }, { maxAttempts: 3 }); + expect(result.status).toBe('outcomeUnknown'); + expect(provider.send).toHaveBeenCalledTimes(1); + }); + + it('retries up to maxAttempts for notSent then stops', async () => { + vi.useFakeTimers(); + const provider = { + name: 'sendgrid', + isConfigured: () => true, + send: vi.fn().mockResolvedValue({ + status: 'notSent', + sent: false, + provider: 'sendgrid', + error: 'SendGrid HTTP 503', + } satisfies EmailResult), + }; + const service = new EmailService([provider]); + const sending = service.send({ + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + }, { maxAttempts: 3 }); + await vi.runAllTimersAsync(); + const result = await sending; + expect(result.status).toBe('notSent'); + expect(provider.send).toHaveBeenCalledTimes(3); + vi.useRealTimers(); + }); }); describe('SendgridEmailProvider', () => { @@ -56,6 +174,7 @@ describe('SendgridEmailProvider', () => { const originalEnv = { ...process.env }; afterEach(() => { + vi.useRealTimers(); global.fetch = originalFetch; process.env = { ...originalEnv }; vi.resetModules(); @@ -73,6 +192,7 @@ describe('SendgridEmailProvider', () => { html: 'h', text: 't', }); + expect(result.status).toBe('notSent'); expect(result.sent).toBe(false); }); @@ -92,11 +212,117 @@ describe('SendgridEmailProvider', () => { subject: 'Confirm', html: '

Hi

', text: 'Hi', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); + expect(result.status).toBe('sent'); expect(result.sent).toBe(true); expect(global.fetch).toHaveBeenCalledWith( 'https://api.sendgrid.com/v3/mail/send', expect.objectContaining({ method: 'POST' }), ); + const init = vi.mocked(global.fetch).mock.calls[0]?.[1]; + const payload = JSON.parse(String(init?.body)); + expect(payload.personalizations[0]).toMatchObject({ + headers: { 'Message-ID': '' }, + custom_args: { haip_idempotency_key: 'stable-delivery-1' }, + }); + }); + + it('returns at the hard HTTP deadline even when fetch ignores abort', async () => { + vi.useFakeTimers(); + process.env['SENDGRID_API_KEY'] = 'SG.test'; + process.env['SENDGRID_FROM'] = 'hotel@example.com'; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, _reject) => { + init?.signal?.addEventListener('abort', () => undefined, { once: true }); + })) as any; + + const { SendgridEmailProvider } = await import('./providers/sendgrid-email.provider'); + const provider = new SendgridEmailProvider(); + const sending = provider.send({ + to: 'guest@example.com', + subject: 'Confirm', + html: '

Hi

', + text: 'Hi', + }, { timeoutMs: 100 }); + await vi.waitFor(() => expect(global.fetch).toHaveBeenCalledOnce()); + const signal = vi.mocked(global.fetch).mock.calls[0]?.[1]?.signal; + expect(signal).toBeInstanceOf(AbortSignal); + + await vi.advanceTimersByTimeAsync(100); + await expect(sending).resolves.toMatchObject({ + status: 'outcomeUnknown', + sent: false, + error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); + }); + + it('swallows detached fetch rejection after the hard deadline', async () => { + vi.useFakeTimers(); + process.env['SENDGRID_API_KEY'] = 'SG.test'; + process.env['SENDGRID_FROM'] = 'hotel@example.com'; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => { + queueMicrotask(() => reject(Object.assign(new Error('aborted'), { name: 'AbortError' }))); + }, { once: true }); + })) as any; + + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown) => unhandled.push(reason); + process.on('unhandledRejection', onUnhandled); + + const { SendgridEmailProvider } = await import('./providers/sendgrid-email.provider'); + const provider = new SendgridEmailProvider(); + const sending = provider.send({ + to: 'guest@example.com', + subject: 'Confirm', + html: '

Hi

', + text: 'Hi', + }, { timeoutMs: 100 }); + await vi.waitFor(() => expect(global.fetch).toHaveBeenCalledOnce()); + await vi.advanceTimersByTimeAsync(100); + await expect(sending).resolves.toMatchObject({ status: 'outcomeUnknown' }); + await vi.advanceTimersByTimeAsync(100); + await Promise.resolve(); + expect(unhandled).toEqual([]); + process.off('unhandledRejection', onUnhandled); + }); +}); + +describe('boundedEmailFetch', () => { + const originalFetch = global.fetch; + + afterEach(() => { + vi.useRealTimers(); + global.fetch = originalFetch; + vi.resetModules(); + }); + + it('returns at timeout without waiting for a hanging response body', async () => { + vi.useFakeTimers(); + let bodySettled = false; + global.fetch = vi.fn((_url, init) => Promise.resolve({ + ok: true, + json: () => new Promise((_resolve, _reject) => { + init?.signal?.addEventListener('abort', () => undefined, { once: true }); + }), + })) as any; + + const { boundedEmailFetch, EmailTransportTimeoutError } = await import('./providers/bounded-email-transport'); + const work = boundedEmailFetch( + 'https://example.test/send', + { method: 'POST' }, + { timeoutMs: 100 }, + async (response) => { + await response.json(); + bodySettled = true; + return 'done'; + }, + ); + const expectation = expect(work).rejects.toBeInstanceOf(EmailTransportTimeoutError); + await vi.advanceTimersByTimeAsync(100); + await expectation; + expect(bodySettled).toBe(false); }); }); diff --git a/apps/api/src/modules/agent/guest-comms/email.service.ts b/apps/api/src/modules/agent/guest-comms/email.service.ts index e6e6689c..94320659 100644 --- a/apps/api/src/modules/agent/guest-comms/email.service.ts +++ b/apps/api/src/modules/agent/guest-comms/email.service.ts @@ -1,13 +1,23 @@ import { Inject, Injectable, Logger } from '@nestjs/common'; -import type { EmailMessage, EmailProvider, EmailResult } from './email-provider.interface'; +import type { + EmailMessage, + EmailProvider, + EmailResult, + EmailSendOptions, +} from './email-provider.interface'; import { EMAIL_PROVIDERS } from './email-provider.interface'; -export type { EmailMessage, EmailResult } from './email-provider.interface'; +export type { EmailMessage, EmailResult, EmailSendOptions } from './email-provider.interface'; + +const DEFAULT_EMAIL_SEND_MAX_ATTEMPTS = 3; +const EMAIL_SEND_RETRY_BASE_DELAY_MS = 250; /** * Email transport service — SendGrid, Mailgun, SES gateway, SMTP, or console fallback. * * Provider order: first configured among SendGrid → Mailgun → SES → SMTP, else console. + * Automatically retries definitely-not-sent failures only; never auto-retries + * `outcomeUnknown` (provider may have accepted mail before the response was lost). */ @Injectable() export class EmailService { @@ -19,18 +29,43 @@ export class EmailService { return this.providers.some((p) => p.name !== 'console' && p.isConfigured()); } - async send(message: EmailMessage): Promise { + async send(message: EmailMessage, options?: EmailSendOptions): Promise { const provider = this.activeProvider(); - const result = await provider.send(message); - if (!result.provider) { - return { ...result, provider: provider.name }; + const maxAttempts = Math.max(1, Math.floor(options?.maxAttempts ?? DEFAULT_EMAIL_SEND_MAX_ATTEMPTS)); + let lastResult: EmailResult | undefined; + + for (let attempt = 1; attempt <= maxAttempts; attempt++) { + const raw = await provider.send(message, options); + lastResult = raw.provider ? raw : { ...raw, provider: provider.name }; + + if (lastResult.status !== 'notSent' || attempt === maxAttempts) { + break; + } + + this.logger.warn( + `Email to ${message.to} not delivered via ${lastResult.provider} (attempt ${attempt}/${maxAttempts}): ${lastResult.error}; retrying`, + ); + await this.retryDelay(attempt); } - if (!result.sent) { + + const result = lastResult!; + if (result.status === 'notSent') { this.logger.warn(`Email to ${message.to} not delivered via ${result.provider}: ${result.error}`); + } else if (result.status === 'outcomeUnknown') { + this.logger.warn( + `Email to ${message.to} outcome unknown via ${result.provider}: ${result.error} — not auto-retrying`, + ); } return result; } + private async retryDelay(attempt: number): Promise { + await new Promise((resolve) => { + const timer = setTimeout(resolve, EMAIL_SEND_RETRY_BASE_DELAY_MS * attempt); + timer.unref?.(); + }); + } + private activeProvider(): EmailProvider { const real = this.providers.find((p) => p.name !== 'console' && p.isConfigured()); if (real) return real; diff --git a/apps/api/src/modules/agent/guest-comms/guest-communication.agent.spec.ts b/apps/api/src/modules/agent/guest-comms/guest-communication.agent.spec.ts index 3787d715..4a36a8c6 100644 --- a/apps/api/src/modules/agent/guest-comms/guest-communication.agent.spec.ts +++ b/apps/api/src/modules/agent/guest-comms/guest-communication.agent.spec.ts @@ -87,7 +87,7 @@ describe('GuestCommunicationAgent', () => { it('execute sends email when provider is configured', async () => { emailService.isConfigured.mockReturnValue(true); - emailService.send.mockResolvedValue({ sent: true, provider: 'sendgrid' }); + emailService.send.mockResolvedValue({ status: 'sent', sent: true, provider: 'sendgrid' }); const result = await agent.execute({ recommendation: { diff --git a/apps/api/src/modules/agent/guest-comms/providers/bounded-email-transport.ts b/apps/api/src/modules/agent/guest-comms/providers/bounded-email-transport.ts new file mode 100644 index 00000000..982cbfde --- /dev/null +++ b/apps/api/src/modules/agent/guest-comms/providers/bounded-email-transport.ts @@ -0,0 +1,79 @@ +import type { EmailResult, EmailSendOptions } from '../email-provider.interface'; + +export const DEFAULT_EMAIL_SEND_TIMEOUT_MS = 60_000; + +export class EmailTransportTimeoutError extends Error { + constructor() { + super('Email transport timed out'); + this.name = 'EmailTransportTimeoutError'; + } +} + +export function emailSendTimeoutMs(options?: EmailSendOptions): number { + const requested = options?.timeoutMs; + if (!Number.isFinite(requested) || requested === undefined) { + return DEFAULT_EMAIL_SEND_TIMEOUT_MS; + } + return Math.max(1, Math.floor(requested)); +} + +export function sentEmailResult( + provider: string, + messageId?: string, +): EmailResult { + return { status: 'sent', sent: true, provider, messageId }; +} + +export function notSentEmailResult( + provider: string, + error?: string, +): EmailResult { + return { status: 'notSent', sent: false, provider, error }; +} + +export function unknownTimeoutResult(provider: string): EmailResult { + return { + status: 'outcomeUnknown', + sent: false, + provider, + error: 'Email transport timed out', + }; +} + +/** + * Hard outer deadline for HTTP email sends, including response-body consumption. + * Returns at `timeoutMs` even when fetch/abort is ignored; in-flight work + * continues detached with a rejection handler so callers never see an + * unhandled rejection. + */ +export async function boundedEmailFetch( + input: string, + init: RequestInit, + options: EmailSendOptions | undefined, + consume: (response: Response) => Promise | T, +): Promise { + const controller = new AbortController(); + const timeoutMs = emailSendTimeoutMs(options); + + const work = (async (): Promise => { + const response = await fetch(input, { ...init, signal: controller.signal }); + return await consume(response); + })(); + + work.catch(() => undefined); + + let timer: ReturnType | undefined; + const deadline = new Promise((_, reject) => { + timer = setTimeout(() => { + controller.abort(); + reject(new EmailTransportTimeoutError()); + }, timeoutMs); + timer.unref?.(); + }); + + try { + return await Promise.race([work, deadline]); + } finally { + if (timer) clearTimeout(timer); + } +} diff --git a/apps/api/src/modules/agent/guest-comms/providers/console-email.provider.ts b/apps/api/src/modules/agent/guest-comms/providers/console-email.provider.ts index 9821af49..f9a4bc4a 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/console-email.provider.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/console-email.provider.ts @@ -1,5 +1,6 @@ import { Injectable, Logger } from '@nestjs/common'; import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider.interface'; +import { notSentEmailResult } from './bounded-email-transport'; /** * Development fallback — logs the message instead of sending. @@ -18,10 +19,8 @@ export class ConsoleEmailProvider implements EmailProvider { `[Email:console] → ${message.to} | ${message.subject}\n${message.text.slice(0, 200)}`, ); return { - sent: false, - provider: this.name, - messageId: `console-${Date.now()}`, - error: 'No email provider configured — message logged only', + ...notSentEmailResult(this.name, 'No email provider configured — message logged only'), + messageId: message.messageId ?? `console-${Date.now()}`, }; } } diff --git a/apps/api/src/modules/agent/guest-comms/providers/mailgun-email.provider.ts b/apps/api/src/modules/agent/guest-comms/providers/mailgun-email.provider.ts index 2a395095..c216d74d 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/mailgun-email.provider.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/mailgun-email.provider.ts @@ -1,5 +1,17 @@ import { Injectable, Logger } from '@nestjs/common'; -import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider.interface'; +import type { + EmailMessage, + EmailProvider, + EmailResult, + EmailSendOptions, +} from '../email-provider.interface'; +import { + boundedEmailFetch, + EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * Mailgun Messages API adapter. @@ -24,9 +36,9 @@ export class MailgunEmailProvider implements EmailProvider { return Boolean(this.apiKey && this.domain && this.defaultFrom); } - async send(message: EmailMessage): Promise { + async send(message: EmailMessage, options?: EmailSendOptions): Promise { if (!this.isConfigured()) { - return { sent: false, provider: this.name, error: 'Mailgun not configured' }; + return notSentEmailResult(this.name, 'Mailgun not configured'); } const form = new URLSearchParams(); @@ -35,33 +47,43 @@ export class MailgunEmailProvider implements EmailProvider { form.set('subject', message.subject); form.set('text', message.text); form.set('html', message.html); + if (message.messageId) form.set('h:Message-Id', message.messageId); + if (message.idempotencyKey) { + form.set('v:haip-idempotency-key', message.idempotencyKey); + } try { const auth = Buffer.from(`api:${this.apiKey}`).toString('base64'); - const res = await fetch(`${this.apiBase}/v3/${this.domain}/messages`, { - method: 'POST', - headers: { - Authorization: `Basic ${auth}`, - 'Content-Type': 'application/x-www-form-urlencoded', + const { response: res, payload } = await boundedEmailFetch( + `${this.apiBase}/v3/${this.domain}/messages`, + { + method: 'POST', + headers: { + Authorization: `Basic ${auth}`, + 'Content-Type': 'application/x-www-form-urlencoded', + }, + body: form.toString(), }, - body: form.toString(), - }); - const payload = (await res.json().catch(() => ({}))) as { - id?: string; - message?: string; - }; + options, + async (response) => ({ + response, + payload: (await response.json().catch(() => ({}))) as { + id?: string; + message?: string; + }, + }), + ); if (!res.ok) { - return { - sent: false, - provider: this.name, - error: payload.message ?? `Mailgun HTTP ${res.status}`, - }; + return notSentEmailResult(this.name, payload.message ?? `Mailgun HTTP ${res.status}`); } this.logger.log(`Email sent via Mailgun to ${message.to}`); - return { sent: true, provider: this.name, messageId: payload.id }; + return sentEmailResult(this.name, payload.id); } catch (error: any) { + if (error instanceof EmailTransportTimeoutError) { + return unknownTimeoutResult(this.name); + } this.logger.error(`Mailgun send failed: ${error.message}`); - return { sent: false, provider: this.name, error: error.message }; + return notSentEmailResult(this.name, error.message); } } } diff --git a/apps/api/src/modules/agent/guest-comms/providers/mailgun-ses.provider.spec.ts b/apps/api/src/modules/agent/guest-comms/providers/mailgun-ses.provider.spec.ts index c430236f..eda9119c 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/mailgun-ses.provider.spec.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/mailgun-ses.provider.spec.ts @@ -5,6 +5,7 @@ describe('MailgunEmailProvider', () => { const originalEnv = { ...process.env }; afterEach(() => { + vi.useRealTimers(); global.fetch = originalFetch; process.env = { ...originalEnv }; vi.resetModules(); @@ -33,9 +34,73 @@ describe('MailgunEmailProvider', () => { subject: 'S', html: 'h', text: 't', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); + expect(result.status).toBe('sent'); expect(result.sent).toBe(true); expect(result.messageId).toBe(''); + const init = vi.mocked(global.fetch).mock.calls[0]?.[1]; + const form = new URLSearchParams(String(init?.body)); + expect(form.get('h:Message-Id')).toBe(''); + expect(form.get('v:haip-idempotency-key')).toBe('stable-delivery-1'); + }); + + it('returns at the hard HTTP deadline even when fetch ignores abort', async () => { + vi.useFakeTimers(); + process.env['MAILGUN_API_KEY'] = 'key'; + process.env['MAILGUN_DOMAIN'] = 'mg.example.com'; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, _reject) => { + init?.signal?.addEventListener('abort', () => undefined, { once: true }); + })) as any; + + const { MailgunEmailProvider } = await import('./mailgun-email.provider'); + const provider = new MailgunEmailProvider(); + const sending = provider.send({ + to: 'a@b.com', subject: 'S', html: 'h', text: 't', + }, { timeoutMs: 100 }); + await vi.waitFor(() => expect(global.fetch).toHaveBeenCalledOnce()); + const signal = vi.mocked(global.fetch).mock.calls[0]?.[1]?.signal; + expect(signal).toBeInstanceOf(AbortSignal); + await vi.advanceTimersByTimeAsync(100); + + await expect(sending).resolves.toMatchObject({ + status: 'outcomeUnknown', + sent: false, + error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); + }); + + it('swallows detached body rejection after the hard deadline', async () => { + vi.useFakeTimers(); + process.env['MAILGUN_API_KEY'] = 'key'; + process.env['MAILGUN_DOMAIN'] = 'mg.example.com'; + global.fetch = vi.fn((_url, init) => Promise.resolve({ + ok: true, + json: () => new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => { + queueMicrotask(() => reject(Object.assign(new Error('aborted body'), { name: 'AbortError' }))); + }, { once: true }); + }), + })) as any; + + const unhandled: unknown[] = []; + const onUnhandled = (reason: unknown) => unhandled.push(reason); + process.on('unhandledRejection', onUnhandled); + + const { MailgunEmailProvider } = await import('./mailgun-email.provider'); + const provider = new MailgunEmailProvider(); + const sending = provider.send({ + to: 'a@b.com', subject: 'S', html: 'h', text: 't', + }, { timeoutMs: 100 }); + await vi.waitFor(() => expect(global.fetch).toHaveBeenCalledOnce()); + await vi.advanceTimersByTimeAsync(100); + await expect(sending).resolves.toMatchObject({ status: 'outcomeUnknown' }); + await vi.advanceTimersByTimeAsync(100); + await Promise.resolve(); + expect(unhandled).toEqual([]); + process.off('unhandledRejection', onUnhandled); }); }); @@ -44,6 +109,7 @@ describe('SesEmailProvider', () => { const originalEnv = { ...process.env }; afterEach(() => { + vi.useRealTimers(); global.fetch = originalFetch; process.env = { ...originalEnv }; vi.resetModules(); @@ -73,11 +139,49 @@ describe('SesEmailProvider', () => { subject: 'S', html: 'h', text: 't', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); expect(result).toEqual({ + status: 'sent', sent: true, provider: 'amazon-ses', messageId: 'ses-1', }); + const init = vi.mocked(global.fetch).mock.calls[0]?.[1]; + expect(init?.headers).toMatchObject({ + 'X-HAIP-Idempotency-Key': 'stable-delivery-1', + }); + const payload = JSON.parse(String(init?.body)); + expect(payload.Content.Simple.Headers).toContainEqual({ + Name: 'X-HAIP-Message-ID', Value: '', + }); + }); + + it('returns at the hard HTTP deadline even when fetch ignores abort', async () => { + vi.useFakeTimers(); + process.env['SES_ENDPOINT'] = 'http://localhost:4566'; + process.env['SES_API_KEY'] = 'local'; + process.env['SES_FROM'] = 'noreply@example.com'; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, _reject) => { + init?.signal?.addEventListener('abort', () => undefined, { once: true }); + })) as any; + + const { SesEmailProvider } = await import('./ses-email.provider'); + const provider = new SesEmailProvider(); + const sending = provider.send({ + to: 'a@b.com', subject: 'S', html: 'h', text: 't', + }, { timeoutMs: 100 }); + await vi.waitFor(() => expect(global.fetch).toHaveBeenCalledOnce()); + const signal = vi.mocked(global.fetch).mock.calls[0]?.[1]?.signal; + expect(signal).toBeInstanceOf(AbortSignal); + await vi.advanceTimersByTimeAsync(100); + + await expect(sending).resolves.toMatchObject({ + status: 'outcomeUnknown', + sent: false, + error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); }); }); diff --git a/apps/api/src/modules/agent/guest-comms/providers/sendgrid-email.provider.ts b/apps/api/src/modules/agent/guest-comms/providers/sendgrid-email.provider.ts index 9bf80336..a018222b 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/sendgrid-email.provider.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/sendgrid-email.provider.ts @@ -1,5 +1,17 @@ import { Injectable, Logger } from '@nestjs/common'; -import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider.interface'; +import type { + EmailMessage, + EmailProvider, + EmailResult, + EmailSendOptions, +} from '../email-provider.interface'; +import { + boundedEmailFetch, + EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * SendGrid Email API reference adapter. @@ -32,14 +44,23 @@ export class SendgridEmailProvider implements EmailProvider { return Boolean(this.apiKey && this.defaultFrom); } - async send(message: EmailMessage): Promise { + async send(message: EmailMessage, options?: EmailSendOptions): Promise { if (!this.isConfigured()) { - return { sent: false, provider: this.name, error: 'SendGrid not configured' }; + return notSentEmailResult(this.name, 'SendGrid not configured'); } const from = message.from ?? this.defaultFrom!; + const personalization: { + to: Array<{ email: string }>; + headers?: Record; + custom_args?: Record; + } = { to: [{ email: message.to }] }; + if (message.messageId) personalization.headers = { 'Message-ID': message.messageId }; + if (message.idempotencyKey) { + personalization.custom_args = { haip_idempotency_key: message.idempotencyKey }; + } const payload = { - personalizations: [{ to: [{ email: message.to }] }], + personalizations: [personalization], from: { email: from }, subject: message.subject, content: [ @@ -49,27 +70,37 @@ export class SendgridEmailProvider implements EmailProvider { }; try { - const res = await fetch('https://api.sendgrid.com/v3/mail/send', { - method: 'POST', - headers: { - Authorization: `Bearer ${this.apiKey}`, - 'Content-Type': 'application/json', + const { response: res, failureBody } = await boundedEmailFetch( + 'https://api.sendgrid.com/v3/mail/send', + { + method: 'POST', + headers: { + Authorization: `Bearer ${this.apiKey}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify(payload), }, - body: JSON.stringify(payload), - }); + options, + async (response) => ({ + response, + failureBody: response.ok ? undefined : await response.text(), + }), + ); if (!res.ok) { - const body = await res.text(); - this.logger.error(`SendGrid send failed (${res.status}): ${body}`); - return { sent: false, provider: this.name, error: `SendGrid HTTP ${res.status}` }; + this.logger.error(`SendGrid send failed (${res.status}): ${failureBody ?? ''}`); + return notSentEmailResult(this.name, `SendGrid HTTP ${res.status}`); } const messageId = res.headers.get('x-message-id') ?? undefined; this.logger.log(`Email sent via SendGrid to ${message.to}`); - return { sent: true, provider: this.name, messageId }; + return sentEmailResult(this.name, messageId); } catch (error: any) { + if (error instanceof EmailTransportTimeoutError) { + return unknownTimeoutResult(this.name); + } this.logger.error(`SendGrid send failed: ${error.message}`); - return { sent: false, provider: this.name, error: error.message }; + return notSentEmailResult(this.name, error.message); } } } diff --git a/apps/api/src/modules/agent/guest-comms/providers/ses-email.provider.ts b/apps/api/src/modules/agent/guest-comms/providers/ses-email.provider.ts index 4763e5f4..6a3a9a51 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/ses-email.provider.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/ses-email.provider.ts @@ -1,5 +1,17 @@ import { Injectable, Logger } from '@nestjs/common'; -import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider.interface'; +import type { + EmailMessage, + EmailProvider, + EmailResult, + EmailSendOptions, +} from '../email-provider.interface'; +import { + boundedEmailFetch, + EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * Amazon SES outbound adapter via an explicit HTTPS gateway. @@ -22,13 +34,12 @@ export class SesEmailProvider implements EmailProvider { return Boolean(this.from && this.endpoint && this.apiKey); } - async send(message: EmailMessage): Promise { + async send(message: EmailMessage, options?: EmailSendOptions): Promise { if (!this.isConfigured()) { - return { - sent: false, - provider: this.name, - error: 'Amazon SES gateway not configured (set SES_ENDPOINT + SES_API_KEY + SES_FROM)', - }; + return notSentEmailResult( + this.name, + 'Amazon SES gateway not configured (set SES_ENDPOINT + SES_API_KEY + SES_FROM)', + ); } const payload = { @@ -37,6 +48,9 @@ export class SesEmailProvider implements EmailProvider { Content: { Simple: { Subject: { Data: message.subject }, + ...(message.messageId + ? { Headers: [{ Name: 'X-HAIP-Message-ID', Value: message.messageId }] } + : {}), Body: { Text: { Data: message.text }, Html: { Data: message.html }, @@ -46,31 +60,40 @@ export class SesEmailProvider implements EmailProvider { }; try { - const res = await fetch(`${this.endpoint!.replace(/\/$/, '')}/v2/email/outbound-emails`, { - method: 'POST', - headers: { - Authorization: `Bearer ${this.apiKey}`, - 'Content-Type': 'application/json', - 'X-SES-Region': this.region, + const { response: res, body } = await boundedEmailFetch( + `${this.endpoint!.replace(/\/$/, '')}/v2/email/outbound-emails`, + { + method: 'POST', + headers: { + Authorization: `Bearer ${this.apiKey}`, + 'Content-Type': 'application/json', + 'X-SES-Region': this.region, + ...(message.idempotencyKey + ? { 'X-HAIP-Idempotency-Key': message.idempotencyKey } + : {}), + }, + body: JSON.stringify(payload), }, - body: JSON.stringify(payload), - }); - const body = (await res.json().catch(() => ({}))) as { - MessageId?: string; - message?: string; - }; + options, + async (response) => ({ + response, + body: (await response.json().catch(() => ({}))) as { + MessageId?: string; + message?: string; + }, + }), + ); if (!res.ok) { - return { - sent: false, - provider: this.name, - error: body.message ?? `SES HTTP ${res.status}`, - }; + return notSentEmailResult(this.name, body.message ?? `SES HTTP ${res.status}`); } this.logger.log(`Email sent via SES gateway to ${message.to}`); - return { sent: true, provider: this.name, messageId: body.MessageId }; + return sentEmailResult(this.name, body.MessageId); } catch (error: any) { + if (error instanceof EmailTransportTimeoutError) { + return unknownTimeoutResult(this.name); + } this.logger.error(`SES send failed: ${error.message}`); - return { sent: false, provider: this.name, error: error.message }; + return notSentEmailResult(this.name, error.message); } } } diff --git a/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.spec.ts b/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.spec.ts new file mode 100644 index 00000000..c539b0f2 --- /dev/null +++ b/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.spec.ts @@ -0,0 +1,218 @@ +import { createServer, type Socket } from 'node:net'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { SmtpEmailProvider } from './smtp-email.provider'; + +describe('SmtpEmailProvider bounded send', () => { + const originalEnv = { ...process.env }; + + afterEach(() => { + process.env = { ...originalEnv }; + }); + + it('settles and closes a connection that never sends its SMTP greeting', async () => { + const sockets = new Set(); + const server = createServer((socket) => { + sockets.add(socket); + socket.on('close', () => sockets.delete(socket)); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('SMTP test server did not bind'); + process.env['SMTP_HOST'] = '127.0.0.1'; + process.env['SMTP_PORT'] = String(address.port); + delete process.env['SMTP_USER']; + delete process.env['SMTP_PASS']; + const provider = new SmtpEmailProvider(); + const didNotSettle = Symbol('did-not-settle'); + + try { + const result = await Promise.race([ + provider.send({ + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + messageId: '', + }, { timeoutMs: 50 }), + new Promise((resolve) => { + setTimeout(() => resolve(didNotSettle), 500); + }), + ]); + + expect(result).not.toBe(didNotSettle); + expect(result).toMatchObject({ + status: 'outcomeUnknown', + sent: false, + provider: 'smtp', + error: 'Email transport timed out', + }); + await vi.waitFor(() => expect(sockets.size).toBe(0), { timeout: 500 }); + } finally { + for (const socket of sockets) socket.destroy(); + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); + } + }); + + it('hard-closes an active SMTP transaction before a later send can connect', async () => { + const sockets = new Set(); + const heartbeats = new Map>(); + let sessionsAtData = 0; + let maxConcurrentConnections = 0; + const server = createServer({ allowHalfOpen: true }, (socket) => { + sockets.add(socket); + maxConcurrentConnections = Math.max(maxConcurrentConnections, sockets.size); + let input = ''; + let readingData = false; + + socket.write('220 smtp.test ESMTP ready\r\n'); + socket.on('data', (chunk) => { + input += chunk.toString('utf8'); + if (readingData) { + const dataEnd = input.indexOf('\r\n.\r\n'); + if (dataEnd < 0) return; + input = input.slice(dataEnd + 5); + readingData = false; + sessionsAtData += 1; + const heartbeat = setInterval(() => { + if (!socket.destroyed && socket.writable) socket.write(' '); + }, 5); + heartbeats.set(socket, heartbeat); + return; + } + + let lineEnd: number; + while ((lineEnd = input.indexOf('\r\n')) >= 0) { + const command = input.slice(0, lineEnd); + input = input.slice(lineEnd + 2); + if (/^EHLO /i.test(command)) { + socket.write('250-smtp.test\r\n250 PIPELINING\r\n'); + } else if (/^MAIL FROM:/i.test(command) || /^RCPT TO:/i.test(command)) { + socket.write('250 2.1.0 Ok\r\n'); + } else if (/^DATA$/i.test(command)) { + readingData = true; + socket.write('354 End data with .\r\n'); + if (input.includes('\r\n.\r\n')) { + socket.emit('data', Buffer.alloc(0)); + } + break; + } else if (/^QUIT$/i.test(command)) { + socket.write('221 2.0.0 Bye\r\n'); + socket.end(); + } + } + }); + socket.on('error', () => undefined); + socket.on('close', () => { + sockets.delete(socket); + const heartbeat = heartbeats.get(socket); + if (heartbeat) clearInterval(heartbeat); + heartbeats.delete(socket); + }); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('SMTP test server did not bind'); + process.env['SMTP_HOST'] = '127.0.0.1'; + process.env['SMTP_PORT'] = String(address.port); + delete process.env['SMTP_USER']; + delete process.env['SMTP_PASS']; + const provider = new SmtpEmailProvider(); + const didNotSettle = Symbol('did-not-settle'); + const send = (suffix: string) => provider.send({ + to: 'guest@example.com', + subject: `Hi ${suffix}`, + html: '

Hi

', + text: 'Hi', + messageId: ``, + }, { timeoutMs: 50 }); + + try { + for (const suffix of ['one', 'two']) { + const result = await Promise.race([ + send(suffix), + new Promise((resolve) => { + setTimeout(() => resolve(didNotSettle), 500); + }), + ]); + + expect(result).not.toBe(didNotSettle); + expect(result).toMatchObject({ + status: 'outcomeUnknown', + sent: false, + provider: 'smtp', + error: 'Email transport timed out', + }); + await vi.waitFor(() => expect(sockets.size).toBe(0), { timeout: 500 }); + } + expect(sessionsAtData).toBe(2); + expect(maxConcurrentConnections).toBe(1); + } finally { + for (const heartbeat of heartbeats.values()) clearInterval(heartbeat); + for (const socket of sockets) socket.destroy(); + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); + } + }); + + it('returns at the deadline without cancelling the scheduled late close', async () => { + const sockets = new Set(); + const server = createServer({ allowHalfOpen: true }, (socket) => { + sockets.add(socket); + let input = ''; + let readingData = false; + + socket.write('220 smtp.test ESMTP ready\r\n'); + socket.on('data', (chunk) => { + input += chunk.toString('utf8'); + if (readingData) return; + + let lineEnd: number; + while ((lineEnd = input.indexOf('\r\n')) >= 0) { + const command = input.slice(0, lineEnd); + input = input.slice(lineEnd + 2); + if (/^EHLO /i.test(command)) { + socket.write('250-smtp.test\r\n250 PIPELINING\r\n'); + } else if (/^MAIL FROM:/i.test(command) || /^RCPT TO:/i.test(command)) { + socket.write('250 2.1.0 Ok\r\n'); + } else if (/^DATA$/i.test(command)) { + readingData = true; + socket.write('354 End data with .\r\n'); + break; + } + } + }); + socket.on('error', () => undefined); + socket.on('close', () => sockets.delete(socket)); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('SMTP test server did not bind'); + process.env['SMTP_HOST'] = '127.0.0.1'; + process.env['SMTP_PORT'] = String(address.port); + delete process.env['SMTP_USER']; + delete process.env['SMTP_PASS']; + + const provider = new SmtpEmailProvider(); + const clearImmediateSpy = vi.spyOn(global, 'clearImmediate'); + + try { + const result = await provider.send({ + to: 'guest@example.com', + subject: 'Hi', + html: '

Hi

', + text: 'Hi', + }, { timeoutMs: 50 }); + expect(result.status).toBe('outcomeUnknown'); + expect(clearImmediateSpy).not.toHaveBeenCalled(); + } finally { + clearImmediateSpy.mockRestore(); + for (const socket of sockets) socket.destroy(); + await new Promise((resolve, reject) => { + server.close((error) => error ? reject(error) : resolve()); + }); + } + }); +}); diff --git a/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.ts b/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.ts index 4dab8d42..8f0b2c47 100644 --- a/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.ts +++ b/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.ts @@ -1,5 +1,40 @@ import { Injectable, Logger } from '@nestjs/common'; -import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider.interface'; +import type { + EmailMessage, + EmailProvider, + EmailResult, + EmailSendOptions, +} from '../email-provider.interface'; +import { + emailSendTimeoutMs, + notSentEmailResult, + sentEmailResult, + unknownTimeoutResult, +} from './bounded-email-transport'; + +interface OwnedSmtpPoolResource { + connection?: { + _socket?: OwnedSmtpConnectionSocket; + }; + close?: () => void; +} + +interface OwnedSmtpSocket { + destroyed?: boolean; + destroy?: () => void; +} + +interface OwnedSmtpConnectionSocket extends OwnedSmtpSocket { + socket?: OwnedSmtpSocket; +} + +interface OwnedSmtpTransport { + close?: () => void; + transporter?: { + _connections?: OwnedSmtpPoolResource[]; + }; + sendMail: (message: Record) => Promise<{ messageId: string }>; +} /** * SMTP transport (nodemailer) — configured via SMTP_HOST, SMTP_PORT, SMTP_USER, SMTP_PASS, SMTP_FROM. @@ -8,7 +43,8 @@ import type { EmailMessage, EmailProvider, EmailResult } from '../email-provider export class SmtpEmailProvider implements EmailProvider { readonly name = 'smtp'; private readonly logger = new Logger(SmtpEmailProvider.name); - private transport: any = null; + private nodemailer: any = null; + private transportConfig: Record | null = null; constructor() { this.initTransport(); @@ -28,12 +64,13 @@ export class SmtpEmailProvider implements EmailProvider { try { // eslint-disable-next-line @typescript-eslint/no-require-imports const nodemailer = require('nodemailer'); - this.transport = nodemailer.createTransport({ + this.nodemailer = nodemailer; + this.transportConfig = { host, port: parseInt(port, 10), secure: parseInt(port, 10) === 465, auth: user && pass ? { user, pass } : undefined, - }); + }; this.logger.log(`SMTP email provider configured: ${host}:${port}`); } catch { this.logger.warn('nodemailer not available — SMTP email provider disabled'); @@ -41,29 +78,94 @@ export class SmtpEmailProvider implements EmailProvider { } isConfigured(): boolean { - return this.transport !== null; + return this.nodemailer !== null && this.transportConfig !== null; } - async send(message: EmailMessage): Promise { - if (!this.transport) { - return { sent: false, provider: this.name, error: 'SMTP not configured' }; + async send(message: EmailMessage, options?: EmailSendOptions): Promise { + if (!this.isConfigured()) { + return notSentEmailResult(this.name, 'SMTP not configured'); } - try { - const from = message.from ?? process.env['SMTP_FROM'] ?? 'noreply@haip.dev'; - const info = await this.transport.sendMail({ - from, - to: message.to, - subject: message.subject, - html: message.html, - text: message.text, - }); - - this.logger.log(`Email sent via SMTP to ${message.to}: ${info.messageId}`); - return { sent: true, provider: this.name, messageId: info.messageId }; - } catch (error: any) { - this.logger.error(`SMTP send failed to ${message.to}: ${error.message}`); - return { sent: false, provider: this.name, error: error.message }; + const timeoutMs = emailSendTimeoutMs(options); + const transport = this.nodemailer.createTransport({ + ...this.transportConfig, + pool: true, + maxConnections: 1, + maxMessages: 1, + maxRequeues: 0, + connectionTimeout: timeoutMs, + greetingTimeout: timeoutMs, + socketTimeout: timeoutMs, + dnsTimeout: timeoutMs, + }) as OwnedSmtpTransport; + let timedOut = false; + const from = message.from ?? process.env['SMTP_FROM'] ?? 'noreply@haip.dev'; + const mailPayload = { + from, + to: message.to, + subject: message.subject, + html: message.html, + text: message.text, + messageId: message.messageId, + headers: message.idempotencyKey + ? { 'X-HAIP-Idempotency-Key': message.idempotencyKey } + : undefined, + }; + + const sendMailPromise = transport.sendMail(mailPayload).then( + (info) => { + if (timedOut) return unknownTimeoutResult(this.name); + this.logger.log(`Email sent via SMTP to ${message.to}: ${info.messageId}`); + return sentEmailResult(this.name, info.messageId); + }, + (error: any) => { + if ( + timedOut + || error?.code === 'ETIMEDOUT' + || /timed?\s*out|greeting never received/i.test(String(error?.message)) + ) { + return unknownTimeoutResult(this.name); + } + this.logger.error(`SMTP send failed to ${message.to}: ${error.message}`); + return notSentEmailResult(this.name, error.message); + }, + ); + + sendMailPromise.finally(() => { + this.closeOwnedTransport(transport); + }); + + const deadlinePromise = new Promise((resolve) => { + const timeout = setTimeout(() => { + timedOut = true; + this.closeOwnedTransport(transport); + const lateClose = setImmediate(() => this.closeOwnedTransport(transport)); + lateClose.unref?.(); + resolve(unknownTimeoutResult(this.name)); + }, timeoutMs); + timeout.unref?.(); + sendMailPromise.finally(() => clearTimeout(timeout)); + }); + + return await Promise.race([sendMailPromise, deadlinePromise]); + } + + /** + * Hard-close helper for per-send Nodemailer pools. Uses Nodemailer-internal + * pool/socket fields (`_connections`, `_socket`) — version-sensitive; covered + * by smtp-email.provider.spec integration tests. + */ + private closeOwnedTransport(transport: OwnedSmtpTransport): void { + const resources = [...(transport.transporter?._connections ?? [])]; + const sockets = resources.map((resource) => { + const wrappedSocket = resource.connection?._socket; + return wrappedSocket?.socket ?? wrappedSocket; + }); + + transport.close?.(); + for (const resource of resources) resource.close?.(); + for (const socket of sockets) { + if (!socket?.destroyed) socket?.destroy?.(); } } } diff --git a/docs/test-stats.json b/docs/test-stats.json index 570667f0..8f0b7de6 100644 --- a/docs/test-stats.json +++ b/docs/test-stats.json @@ -1,5 +1,5 @@ { - "tests": 1578, - "files": 219, - "updatedAt": "2026-08-27T01:42:47.833Z" + "tests": 1591, + "files": 220, + "updatedAt": "2026-08-27T08:48:14.423Z" } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index bb3b9352..55d5c7b8 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -92,6 +92,9 @@ importers: jwks-rsa: specifier: ^4.0.1 version: 4.0.1 + nodemailer: + specifier: ^9.0.5 + version: 9.0.5 passport: specifier: ^0.7.0 version: 0.7.0 @@ -4174,6 +4177,10 @@ packages: node-releases@2.0.37: resolution: {integrity: sha512-1h5gKZCF+pO/o3Iqt5Jp7wc9rH3eJJ0+nh/CIoiRwjRxde/hAHyLPXYN4V3CqKAbiZPSeJFSWHmJsbkicta0Eg==} + nodemailer@9.0.5: + resolution: {integrity: sha512-wvjiKvjczmsN7U/8006JOdXubgBk2XFAbioDMbT+sM7cPs0QrhJTa6KBRX7P5REGGkDcLUz/EarWidb8G8C1jQ==} + engines: {node: '>=6.0.0'} + normalize-path@3.0.0: resolution: {integrity: sha512-6eZs5Ls3WtCisHWp9S2GUy8dqkpGi4BVSz3GaqiE6ezub0512ESztXUwUB6C6IKbQkY2Pnb/mD4WYojCRwcwLA==} engines: {node: '>=0.10.0'} @@ -9143,6 +9150,8 @@ snapshots: node-releases@2.0.37: {} + nodemailer@9.0.5: {} + normalize-path@3.0.0: {} normalize-url@8.1.1: {}