diff --git a/README.md b/README.md
index 74fc8cba..7afbd998 100644
--- a/README.md
+++ b/README.md
@@ -14,7 +14,7 @@
-
+
@@ -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: {}