From ba5b993c8eca66e9e92833bc83f642cd38cbe838 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 26 Aug 2026 18:32:43 +0000 Subject: [PATCH 1/5] feat(email): add bounded timeouts for outbound email providers Add bounded-email-transport with configurable connect/send timeouts. Extend EmailSendOptions across providers (SMTP, SES, Mailgun, SendGrid). Add SMTP provider unit tests. Co-authored-by: Agus --- apps/api/package.json | 2 + .../guest-comms/email-provider.interface.ts | 13 +- .../agent/guest-comms/email.service.spec.ts | 68 ++++++++ .../agent/guest-comms/email.service.ts | 13 +- .../providers/bounded-email-transport.ts | 63 +++++++ .../providers/console-email.provider.ts | 2 +- .../providers/mailgun-email.provider.ts | 50 ++++-- .../providers/mailgun-ses.provider.spec.ts | 117 +++++++++++++ .../providers/sendgrid-email.provider.ts | 53 ++++-- .../providers/ses-email.provider.ts | 56 ++++-- .../providers/smtp-email.provider.spec.ts | 162 ++++++++++++++++++ .../providers/smtp-email.provider.ts | 109 +++++++++++- pnpm-lock.yaml | 55 ++++++ 13 files changed, 710 insertions(+), 53 deletions(-) create mode 100644 apps/api/src/modules/agent/guest-comms/providers/bounded-email-transport.ts create mode 100644 apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.spec.ts diff --git a/apps/api/package.json b/apps/api/package.json index a1ed495b..e1fcf582 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -31,6 +31,7 @@ "@nestjs/websockets": "^10.0.0", "@telivityhaip/database": "workspace:*", "@telivityhaip/shared": "workspace:^", + "@telivityhaip/booking-requests": "workspace:*", "@types/jsonwebtoken": "^9.0.10", "bullmq": "^5.81.1", "class-transformer": "^0.5.1", @@ -42,6 +43,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..b70f4367 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,6 +4,10 @@ export interface EmailMessage { html: string; text: string; from?: string; + /** Stable caller identity for providers/gateways that support deduplication. */ + idempotencyKey?: string; + /** Stable RFC Message-ID reused when an at-least-once transport is retried. */ + messageId?: string; } export interface EmailResult { @@ -11,12 +15,19 @@ export interface EmailResult { messageId?: string; provider?: string; error?: string; + /** True when the transport may have accepted mail before timing out. */ + outcomeUnknown?: boolean; +} + +export interface EmailSendOptions { + /** Hard upper bound requested by the caller for transport settlement. */ + timeoutMs?: 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..80d6ba57 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 @@ -36,6 +36,28 @@ describe('EmailService', () => { 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({ sent: true, messageId: 'provider-id' }), + }; + 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() }; @@ -56,6 +78,7 @@ describe('SendgridEmailProvider', () => { const originalEnv = { ...process.env }; afterEach(() => { + vi.useRealTimers(); global.fetch = originalFetch; process.env = { ...originalEnv }; vi.resetModules(); @@ -92,11 +115,56 @@ describe('SendgridEmailProvider', () => { subject: 'Confirm', html: '

Hi

', text: 'Hi', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); 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('aborts and awaits settlement of a bounded SendGrid request', async () => { + vi.useFakeTimers(); + process.env['SENDGRID_API_KEY'] = 'SG.test'; + process.env['SENDGRID_FROM'] = 'hotel@example.com'; + let fetchSettled = false; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { + const signal = init?.signal; + signal?.addEventListener('abort', () => { + queueMicrotask(() => { + fetchSettled = true; + reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); + }); + }, { 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({ + sent: false, + outcomeUnknown: true, + error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); + expect(fetchSettled).toBe(true); }); }); 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..fd15f059 100644 --- a/apps/api/src/modules/agent/guest-comms/email.service.ts +++ b/apps/api/src/modules/agent/guest-comms/email.service.ts @@ -1,8 +1,13 @@ 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'; /** * Email transport service — SendGrid, Mailgun, SES gateway, SMTP, or console fallback. @@ -19,9 +24,9 @@ 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); + const result = await provider.send(message, options); if (!result.provider) { return { ...result, provider: provider.name }; } 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..c42e3aac --- /dev/null +++ b/apps/api/src/modules/agent/guest-comms/providers/bounded-email-transport.ts @@ -0,0 +1,63 @@ +import type { 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)); +} + +/** + * Aborts an HTTP transport at the deadline but does not return until fetch has + * actually settled, so callers never make the delivery retry-eligible while + * the original in-process request is still live. + */ +export async function boundedEmailFetch( + input: string, + init: RequestInit, + options: EmailSendOptions | undefined, + consume: (response: Response) => Promise | T, +): Promise { + const controller = new AbortController(); + let timedOut = false; + const timeout = setTimeout(() => { + timedOut = true; + controller.abort(); + }, emailSendTimeoutMs(options)); + timeout.unref?.(); + try { + const response = await fetch(input, { ...init, signal: controller.signal }); + const result = await consume(response); + if (timedOut) throw new EmailTransportTimeoutError(); + return result; + } catch (error: unknown) { + if (timedOut) throw new EmailTransportTimeoutError(); + throw error; + } finally { + clearTimeout(timeout); + } +} + +export function unknownTimeoutResult(provider: string): { + sent: false; + provider: string; + error: string; + outcomeUnknown: true; +} { + return { + sent: false, + provider, + error: 'Email transport timed out', + outcomeUnknown: true, + }; +} 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..dd4d349a 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 @@ -20,7 +20,7 @@ export class ConsoleEmailProvider implements EmailProvider { return { sent: false, provider: this.name, - messageId: `console-${Date.now()}`, + messageId: message.messageId ?? `console-${Date.now()}`, error: 'No email provider configured — message logged only', }; } 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..c3d6c64c 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,15 @@ 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, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * Mailgun Messages API adapter. @@ -24,7 +34,7 @@ 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' }; } @@ -35,21 +45,32 @@ 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, @@ -60,6 +81,9 @@ export class MailgunEmailProvider implements EmailProvider { this.logger.log(`Email sent via Mailgun to ${message.to}`); return { sent: true, provider: this.name, messageId: 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 }; } 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..29ff0d8d 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,79 @@ describe('MailgunEmailProvider', () => { subject: 'S', html: 'h', text: 't', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); 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('aborts and awaits settlement of a bounded Mailgun request', async () => { + vi.useFakeTimers(); + process.env['MAILGUN_API_KEY'] = 'key'; + process.env['MAILGUN_DOMAIN'] = 'mg.example.com'; + let settled = false; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => { + queueMicrotask(() => { + settled = true; + reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); + }); + }, { 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({ + sent: false, outcomeUnknown: true, error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); + expect(settled).toBe(true); + }); + + it('keeps the bound active until the Mailgun response body settles', async () => { + vi.useFakeTimers(); + process.env['MAILGUN_API_KEY'] = 'key'; + process.env['MAILGUN_DOMAIN'] = 'mg.example.com'; + let bodySettled = false; + global.fetch = vi.fn((_url, init) => Promise.resolve({ + ok: true, + json: () => new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => { + queueMicrotask(() => { + bodySettled = true; + reject(Object.assign(new Error('aborted body'), { name: 'AbortError' })); + }); + }, { 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; + await vi.advanceTimersByTimeAsync(100); + + expect(signal?.aborted).toBe(true); + await expect(sending).resolves.toMatchObject({ + sent: false, outcomeUnknown: true, error: 'Email transport timed out', + }); + expect(bodySettled).toBe(true); }); }); @@ -44,6 +115,7 @@ describe('SesEmailProvider', () => { const originalEnv = { ...process.env }; afterEach(() => { + vi.useRealTimers(); global.fetch = originalFetch; process.env = { ...originalEnv }; vi.resetModules(); @@ -73,11 +145,56 @@ describe('SesEmailProvider', () => { subject: 'S', html: 'h', text: 't', + idempotencyKey: 'stable-delivery-1', + messageId: '', }); expect(result).toEqual({ 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: '', + }); + expect(payload.Content.Simple.Headers).not.toContainEqual(expect.objectContaining({ + Name: 'Message-ID', + })); + }); + + it('aborts and awaits settlement of a bounded SES gateway request', async () => { + vi.useFakeTimers(); + process.env['SES_ENDPOINT'] = 'http://localhost:4566'; + process.env['SES_API_KEY'] = 'local'; + process.env['SES_FROM'] = 'noreply@example.com'; + let settled = false; + global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => { + queueMicrotask(() => { + settled = true; + reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); + }); + }, { 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({ + sent: false, outcomeUnknown: true, error: 'Email transport timed out', + }); + expect(signal?.aborted).toBe(true); + expect(settled).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..133f0970 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,15 @@ 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, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * SendGrid Email API reference adapter. @@ -32,14 +42,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' }; } 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,18 +68,25 @@ 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}`); + this.logger.error(`SendGrid send failed (${res.status}): ${failureBody ?? ''}`); return { sent: false, provider: this.name, error: `SendGrid HTTP ${res.status}` }; } @@ -68,6 +94,9 @@ export class SendgridEmailProvider implements EmailProvider { this.logger.log(`Email sent via SendGrid to ${message.to}`); return { sent: true, provider: 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 }; } 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..970cec0c 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,15 @@ 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, + unknownTimeoutResult, +} from './bounded-email-transport'; /** * Amazon SES outbound adapter via an explicit HTTPS gateway. @@ -22,7 +32,7 @@ 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, @@ -37,6 +47,11 @@ export class SesEmailProvider implements EmailProvider { Content: { Simple: { Subject: { Data: message.subject }, + // SES assigns and overwrites the RFC Message-ID. Preserve our stable + // logical identity in a permitted custom header for gateway replay. + ...(message.messageId + ? { Headers: [{ Name: 'X-HAIP-Message-ID', Value: message.messageId }] } + : {}), Body: { Text: { Data: message.text }, Html: { Data: message.html }, @@ -46,19 +61,29 @@ 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, @@ -69,6 +94,9 @@ export class SesEmailProvider implements EmailProvider { this.logger.log(`Email sent via SES gateway to ${message.to}`); return { sent: true, provider: this.name, messageId: 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 }; } 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..5a042fc9 --- /dev/null +++ b/apps/api/src/modules/agent/guest-comms/providers/smtp-email.provider.spec.ts @@ -0,0 +1,162 @@ +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({ + sent: false, + provider: 'smtp', + outcomeUnknown: true, + 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; + // Keep TCP traffic flowing without completing the SMTP DATA reply. + // This defeats Nodemailer's socket-inactivity timeout so only the + // provider's owned hard deadline can end the live transaction. + 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({ + sent: false, + provider: 'smtp', + outcomeUnknown: true, + 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()); + }); + } + }); +}); 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..6b49ddf4 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,38 @@ 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, + 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 +41,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 +62,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 +76,87 @@ 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) { + async send(message: EmailMessage, options?: EmailSendOptions): Promise { + if (!this.isConfigured()) { return { sent: false, provider: this.name, error: 'SMTP not configured' }; } + 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; + let lateClose: ReturnType | undefined; + const timeout = setTimeout(() => { + timedOut = true; + this.closeOwnedTransport(transport); + // Pool resource setup itself is asynchronous. Re-close on the next turn + // so a resource created at the deadline cannot outlive this send. + lateClose = setImmediate(() => this.closeOwnedTransport(transport)); + lateClose.unref?.(); + }, timeoutMs); + timeout.unref?.(); try { const from = message.from ?? process.env['SMTP_FROM'] ?? 'noreply@haip.dev'; - const info = await this.transport.sendMail({ + const info = await transport.sendMail({ 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, }); + if (timedOut) return unknownTimeoutResult(this.name); this.logger.log(`Email sent via SMTP to ${message.to}: ${info.messageId}`); return { sent: true, provider: this.name, messageId: info.messageId }; } catch (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 { sent: false, provider: this.name, error: error.message }; + } finally { + clearTimeout(timeout); + if (lateClose) clearImmediate(lateClose); + // This runs only after sendMail has settled. Destroying again here makes + // return from send() the ownership boundary for every per-send socket. + this.closeOwnedTransport(transport); + } + } + + private closeOwnedTransport(transport: OwnedSmtpTransport): void { + const resources = [...(transport.transporter?._connections ?? [])]; + const sockets = resources.map((resource) => { + const wrappedSocket = resource.connection?._socket; + return wrappedSocket?.socket ?? wrappedSocket; + }); + + // Marks the pool closed and fails any work that has not acquired a resource. + transport.close?.(); + for (const resource of resources) resource.close?.(); + // SMTPConnection.close() is graceful after greeting. A hard deadline also + // destroys the owned socket so an active half-open transaction cannot live. + for (const socket of sockets) { + if (!socket?.destroyed) socket?.destroy?.(); } } } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index bb3b9352..6ca2016f 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -53,6 +53,9 @@ importers: '@nestjs/websockets': specifier: ^10.0.0 version: 10.4.22(@nestjs/common@10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@10.4.22)(@nestjs/platform-socket.io@10.4.22)(reflect-metadata@0.2.2)(rxjs@7.8.2) + '@telivityhaip/booking-requests': + specifier: workspace:* + version: link:../../packages/booking-requests '@telivityhaip/database': specifier: workspace:* version: link:../../packages/database @@ -92,6 +95,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 @@ -324,6 +330,49 @@ importers: specifier: ^2.1.0 version: 2.1.9(@types/node@25.5.2)(jsdom@25.0.1)(terser@5.46.1) + packages/booking-requests: + dependencies: + '@telivityhaip/database': + specifier: workspace:* + version: link:../database + '@telivityhaip/shared': + specifier: workspace:^ + version: link:../shared + decimal.js: + specifier: ^10.4.3 + version: 10.6.0 + drizzle-orm: + specifier: ^0.39.0 + version: 0.39.3(postgres@3.4.9) + devDependencies: + '@nestjs/common': + specifier: ^10.4.0 + version: 10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2) + '@nestjs/core': + specifier: ^10.4.0 + version: 10.4.22(@nestjs/common@10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/platform-express@10.4.22)(@nestjs/websockets@10.4.22)(reflect-metadata@0.2.2)(rxjs@7.8.2) + '@types/node': + specifier: ^22.0.0 + version: 22.19.17 + eslint: + specifier: ^9.0.0 + version: 9.39.4(jiti@1.21.7) + postgres: + specifier: ^3.4.5 + version: 3.4.9 + tsup: + specifier: ^8.3.0 + version: 8.5.1(@swc/core@1.15.24)(jiti@1.21.7)(postcss@8.5.8)(tsx@4.21.0)(typescript@5.9.3) + tsx: + specifier: ^4.19.0 + version: 4.21.0 + typescript: + specifier: ^5.7.0 + version: 5.9.3 + vitest: + specifier: ^2.1.0 + version: 2.1.9(@types/node@22.19.17)(jsdom@25.0.1)(terser@5.46.1) + packages/database: dependencies: drizzle-orm: @@ -4174,6 +4223,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 +9196,8 @@ snapshots: node-releases@2.0.37: {} + nodemailer@9.0.5: {} + normalize-path@3.0.0: {} normalize-url@8.1.1: {} From ddfc90e5c56123c689076f5b3ca392d5e43281e8 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Wed, 26 Aug 2026 19:57:02 +0000 Subject: [PATCH 2/5] chore: sync README test counts for CI Co-authored-by: telivity-otaip --- README.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 0459cd55..2a3549b2 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ NestJS PostgreSQL Apache 2.0 License -1568 Tests Passing 12 AI Agents +1575 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 (1568 tests across 218 test files) | Unit and integration tests || Build | tsup (packages) + Vite (dashboard) + nest build (API) | Fast builds | +| Testing | Vitest (1575 tests across 219 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 (1568 tests across 218 test files) +# All tests (1575 tests across 219 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 (1568 tests, 218 files) +pnpm test # Run all tests (1575 tests, 219 files) pnpm lint # ESLint ``` From bcd63cac5657621f65addc95c27ac8cf053beee7 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 27 Aug 2026 00:49:12 +0000 Subject: [PATCH 3/5] =?UTF-8?q?fix(email):=20address=20#349=20review=20?= =?UTF-8?q?=E2=80=94=20drop=20unrelated=20dep,=20clarify=20timeout=20contr?= =?UTF-8?q?act?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Remove @telivityhaip/booking-requests from the email-only PR. Document cooperative HTTP deadlines and idempotency metadata. SMTP returns at the deadline via Promise.race while hard-closing owned pool sockets. Co-authored-by: telivity-otaip --- apps/api/package.json | 1 - .../guest-comms/email-provider.interface.ts | 18 +++- .../providers/bounded-email-transport.ts | 8 +- .../providers/smtp-email.provider.ts | 92 +++++++++++-------- pnpm-lock.yaml | 46 ---------- 5 files changed, 74 insertions(+), 91 deletions(-) diff --git a/apps/api/package.json b/apps/api/package.json index e1fcf582..67ea9504 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -31,7 +31,6 @@ "@nestjs/websockets": "^10.0.0", "@telivityhaip/database": "workspace:*", "@telivityhaip/shared": "workspace:^", - "@telivityhaip/booking-requests": "workspace:*", "@types/jsonwebtoken": "^9.0.10", "bullmq": "^5.81.1", "class-transformer": "^0.5.1", 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 b70f4367..831d8190 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,9 +4,16 @@ export interface EmailMessage { html: string; text: string; from?: string; - /** Stable caller identity for providers/gateways that support deduplication. */ + /** + * Stable caller identity forwarded to providers that support deduplication + * (e.g. Mailgun `v:haip-idempotency-key`). Safe to retry when the transport + * returns `outcomeUnknown` after a cooperative deadline. + */ idempotencyKey?: string; - /** Stable RFC Message-ID reused when an at-least-once transport is retried. */ + /** + * Stable RFC Message-ID reused across retries so duplicate deliveries share + * the same provider-visible message identity when supported. + */ messageId?: string; } @@ -20,7 +27,12 @@ export interface EmailResult { } export interface EmailSendOptions { - /** Hard upper bound requested by the caller for transport settlement. */ + /** + * Cooperative send deadline in milliseconds. HTTP transports abort at the + * deadline but await settlement before returning `outcomeUnknown`; SMTP + * closes owned sockets and returns at the deadline even if `sendMail` is + * still settling. + */ timeoutMs?: number; } 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 index c42e3aac..9b50e4dd 100644 --- 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 @@ -18,9 +18,11 @@ export function emailSendTimeoutMs(options?: EmailSendOptions): number { } /** - * Aborts an HTTP transport at the deadline but does not return until fetch has - * actually settled, so callers never make the delivery retry-eligible while - * the original in-process request is still live. + * Cooperative HTTP send deadline: aborts fetch at `timeoutMs` but does not + * return until the in-flight request settles (success, error, or abort). + * This prevents marking a delivery retry-eligible while the original socket + * may still complete server-side. Callers receive `outcomeUnknown` only after + * settlement when the deadline fired first. */ export async function boundedEmailFetch( input: string, 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 6b49ddf4..4b1c85e0 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 @@ -98,51 +98,67 @@ export class SmtpEmailProvider implements EmailProvider { }) as OwnedSmtpTransport; let timedOut = false; let lateClose: ReturnType | undefined; - const timeout = setTimeout(() => { - timedOut = true; - this.closeOwnedTransport(transport); - // Pool resource setup itself is asynchronous. Re-close on the next turn - // so a resource created at the deadline cannot outlive this send. - lateClose = setImmediate(() => this.closeOwnedTransport(transport)); - lateClose.unref?.(); - }, timeoutMs); - timeout.unref?.(); + 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 { sent: true, provider: this.name, messageId: info.messageId } satisfies EmailResult; + }, + (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 { sent: false, provider: this.name, error: error.message } satisfies EmailResult; + }, + ); + + const deadlinePromise = new Promise((resolve) => { + const timeout = setTimeout(() => { + timedOut = true; + this.closeOwnedTransport(transport); + // Pool resource setup itself is asynchronous. Re-close on the next turn + // so a resource created at the deadline cannot outlive this send. + lateClose = setImmediate(() => this.closeOwnedTransport(transport)); + lateClose.unref?.(); + resolve(unknownTimeoutResult(this.name)); + }, timeoutMs); + timeout.unref?.(); + sendMailPromise.finally(() => clearTimeout(timeout)); + }); + try { - const from = message.from ?? process.env['SMTP_FROM'] ?? 'noreply@haip.dev'; - const info = await transport.sendMail({ - 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, - }); - - if (timedOut) return unknownTimeoutResult(this.name); - this.logger.log(`Email sent via SMTP to ${message.to}: ${info.messageId}`); - return { sent: true, provider: this.name, messageId: info.messageId }; - } catch (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 { sent: false, provider: this.name, error: error.message }; + return await Promise.race([sendMailPromise, deadlinePromise]); } finally { - clearTimeout(timeout); if (lateClose) clearImmediate(lateClose); - // This runs only after sendMail has settled. Destroying again here makes - // return from send() the ownership boundary for every per-send socket. + // Destroy again after settlement so return from send() is the ownership + // boundary for every per-send socket. this.closeOwnedTransport(transport); } } + /** + * 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) => { diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 6ca2016f..55d5c7b8 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -53,9 +53,6 @@ importers: '@nestjs/websockets': specifier: ^10.0.0 version: 10.4.22(@nestjs/common@10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@10.4.22)(@nestjs/platform-socket.io@10.4.22)(reflect-metadata@0.2.2)(rxjs@7.8.2) - '@telivityhaip/booking-requests': - specifier: workspace:* - version: link:../../packages/booking-requests '@telivityhaip/database': specifier: workspace:* version: link:../../packages/database @@ -330,49 +327,6 @@ importers: specifier: ^2.1.0 version: 2.1.9(@types/node@25.5.2)(jsdom@25.0.1)(terser@5.46.1) - packages/booking-requests: - dependencies: - '@telivityhaip/database': - specifier: workspace:* - version: link:../database - '@telivityhaip/shared': - specifier: workspace:^ - version: link:../shared - decimal.js: - specifier: ^10.4.3 - version: 10.6.0 - drizzle-orm: - specifier: ^0.39.0 - version: 0.39.3(postgres@3.4.9) - devDependencies: - '@nestjs/common': - specifier: ^10.4.0 - version: 10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2) - '@nestjs/core': - specifier: ^10.4.0 - version: 10.4.22(@nestjs/common@10.4.22(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/platform-express@10.4.22)(@nestjs/websockets@10.4.22)(reflect-metadata@0.2.2)(rxjs@7.8.2) - '@types/node': - specifier: ^22.0.0 - version: 22.19.17 - eslint: - specifier: ^9.0.0 - version: 9.39.4(jiti@1.21.7) - postgres: - specifier: ^3.4.5 - version: 3.4.9 - tsup: - specifier: ^8.3.0 - version: 8.5.1(@swc/core@1.15.24)(jiti@1.21.7)(postcss@8.5.8)(tsx@4.21.0)(typescript@5.9.3) - tsx: - specifier: ^4.19.0 - version: 4.21.0 - typescript: - specifier: ^5.7.0 - version: 5.9.3 - vitest: - specifier: ^2.1.0 - version: 2.1.9(@types/node@22.19.17)(jsdom@25.0.1)(terser@5.46.1) - packages/database: dependencies: drizzle-orm: From 72e3aff15fe82937aa08b08a1e4160cc74108c29 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 27 Aug 2026 08:25:57 +0000 Subject: [PATCH 4/5] fix(email): address Agus review on delivery contract and hard deadlines - Add EmailDeliveryStatus (sent/notSent/outcomeUnknown) to EmailResult - EmailService auto-retries notSent only; never retries outcomeUnknown - HTTP boundedEmailFetch: hard outer deadline incl. body consumption - SMTP: cleanup on sendMail settlement; do not cancel late close in finally - Clarify idempotencyKey/messageId as correlation, not exactly-once - Update guest-comms tests for retry, HTTP deadline, and SMTP late-close Co-authored-by: telivity-otaip --- .../guest-comms/email-provider.interface.ts | 31 ++- .../agent/guest-comms/email.service.spec.ts | 192 ++++++++++++++++-- .../agent/guest-comms/email.service.ts | 38 +++- .../guest-communication.agent.spec.ts | 2 +- .../providers/bounded-email-transport.ts | 82 ++++---- .../providers/console-email.provider.ts | 5 +- .../providers/mailgun-email.provider.ts | 14 +- .../providers/mailgun-ses.provider.spec.ts | 63 +++--- .../providers/sendgrid-email.provider.ts | 10 +- .../providers/ses-email.provider.ts | 23 +-- .../providers/smtp-email.provider.spec.ts | 66 +++++- .../providers/smtp-email.provider.ts | 29 +-- 12 files changed, 398 insertions(+), 157 deletions(-) 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 831d8190..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 @@ -5,35 +5,44 @@ export interface EmailMessage { text: string; from?: string; /** - * Stable caller identity forwarded to providers that support deduplication - * (e.g. Mailgun `v:haip-idempotency-key`). Safe to retry when the transport - * returns `outcomeUnknown` after a cooperative deadline. + * 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 so duplicate deliveries share - * the same provider-visible message identity when supported. + * 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; - /** True when the transport may have accepted mail before timing out. */ - outcomeUnknown?: boolean; } export interface EmailSendOptions { /** - * Cooperative send deadline in milliseconds. HTTP transports abort at the - * deadline but await settlement before returning `outcomeUnknown`; SMTP - * closes owned sockets and returns at the deadline even if `sendMail` is - * still settling. + * 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 { 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 80d6ba57..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,7 +42,8 @@ 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(); }); @@ -40,7 +51,11 @@ describe('EmailService', () => { const provider = { name: 'sendgrid', isConfigured: () => true, - send: vi.fn().mockResolvedValue({ sent: true, messageId: 'provider-id' }), + send: vi.fn().mockResolvedValue({ + status: 'sent', + sent: true, + messageId: 'provider-id', + } satisfies EmailResult), }; const service = new EmailService([provider]); const message = { @@ -71,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', () => { @@ -96,6 +192,7 @@ describe('SendgridEmailProvider', () => { html: 'h', text: 't', }); + expect(result.status).toBe('notSent'); expect(result.sent).toBe(false); }); @@ -118,6 +215,7 @@ describe('SendgridEmailProvider', () => { 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', @@ -131,19 +229,12 @@ describe('SendgridEmailProvider', () => { }); }); - it('aborts and awaits settlement of a bounded SendGrid request', async () => { + 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'; - let fetchSettled = false; - global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { - const signal = init?.signal; - signal?.addEventListener('abort', () => { - queueMicrotask(() => { - fetchSettled = true; - reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); - }); - }, { once: true }); + 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'); @@ -160,11 +251,78 @@ describe('SendgridEmailProvider', () => { await vi.advanceTimersByTimeAsync(100); await expect(sending).resolves.toMatchObject({ + status: 'outcomeUnknown', sent: false, - outcomeUnknown: true, error: 'Email transport timed out', }); expect(signal?.aborted).toBe(true); - expect(fetchSettled).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 fd15f059..94320659 100644 --- a/apps/api/src/modules/agent/guest-comms/email.service.ts +++ b/apps/api/src/modules/agent/guest-comms/email.service.ts @@ -9,10 +9,15 @@ import { EMAIL_PROVIDERS } 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 { @@ -26,16 +31,41 @@ export class EmailService { async send(message: EmailMessage, options?: EmailSendOptions): Promise { const provider = this.activeProvider(); - const result = await provider.send(message, options); - 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 index 9b50e4dd..982cbfde 100644 --- 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 @@ -1,4 +1,4 @@ -import type { EmailSendOptions } from '../email-provider.interface'; +import type { EmailResult, EmailSendOptions } from '../email-provider.interface'; export const DEFAULT_EMAIL_SEND_TIMEOUT_MS = 60_000; @@ -17,12 +17,34 @@ export function emailSendTimeoutMs(options?: EmailSendOptions): number { 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', + }; +} + /** - * Cooperative HTTP send deadline: aborts fetch at `timeoutMs` but does not - * return until the in-flight request settles (success, error, or abort). - * This prevents marking a delivery retry-eligible while the original socket - * may still complete server-side. Callers receive `outcomeUnknown` only after - * settlement when the deadline fired first. + * 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, @@ -31,35 +53,27 @@ export async function boundedEmailFetch( consume: (response: Response) => Promise | T, ): Promise { const controller = new AbortController(); - let timedOut = false; - const timeout = setTimeout(() => { - timedOut = true; - controller.abort(); - }, emailSendTimeoutMs(options)); - timeout.unref?.(); - try { + const timeoutMs = emailSendTimeoutMs(options); + + const work = (async (): Promise => { const response = await fetch(input, { ...init, signal: controller.signal }); - const result = await consume(response); - if (timedOut) throw new EmailTransportTimeoutError(); - return result; - } catch (error: unknown) { - if (timedOut) throw new EmailTransportTimeoutError(); - throw error; + 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 { - clearTimeout(timeout); + if (timer) clearTimeout(timer); } } - -export function unknownTimeoutResult(provider: string): { - sent: false; - provider: string; - error: string; - outcomeUnknown: true; -} { - return { - sent: false, - provider, - error: 'Email transport timed out', - outcomeUnknown: true, - }; -} 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 dd4d349a..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, + ...notSentEmailResult(this.name, 'No email provider configured — message logged only'), messageId: message.messageId ?? `console-${Date.now()}`, - error: 'No email provider configured — message logged only', }; } } 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 c3d6c64c..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 @@ -8,6 +8,8 @@ import type { import { boundedEmailFetch, EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, unknownTimeoutResult, } from './bounded-email-transport'; @@ -36,7 +38,7 @@ export class MailgunEmailProvider implements EmailProvider { 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(); @@ -72,20 +74,16 @@ export class MailgunEmailProvider implements EmailProvider { }), ); 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 29ff0d8d..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 @@ -37,6 +37,7 @@ describe('MailgunEmailProvider', () => { 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]; @@ -45,18 +46,12 @@ describe('MailgunEmailProvider', () => { expect(form.get('v:haip-idempotency-key')).toBe('stable-delivery-1'); }); - it('aborts and awaits settlement of a bounded Mailgun request', async () => { + 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'; - let settled = false; - global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { - init?.signal?.addEventListener('abort', () => { - queueMicrotask(() => { - settled = true; - reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); - }); - }, { once: true }); + 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'); @@ -70,43 +65,42 @@ describe('MailgunEmailProvider', () => { await vi.advanceTimersByTimeAsync(100); await expect(sending).resolves.toMatchObject({ - sent: false, outcomeUnknown: true, error: 'Email transport timed out', + status: 'outcomeUnknown', + sent: false, + error: 'Email transport timed out', }); expect(signal?.aborted).toBe(true); - expect(settled).toBe(true); }); - it('keeps the bound active until the Mailgun response body settles', async () => { + 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'; - let bodySettled = false; global.fetch = vi.fn((_url, init) => Promise.resolve({ ok: true, json: () => new Promise((_resolve, reject) => { init?.signal?.addEventListener('abort', () => { - queueMicrotask(() => { - bodySettled = true; - reject(Object.assign(new Error('aborted body'), { name: 'AbortError' })); - }); + 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()); - const signal = vi.mocked(global.fetch).mock.calls[0]?.[1]?.signal; await vi.advanceTimersByTimeAsync(100); - - expect(signal?.aborted).toBe(true); - await expect(sending).resolves.toMatchObject({ - sent: false, outcomeUnknown: true, error: 'Email transport timed out', - }); - expect(bodySettled).toBe(true); + await expect(sending).resolves.toMatchObject({ status: 'outcomeUnknown' }); + await vi.advanceTimersByTimeAsync(100); + await Promise.resolve(); + expect(unhandled).toEqual([]); + process.off('unhandledRejection', onUnhandled); }); }); @@ -149,6 +143,7 @@ describe('SesEmailProvider', () => { messageId: '', }); expect(result).toEqual({ + status: 'sent', sent: true, provider: 'amazon-ses', messageId: 'ses-1', @@ -161,24 +156,15 @@ describe('SesEmailProvider', () => { expect(payload.Content.Simple.Headers).toContainEqual({ Name: 'X-HAIP-Message-ID', Value: '', }); - expect(payload.Content.Simple.Headers).not.toContainEqual(expect.objectContaining({ - Name: 'Message-ID', - })); }); - it('aborts and awaits settlement of a bounded SES gateway request', async () => { + 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'; - let settled = false; - global.fetch = vi.fn((_url, init) => new Promise((_resolve, reject) => { - init?.signal?.addEventListener('abort', () => { - queueMicrotask(() => { - settled = true; - reject(Object.assign(new Error('aborted'), { name: 'AbortError' })); - }); - }, { once: true }); + 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'); @@ -192,9 +178,10 @@ describe('SesEmailProvider', () => { await vi.advanceTimersByTimeAsync(100); await expect(sending).resolves.toMatchObject({ - sent: false, outcomeUnknown: true, error: 'Email transport timed out', + status: 'outcomeUnknown', + sent: false, + error: 'Email transport timed out', }); expect(signal?.aborted).toBe(true); - expect(settled).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 133f0970..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 @@ -8,6 +8,8 @@ import type { import { boundedEmailFetch, EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, unknownTimeoutResult, } from './bounded-email-transport'; @@ -44,7 +46,7 @@ export class SendgridEmailProvider implements EmailProvider { 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!; @@ -87,18 +89,18 @@ export class SendgridEmailProvider implements EmailProvider { if (!res.ok) { this.logger.error(`SendGrid send failed (${res.status}): ${failureBody ?? ''}`); - return { sent: false, provider: this.name, error: `SendGrid HTTP ${res.status}` }; + 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 970cec0c..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 @@ -8,6 +8,8 @@ import type { import { boundedEmailFetch, EmailTransportTimeoutError, + notSentEmailResult, + sentEmailResult, unknownTimeoutResult, } from './bounded-email-transport'; @@ -34,11 +36,10 @@ export class SesEmailProvider implements EmailProvider { 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 = { @@ -47,8 +48,6 @@ export class SesEmailProvider implements EmailProvider { Content: { Simple: { Subject: { Data: message.subject }, - // SES assigns and overwrites the RFC Message-ID. Preserve our stable - // logical identity in a permitted custom header for gateway replay. ...(message.messageId ? { Headers: [{ Name: 'X-HAIP-Message-ID', Value: message.messageId }] } : {}), @@ -85,20 +84,16 @@ export class SesEmailProvider implements EmailProvider { }), ); 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 index 5a042fc9..c539b0f2 100644 --- 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 @@ -41,9 +41,9 @@ describe('SmtpEmailProvider bounded send', () => { expect(result).not.toBe(didNotSettle); expect(result).toMatchObject({ + status: 'outcomeUnknown', sent: false, provider: 'smtp', - outcomeUnknown: true, error: 'Email transport timed out', }); await vi.waitFor(() => expect(sockets.size).toBe(0), { timeout: 500 }); @@ -75,9 +75,6 @@ describe('SmtpEmailProvider bounded send', () => { input = input.slice(dataEnd + 5); readingData = false; sessionsAtData += 1; - // Keep TCP traffic flowing without completing the SMTP DATA reply. - // This defeats Nodemailer's socket-inactivity timeout so only the - // provider's owned hard deadline can end the live transaction. const heartbeat = setInterval(() => { if (!socket.destroyed && socket.writable) socket.write(' '); }, 5); @@ -142,9 +139,9 @@ describe('SmtpEmailProvider bounded send', () => { expect(result).not.toBe(didNotSettle); expect(result).toMatchObject({ + status: 'outcomeUnknown', sent: false, provider: 'smtp', - outcomeUnknown: true, error: 'Email transport timed out', }); await vi.waitFor(() => expect(sockets.size).toBe(0), { timeout: 500 }); @@ -159,4 +156,63 @@ describe('SmtpEmailProvider bounded send', () => { }); } }); + + 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 4b1c85e0..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 @@ -7,6 +7,8 @@ import type { } from '../email-provider.interface'; import { emailSendTimeoutMs, + notSentEmailResult, + sentEmailResult, unknownTimeoutResult, } from './bounded-email-transport'; @@ -81,7 +83,7 @@ export class SmtpEmailProvider implements EmailProvider { async send(message: EmailMessage, options?: EmailSendOptions): Promise { if (!this.isConfigured()) { - return { sent: false, provider: this.name, error: 'SMTP not configured' }; + return notSentEmailResult(this.name, 'SMTP not configured'); } const timeoutMs = emailSendTimeoutMs(options); @@ -97,7 +99,6 @@ export class SmtpEmailProvider implements EmailProvider { dnsTimeout: timeoutMs, }) as OwnedSmtpTransport; let timedOut = false; - let lateClose: ReturnType | undefined; const from = message.from ?? process.env['SMTP_FROM'] ?? 'noreply@haip.dev'; const mailPayload = { from, @@ -115,7 +116,7 @@ export class SmtpEmailProvider implements EmailProvider { (info) => { if (timedOut) return unknownTimeoutResult(this.name); this.logger.log(`Email sent via SMTP to ${message.to}: ${info.messageId}`); - return { sent: true, provider: this.name, messageId: info.messageId } satisfies EmailResult; + return sentEmailResult(this.name, info.messageId); }, (error: any) => { if ( @@ -126,17 +127,19 @@ export class SmtpEmailProvider implements EmailProvider { return unknownTimeoutResult(this.name); } this.logger.error(`SMTP send failed to ${message.to}: ${error.message}`); - return { sent: false, provider: this.name, error: error.message } satisfies EmailResult; + return notSentEmailResult(this.name, error.message); }, ); + sendMailPromise.finally(() => { + this.closeOwnedTransport(transport); + }); + const deadlinePromise = new Promise((resolve) => { const timeout = setTimeout(() => { timedOut = true; this.closeOwnedTransport(transport); - // Pool resource setup itself is asynchronous. Re-close on the next turn - // so a resource created at the deadline cannot outlive this send. - lateClose = setImmediate(() => this.closeOwnedTransport(transport)); + const lateClose = setImmediate(() => this.closeOwnedTransport(transport)); lateClose.unref?.(); resolve(unknownTimeoutResult(this.name)); }, timeoutMs); @@ -144,14 +147,7 @@ export class SmtpEmailProvider implements EmailProvider { sendMailPromise.finally(() => clearTimeout(timeout)); }); - try { - return await Promise.race([sendMailPromise, deadlinePromise]); - } finally { - if (lateClose) clearImmediate(lateClose); - // Destroy again after settlement so return from send() is the ownership - // boundary for every per-send socket. - this.closeOwnedTransport(transport); - } + return await Promise.race([sendMailPromise, deadlinePromise]); } /** @@ -166,11 +162,8 @@ export class SmtpEmailProvider implements EmailProvider { return wrappedSocket?.socket ?? wrappedSocket; }); - // Marks the pool closed and fails any work that has not acquired a resource. transport.close?.(); for (const resource of resources) resource.close?.(); - // SMTPConnection.close() is graceful after greeting. A hard deadline also - // destroys the owned socket so an active half-open transaction cannot live. for (const socket of sockets) { if (!socket?.destroyed) socket?.destroy?.(); } From cf6d97827880095121d1f36e01b545cea0224938 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 27 Aug 2026 08:49:09 +0000 Subject: [PATCH 5/5] chore: sync README test counts after email transport tests (1591 tests) Co-authored-by: telivity-otaip --- README.md | 8 ++++---- docs/test-stats.json | 4 ++-- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index 2736b8e3..7afbd998 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ NestJS PostgreSQL Apache 2.0 License - 1585 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 (1585 tests across 220 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 (1585 tests across 220 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 (1585 tests, 220 files) +pnpm test # Run all tests (1591 tests, 220 files) pnpm lint # ESLint ``` diff --git a/docs/test-stats.json b/docs/test-stats.json index 85f681f9..8f0b7de6 100644 --- a/docs/test-stats.json +++ b/docs/test-stats.json @@ -1,5 +1,5 @@ { - "tests": 1585, + "tests": 1591, "files": 220, - "updatedAt": "2026-08-27T02:03:00.474Z" + "updatedAt": "2026-08-27T08:48:14.423Z" }