feat(iios): register EmailAdapter for EMAIL + email-shaped egress provider

AdapterRegistry now backs the EMAIL channel with the email-native EmailAdapter
(RFC threading) instead of the generic WebhookAdapter. Adds EmailProvider — an
email-envelope egress provider ({to,subject,text,html,inReplyTo}) registered for
EMAIL when IIOS_PROVIDER_URL_EMAIL is set (sandbox stays the default). Spec:
signed email → normalized interaction on the email thread; a reply joins the
same thread (one conversation); bad signature → rejected (fail-closed).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-07-03 20:28:35 +05:30
parent ae50f56404
commit f0afca89e3
4 changed files with 91 additions and 8 deletions
@@ -1,6 +1,6 @@
import { Injectable } from '@nestjs/common'; import { Injectable } from '@nestjs/common';
import type { ScopeVector } from '@insignia/iios-contracts'; import type { ScopeVector } from '@insignia/iios-contracts';
import { WebhookAdapter, type ChannelAdapter } from '@insignia/iios-adapter-sdk'; import { WebhookAdapter, EmailAdapter, type ChannelAdapter } from '@insignia/iios-adapter-sdk';
export interface AdapterConfig { export interface AdapterConfig {
secret: string; secret: string;
@@ -8,10 +8,9 @@ export interface AdapterConfig {
} }
/** /**
* Registers channel adapters + their per-channel config (secret + scope). In P5 * Registers channel adapters + their per-channel config (secret + scope). EMAIL uses
* one reference WebhookAdapter backs several channelTypes (WEBHOOK/EMAIL/WHATSAPP) * the email-native EmailAdapter (RFC threading); the rest use the reference
* so the fixtures exercise the same anti-corruption pipeline. Secrets come from * WebhookAdapter. Secrets come from ADAPTER_SECRETS (JSON map); scope defaults to demo.
* ADAPTER_SECRETS (JSON map); scope defaults to the demo scope.
*/ */
@Injectable() @Injectable()
export class AdapterRegistry { export class AdapterRegistry {
@@ -29,8 +28,9 @@ export class AdapterRegistry {
appId: process.env.ADAPTER_APP ?? 'portal-demo', appId: process.env.ADAPTER_APP ?? 'portal-demo',
}; };
for (const channelType of ['WEBHOOK', 'EMAIL', 'WHATSAPP', 'PORTAL']) { for (const channelType of ['WEBHOOK', 'EMAIL', 'WHATSAPP', 'PORTAL']) {
const adapter: ChannelAdapter = channelType === 'EMAIL' ? new EmailAdapter() : new WebhookAdapter(channelType);
this.adapters.set(channelType, { this.adapters.set(channelType, {
adapter: new WebhookAdapter(channelType), adapter,
config: { secret: secrets[channelType] ?? 'dev-adapter-secret', scope }, config: { secret: secrets[channelType] ?? 'dev-adapter-secret', scope },
}); });
} }
@@ -6,7 +6,7 @@ import { resetDb } from '../test-utils/reset-db';
import { ProjectionCursorService } from '../projection/projection-cursor.service'; import { ProjectionCursorService } from '../projection/projection-cursor.service';
import { makeFakePorts } from '@insignia/iios-testkit'; import { makeFakePorts } from '@insignia/iios-testkit';
import { IIOS_EVENTS, type CloudEvent } from '@insignia/iios-contracts'; import { IIOS_EVENTS, type CloudEvent } from '@insignia/iios-contracts';
import { signedFixture, emailFixture } from '@insignia/iios-adapter-sdk'; import { signedFixture, emailFixture, emailInboundFixture, emailReplyFixture } from '@insignia/iios-adapter-sdk';
import { AdapterRegistry } from './adapter.registry'; import { AdapterRegistry } from './adapter.registry';
import { InboundService } from './inbound.service'; import { InboundService } from './inbound.service';
import { OutboundService } from './outbound.service'; import { OutboundService } from './outbound.service';
@@ -81,6 +81,46 @@ describe('Adapter inbound (P5)', () => {
}); });
}); });
describe('Email adapter — inbound + threading', () => {
const receiveEmail = async (payload: object) => {
const { body, headers } = signedFixture(payload, SECRET);
return inbound().receive('EMAIL', body, headers);
};
const normalizeAll = async () => {
for (const ev of await rawReceivedEvents()) await projector().onRawReceived(ev);
};
it('a signed inbound email → a normalized interaction with the body on the email thread', async () => {
const ack = await receiveEmail(emailInboundFixture);
expect(ack.status).toBe('RECEIVED');
await normalizeAll();
const raw = await prisma.iiosInboundRawEvent.findUniqueOrThrow({ where: { id: ack.rawEventId } });
expect(raw.status).toBe('NORMALIZED');
const interaction = await prisma.iiosInteraction.findUniqueOrThrow({ where: { id: raw.interactionId! }, include: { parts: true, thread: true } });
expect(interaction.parts[0]?.bodyText).toBe(emailInboundFixture.text);
expect(interaction.thread?.externalThreadRef).toBe('email:<m1@example.com>');
});
it('a reply email (In-Reply-To) lands on the SAME email thread as the original', async () => {
await receiveEmail(emailInboundFixture);
await receiveEmail(emailReplyFixture);
await normalizeAll();
const interactions = await prisma.iiosInteraction.findMany({ orderBy: { occurredAt: 'asc' } });
expect(interactions).toHaveLength(2);
expect(interactions[0]!.threadId).toBe(interactions[1]!.threadId); // reply joined the thread
expect(await prisma.iiosThread.count()).toBe(1); // one email conversation
});
it('bad signature on EMAIL → REJECTED, nothing ingested (fail-closed)', async () => {
const { body } = signedFixture(emailInboundFixture, SECRET);
await expect(inbound().receive('EMAIL', body, { 'x-iios-signature': 'sha256=bad' })).rejects.toThrow();
expect(await prisma.iiosInboundRawEvent.count({ where: { status: 'REJECTED' } })).toBe(1);
expect(await prisma.iiosInteraction.count()).toBe(0);
});
});
describe('Adapter outbound (P5)', () => { describe('Adapter outbound (P5)', () => {
const outbound = (limit?: number): OutboundService => { const outbound = (limit?: number): OutboundService => {
const s = new OutboundService(asService, new CapabilityBroker(makeFakePorts(), new CapabilityProviderRegistry()), new IdempotencyService(asService)); const s = new OutboundService(asService, new CapabilityBroker(makeFakePorts(), new CapabilityProviderRegistry()), new IdempotencyService(asService));
@@ -2,6 +2,7 @@ import { Injectable, NotFoundException } from '@nestjs/common';
import type { CapabilityProvider } from '@insignia/iios-contracts'; import type { CapabilityProvider } from '@insignia/iios-contracts';
import { SandboxProvider } from './sandbox.provider'; import { SandboxProvider } from './sandbox.provider';
import { HttpProvider } from './http.provider'; import { HttpProvider } from './http.provider';
import { EmailProvider } from './email.provider';
const DEFAULT_CHANNELS = ['WEBHOOK', 'EMAIL', 'WHATSAPP', 'PORTAL']; const DEFAULT_CHANNELS = ['WEBHOOK', 'EMAIL', 'WHATSAPP', 'PORTAL'];
@@ -19,7 +20,9 @@ export class CapabilityProviderRegistry {
for (const ch of DEFAULT_CHANNELS) this.register(new SandboxProvider([ch])); for (const ch of DEFAULT_CHANNELS) this.register(new SandboxProvider([ch]));
for (const ch of DEFAULT_CHANNELS) { for (const ch of DEFAULT_CHANNELS) {
const url = process.env[`IIOS_PROVIDER_URL_${ch}`]; const url = process.env[`IIOS_PROVIDER_URL_${ch}`];
if (url) this.register(new HttpProvider(ch, url)); if (!url) continue;
// EMAIL gets an email-shaped envelope provider; other channels use the generic HTTP one.
this.register(ch === 'EMAIL' ? new EmailProvider(url) : new HttpProvider(ch, url));
} }
} }
@@ -0,0 +1,40 @@
import type { CapabilityProvider, CapabilityRequest, ProviderResult } from '@insignia/iios-contracts';
interface EmailPayload {
subject?: string;
text?: string;
html?: string;
inReplyTo?: string;
}
/**
* Email egress provider: shapes the outbound command into an email envelope
* ({to, subject, text, html, inReplyTo}) and POSTs it to a configured send endpoint
* (a vendor mail API / SMTP relay behind a URL). Activated only when
* IIOS_PROVIDER_URL_EMAIL is set; a transport failure is surfaced as FAILED, never thrown.
*/
export class EmailProvider implements CapabilityProvider {
readonly name = 'email-http';
readonly channelTypes = ['EMAIL'];
readonly capabilities = { canSend: true };
constructor(private readonly url: string) {}
async send(req: CapabilityRequest): Promise<ProviderResult> {
const started = Date.now();
const p = (req.payload ?? {}) as EmailPayload;
const email = { to: req.target, subject: p.subject ?? '(no subject)', text: p.text, html: p.html, inReplyTo: p.inReplyTo };
try {
const res = await fetch(this.url, {
method: 'POST',
headers: { 'content-type': 'application/json', 'idempotency-key': req.idempotencyKey },
body: JSON.stringify(email),
});
const latencyMs = Date.now() - started;
if (!res.ok) return { providerRef: `email-${res.status}`, outcome: 'FAILED', errorCode: `HTTP_${res.status}`, latencyMs };
return { providerRef: `email-${req.idempotencyKey}`, outcome: 'SENT', latencyMs };
} catch (err) {
return { providerRef: 'email-error', outcome: 'FAILED', errorCode: (err as Error).message.slice(0, 60), latencyMs: Date.now() - started };
}
}
}