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 @@
-
+
@@ -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 @@
-
+
@@ -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"
}