refactor(14-01): extract DKV inbox providers into shared inbox/ module
- Move ImapProvider, ExchangeInboxProvider, InboxProvider into apps/api/src/inbox/ - Move InboxConfig/InboxAttachment/InboxEmail into new inbox.types.ts - dkv.types.ts re-exports the moved types so existing DKV imports keep compiling - DKV switches import paths to ../inbox/... and imports InboxModule - Pure move + import-path swap: fetchPdfAttachments and all DKV logic unchanged (D-01/D-02) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,452 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
// eslint-disable-next-line @typescript-eslint/no-require-imports
|
||||
const httpntlm = require('httpntlm') as { post: (opts: any, cb: (err: Error | null, res: any) => void) => void };
|
||||
import type { InboxAttachment, InboxConfig, InboxEmail } from './inbox-provider.interface';
|
||||
import type { InboxProvider } from './inbox-provider.interface';
|
||||
|
||||
const MAX_ATTACHMENT_BYTES = 25 * 1024 * 1024; // 25 MB — T-07-05
|
||||
|
||||
// ─── EWS SOAP namespace constants ────────────────────────────────────────────
|
||||
|
||||
const NS_SOAP = 'http://schemas.xmlsoap.org/soap/envelope/';
|
||||
const NS_TYPES = 'http://schemas.microsoft.com/exchange/services/2006/types';
|
||||
const NS_MESSAGES = 'http://schemas.microsoft.com/exchange/services/2006/messages';
|
||||
|
||||
// ─── SOAP envelope builders ───────────────────────────────────────────────────
|
||||
|
||||
function soapEnvelope(body: string): string {
|
||||
return `<?xml version="1.0" encoding="utf-8"?>
|
||||
<soap:Envelope xmlns:soap="${NS_SOAP}"
|
||||
xmlns:t="${NS_TYPES}"
|
||||
xmlns:m="${NS_MESSAGES}">
|
||||
<soap:Body>${body}</soap:Body>
|
||||
</soap:Envelope>`;
|
||||
}
|
||||
|
||||
/** folderElement: either <t:DistinguishedFolderId Id="inbox"/> or <t:FolderId Id="AAA..."/> */
|
||||
function findItemSoap(folderElement: string, maxResults: number, senderFilter?: string): string {
|
||||
const isReadFilter = `<t:IsEqualTo>
|
||||
<t:FieldURI FieldURI="message:IsRead"/>
|
||||
<t:FieldURIOrConstant><t:Constant Value="false"/></t:FieldURIOrConstant>
|
||||
</t:IsEqualTo>`;
|
||||
|
||||
const restriction = senderFilter
|
||||
? `<m:Restriction>
|
||||
<t:And>
|
||||
${isReadFilter}
|
||||
<t:Contains ContainmentMode="Substring" ContainmentComparison="IgnoreCase">
|
||||
<t:FieldURI FieldURI="message:From"/>
|
||||
<t:Constant Value="${escapeXml(senderFilter)}"/>
|
||||
</t:Contains>
|
||||
</t:And>
|
||||
</m:Restriction>`
|
||||
: `<m:Restriction>${isReadFilter}</m:Restriction>`;
|
||||
|
||||
return soapEnvelope(`
|
||||
<m:FindItem Traversal="Shallow">
|
||||
<m:ItemShape>
|
||||
<t:BaseShape>IdOnly</t:BaseShape>
|
||||
<t:AdditionalProperties>
|
||||
<t:FieldURI FieldURI="item:Subject"/>
|
||||
<t:FieldURI FieldURI="message:InternetMessageId"/>
|
||||
<t:FieldURI FieldURI="message:From"/>
|
||||
<t:FieldURI FieldURI="item:DateTimeReceived"/>
|
||||
<t:FieldURI FieldURI="item:HasAttachments"/>
|
||||
</t:AdditionalProperties>
|
||||
</m:ItemShape>
|
||||
<m:IndexedPageItemView MaxEntriesReturned="${maxResults}" Offset="0" BasePoint="Beginning"/>
|
||||
${restriction}
|
||||
<m:ParentFolderIds>
|
||||
${folderElement}
|
||||
</m:ParentFolderIds>
|
||||
</m:FindItem>`);
|
||||
}
|
||||
|
||||
function getItemSoap(itemIds: string[]): string {
|
||||
const ids = itemIds.map((id) => `<t:ItemId Id="${escapeXml(id)}"/>`).join('');
|
||||
return soapEnvelope(`
|
||||
<m:GetItem>
|
||||
<m:ItemShape>
|
||||
<t:BaseShape>IdOnly</t:BaseShape>
|
||||
<t:AdditionalProperties>
|
||||
<t:FieldURI FieldURI="item:Subject"/>
|
||||
<t:FieldURI FieldURI="message:InternetMessageId"/>
|
||||
<t:FieldURI FieldURI="message:From"/>
|
||||
<t:FieldURI FieldURI="item:DateTimeReceived"/>
|
||||
<t:FieldURI FieldURI="item:Attachments"/>
|
||||
</t:AdditionalProperties>
|
||||
</m:ItemShape>
|
||||
<m:ItemIds>${ids}</m:ItemIds>
|
||||
</m:GetItem>`);
|
||||
}
|
||||
|
||||
function markReadSoap(itemId: string, changeKey: string): string {
|
||||
return soapEnvelope(`
|
||||
<m:UpdateItem MessageDisposition="SaveOnly" ConflictResolution="AlwaysOverwrite">
|
||||
<m:ItemChanges>
|
||||
<t:ItemChange>
|
||||
<t:ItemId Id="${escapeXml(itemId)}" ChangeKey="${escapeXml(changeKey)}"/>
|
||||
<t:Updates>
|
||||
<t:SetItemField>
|
||||
<t:FieldURI FieldURI="message:IsRead"/>
|
||||
<t:Message><t:IsRead>true</t:IsRead></t:Message>
|
||||
</t:SetItemField>
|
||||
</t:Updates>
|
||||
</t:ItemChange>
|
||||
</m:ItemChanges>
|
||||
</m:UpdateItem>`);
|
||||
}
|
||||
|
||||
function getAttachmentSoap(attachmentId: string): string {
|
||||
return soapEnvelope(`
|
||||
<m:GetAttachment>
|
||||
<m:AttachmentIds>
|
||||
<t:AttachmentId Id="${escapeXml(attachmentId)}"/>
|
||||
</m:AttachmentIds>
|
||||
</m:GetAttachment>`);
|
||||
}
|
||||
|
||||
// ─── XML helpers ─────────────────────────────────────────────────────────────
|
||||
|
||||
function escapeXml(s: string): string {
|
||||
return s.replace(/&/g, '&').replace(/</g, '<').replace(/>/g, '>').replace(/"/g, '"');
|
||||
}
|
||||
|
||||
/** Extract all text values of a tag name from an XML string (non-recursive, fast). */
|
||||
function extractAll(xml: string, tag: string): string[] {
|
||||
const results: string[] = [];
|
||||
const open = `<${tag}`;
|
||||
const close = `</${tag}>`;
|
||||
let pos = 0;
|
||||
while (pos < xml.length) {
|
||||
const start = xml.indexOf(open, pos);
|
||||
if (start === -1) break;
|
||||
const end = xml.indexOf(close, start);
|
||||
if (end === -1) break;
|
||||
// Get inner text (content between > and </tag>)
|
||||
const innerStart = xml.indexOf('>', start) + 1;
|
||||
results.push(xml.slice(innerStart, end));
|
||||
pos = end + close.length;
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
/** Extract first value of attribute from a tag. */
|
||||
function extractAttr(xml: string, tag: string, attr: string): string {
|
||||
const tagStart = xml.indexOf(`<${tag}`);
|
||||
if (tagStart === -1) return '';
|
||||
const tagEnd = xml.indexOf('>', tagStart);
|
||||
const tagStr = xml.slice(tagStart, tagEnd + 1);
|
||||
const attrMatch = tagStr.match(new RegExp(`${attr}="([^"]*)"`));
|
||||
return attrMatch ? attrMatch[1] : '';
|
||||
}
|
||||
|
||||
/** Extract all matching tag attributes from repeated elements. */
|
||||
function extractAttrs(xml: string, tag: string, attr: string): string[] {
|
||||
const results: string[] = [];
|
||||
const open = `<${tag}`;
|
||||
let pos = 0;
|
||||
while (pos < xml.length) {
|
||||
const start = xml.indexOf(open, pos);
|
||||
if (start === -1) break;
|
||||
const tagEnd = xml.indexOf('>', start);
|
||||
const tagStr = xml.slice(start, tagEnd + 1);
|
||||
const match = tagStr.match(new RegExp(`${attr}="([^"]*)"`));
|
||||
if (match) results.push(match[1]);
|
||||
pos = tagEnd + 1;
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
const DISTINGUISHED_FOLDER_MAP: Record<string, string> = {
|
||||
inbox: 'inbox',
|
||||
posteingang: 'inbox',
|
||||
deleteditems: 'deleteditems',
|
||||
gelöschteelemente: 'deleteditems',
|
||||
trash: 'deleteditems',
|
||||
sentitems: 'sentitems',
|
||||
gesendet: 'sentitems',
|
||||
drafts: 'drafts',
|
||||
entwürfe: 'drafts',
|
||||
junk: 'junkemail',
|
||||
junkemail: 'junkemail',
|
||||
spam: 'junkemail',
|
||||
};
|
||||
|
||||
/** Returns the DistinguishedFolderId if folder is a well-known name, else null (→ needs FindFolder). */
|
||||
function resolveDistinguishedFolder(folder?: string): string | null {
|
||||
const name = (folder ?? 'INBOX').toLowerCase().replace(/[\s_-]/g, '');
|
||||
return DISTINGUISHED_FOLDER_MAP[name] ?? null;
|
||||
}
|
||||
|
||||
/** Build FindFolder SOAP to resolve a subfolder by display name under a parent DistinguishedFolderId. */
|
||||
function findFolderSoap(parentDistinguishedId: string, displayName: string): string {
|
||||
return soapEnvelope(`
|
||||
<m:FindFolder Traversal="Deep">
|
||||
<m:FolderShape>
|
||||
<t:BaseShape>IdOnly</t:BaseShape>
|
||||
</m:FolderShape>
|
||||
<m:Restriction>
|
||||
<t:IsEqualTo>
|
||||
<t:FieldURI FieldURI="folder:DisplayName"/>
|
||||
<t:FieldURIOrConstant>
|
||||
<t:Constant Value="${escapeXml(displayName)}"/>
|
||||
</t:FieldURIOrConstant>
|
||||
</t:IsEqualTo>
|
||||
</m:Restriction>
|
||||
<m:ParentFolderIds>
|
||||
<t:DistinguishedFolderId Id="${escapeXml(parentDistinguishedId)}"/>
|
||||
</m:ParentFolderIds>
|
||||
</m:FindFolder>`);
|
||||
}
|
||||
|
||||
// ─── NTLM HTTP helper ────────────────────────────────────────────────────────
|
||||
|
||||
interface NtlmOptions {
|
||||
url: string;
|
||||
username: string;
|
||||
password: string;
|
||||
domain: string;
|
||||
workstation: string;
|
||||
body: string;
|
||||
headers: Record<string, string>;
|
||||
rejectUnauthorized?: boolean;
|
||||
}
|
||||
|
||||
function ntlmPost(opts: NtlmOptions): Promise<{ statusCode: number; body: string }> {
|
||||
return new Promise((resolve, reject) => {
|
||||
(httpntlm as any).post(opts, (err: Error | null, res: any) => {
|
||||
if (err) return reject(err);
|
||||
resolve({ statusCode: res.statusCode, body: res.body?.toString('utf-8') ?? '' });
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
// ─── Provider ────────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* ExchangeInboxProvider — NTLM-authenticated EWS via raw SOAP over httpntlm.
|
||||
*
|
||||
* Replaces the ews-javascript-api approach which only supports Basic Auth.
|
||||
* Uses httpntlm to perform the NTLM challenge-response handshake transparently.
|
||||
*
|
||||
* Security:
|
||||
* - T-07-03: credentials never logged — only generic error messages
|
||||
* - T-07-05: attachment size checked before buffering (PDF-bomb mitigation)
|
||||
* - EWS XML is escaped before insertion into SOAP envelopes
|
||||
*/
|
||||
@Injectable()
|
||||
export class ExchangeInboxProvider implements InboxProvider {
|
||||
private readonly logger = new Logger(ExchangeInboxProvider.name);
|
||||
|
||||
async fetchPdfAttachments(config: InboxConfig): Promise<InboxEmail[]> {
|
||||
try {
|
||||
return await this.fetchViaEws(config);
|
||||
} catch (error) {
|
||||
this.logger.error(`EWS inbox fetch failed: ${(error as Error).message}`);
|
||||
return [];
|
||||
}
|
||||
}
|
||||
|
||||
async testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }> {
|
||||
try {
|
||||
const folderElement = await this.resolveFolderElement(config);
|
||||
const soap = findItemSoap(folderElement, 1);
|
||||
const res = await this.ewsPost(config, soap, 'FindItem');
|
||||
|
||||
if (res.statusCode === 401) {
|
||||
return { success: false, message: '401 Unauthorized — credentials rejected or NTLM not allowed' };
|
||||
}
|
||||
if (res.statusCode !== 200) {
|
||||
return { success: false, message: `HTTP ${res.statusCode}` };
|
||||
}
|
||||
if (res.body.includes('ResponseClass="Error"')) {
|
||||
const msg = extractAll(res.body, 'm:MessageText')[0] ?? extractAll(res.body, 'MessageText')[0] ?? 'EWS error';
|
||||
return { success: false, message: msg };
|
||||
}
|
||||
return { success: true };
|
||||
} catch (err) {
|
||||
const message = (err as Error).message;
|
||||
this.logger.error(`EWS connection test failed: ${message}`);
|
||||
return { success: false, message };
|
||||
}
|
||||
}
|
||||
|
||||
// ─── Private ───────────────────────────────────────────────────────────────
|
||||
|
||||
/** Resolve folder config to a EWS ParentFolderIds XML element.
|
||||
* Well-known names → <t:DistinguishedFolderId>.
|
||||
* Custom names (e.g. "DKV" or "INBOX/DKV") → FindFolder deep search under msgfolderroot → <t:FolderId>.
|
||||
* Searches msgfolderroot (full mailbox) so folders at root level are found, not just inbox subfolders.
|
||||
*/
|
||||
private async resolveFolderElement(config: InboxConfig): Promise<string> {
|
||||
const folderCfg = config.folder ?? 'INBOX';
|
||||
// Strip leading "INBOX/" prefix — EWS FindFolder searches deep, name alone suffices
|
||||
const displayName = folderCfg.replace(/^INBOX\//i, '').trim();
|
||||
const distinguished = resolveDistinguishedFolder(displayName);
|
||||
if (distinguished) {
|
||||
return `<t:DistinguishedFolderId Id="${escapeXml(distinguished)}"/>`;
|
||||
}
|
||||
// Custom subfolder: search entire mailbox (msgfolderroot) so top-level folders are found too
|
||||
const ffSoap = findFolderSoap('msgfolderroot', displayName);
|
||||
const ffRes = await this.ewsPost(config, ffSoap, 'FindFolder');
|
||||
if (ffRes.statusCode !== 200) {
|
||||
this.logger.warn(`EWS FindFolder HTTP ${ffRes.statusCode} for "${displayName}" — falling back to inbox`);
|
||||
return `<t:DistinguishedFolderId Id="inbox"/>`;
|
||||
}
|
||||
this.logger.debug(`EWS FindFolder response for "${displayName}": ${ffRes.body.slice(0, 500)}`);
|
||||
const folderId = extractAttr(ffRes.body, 't:FolderId', 'Id');
|
||||
if (!folderId) {
|
||||
this.logger.warn(`EWS FindFolder: subfolder "${displayName}" not found under msgfolderroot — falling back to inbox`);
|
||||
return `<t:DistinguishedFolderId Id="inbox"/>`;
|
||||
}
|
||||
this.logger.log(`EWS FindFolder: resolved "${displayName}" → FolderId ${folderId.slice(0, 20)}…`);
|
||||
return `<t:FolderId Id="${escapeXml(folderId)}"/>`;
|
||||
}
|
||||
|
||||
private async fetchViaEws(config: InboxConfig): Promise<InboxEmail[]> {
|
||||
const folderElement = await this.resolveFolderElement(config);
|
||||
|
||||
// 1. FindItem — get IDs of emails with attachments
|
||||
const findSoap = findItemSoap(folderElement, 50, config.senderFilter);
|
||||
const findRes = await this.ewsPost(config, findSoap, 'FindItem');
|
||||
if (findRes.statusCode !== 200) {
|
||||
this.logger.warn(`EWS FindItem returned HTTP ${findRes.statusCode}`);
|
||||
return [];
|
||||
}
|
||||
|
||||
// Parse item IDs and HasAttachments flag from FindItem response
|
||||
const rawIds = extractAttrs(findRes.body, 't:ItemId', 'Id');
|
||||
if (rawIds.length === 0) return [];
|
||||
|
||||
// Only fetch items that have attachments
|
||||
const hasAttachFlags = extractAll(findRes.body, 't:HasAttachments');
|
||||
const itemIds = rawIds.filter((_, i) => hasAttachFlags[i] === 'true');
|
||||
if (itemIds.length === 0) return [];
|
||||
|
||||
const results: InboxEmail[] = [];
|
||||
|
||||
// 2. GetItem in batches of 10 to load attachment metadata
|
||||
for (let i = 0; i < itemIds.length; i += 10) {
|
||||
const batch = itemIds.slice(i, i + 10);
|
||||
const getRes = await this.ewsPost(config, getItemSoap(batch), 'GetItem');
|
||||
if (getRes.statusCode !== 200) continue;
|
||||
|
||||
// Parse each Message element from GetItem response
|
||||
const messageBlocks = this.splitMessageBlocks(getRes.body);
|
||||
|
||||
for (const block of messageBlocks) {
|
||||
const uid = extractAttr(block, 't:ItemId', 'Id');
|
||||
const changeKey = extractAttr(block, 't:ItemId', 'ChangeKey');
|
||||
const subject = extractAll(block, 't:Subject')[0] ?? '';
|
||||
const messageId = extractAll(block, 't:InternetMessageId')[0] ?? '';
|
||||
const from = (extractAttr(block, 't:Mailbox', 'SmtpAddress') ||
|
||||
extractAll(block, 't:EmailAddress')[0]) ?? '';
|
||||
const dateStr = extractAll(block, 't:DateTimeReceived')[0] ?? '';
|
||||
const date = dateStr ? new Date(dateStr) : new Date();
|
||||
|
||||
// Filter by sender if provided (client-side fallback for case-sensitivity)
|
||||
if (config.senderFilter && from &&
|
||||
!from.toLowerCase().includes(config.senderFilter.toLowerCase())) {
|
||||
continue;
|
||||
}
|
||||
|
||||
// Collect PDF attachment IDs by iterating each FileAttachment block.
|
||||
// extractAttrs(block, 't:FileAttachment', 'Id') was wrong — the Id lives
|
||||
// inside a child <t:AttachmentId Id="..."/> not on the FileAttachment tag.
|
||||
const attachmentIds: string[] = [];
|
||||
const FA_OPEN = '<t:FileAttachment>';
|
||||
const FA_CLOSE = '</t:FileAttachment>';
|
||||
let faPos = 0;
|
||||
while (faPos < block.length) {
|
||||
const faStart = block.indexOf(FA_OPEN, faPos);
|
||||
if (faStart === -1) break;
|
||||
const faEnd = block.indexOf(FA_CLOSE, faStart);
|
||||
if (faEnd === -1) break;
|
||||
const faBlock = block.slice(faStart, faEnd);
|
||||
const attId = extractAttr(faBlock, 't:AttachmentId', 'Id');
|
||||
const ct = (extractAll(faBlock, 't:ContentType')[0] ?? '').toLowerCase();
|
||||
const nm = (extractAll(faBlock, 't:Name')[0] ?? '').toLowerCase();
|
||||
if (attId && (ct.includes('pdf') || nm.endsWith('.pdf'))) {
|
||||
attachmentIds.push(attId);
|
||||
}
|
||||
faPos = faEnd + FA_CLOSE.length;
|
||||
}
|
||||
|
||||
if (attachmentIds.length === 0) continue;
|
||||
|
||||
// 3. GetAttachment for each PDF
|
||||
const attachments: InboxAttachment[] = [];
|
||||
for (const attId of attachmentIds) {
|
||||
const attRes = await this.ewsPost(config, getAttachmentSoap(attId), 'GetAttachment');
|
||||
if (attRes.statusCode !== 200) continue;
|
||||
|
||||
const name = extractAll(attRes.body, 't:Name')[0] ?? 'attachment.pdf';
|
||||
const content = extractAll(attRes.body, 't:Content')[0] ?? '';
|
||||
if (!content) continue;
|
||||
|
||||
const buffer = Buffer.from(content, 'base64');
|
||||
if (buffer.length > MAX_ATTACHMENT_BYTES) {
|
||||
this.logger.warn(`Skipped EWS attachment "${name}" — exceeds size limit (T-07-05)`);
|
||||
continue;
|
||||
}
|
||||
attachments.push({ filename: name, contentType: 'application/pdf', buffer });
|
||||
}
|
||||
|
||||
if (attachments.length === 0) continue;
|
||||
|
||||
results.push({ uid, messageId, subject, from, date, attachments });
|
||||
|
||||
// Mark as read so subsequent polls skip this message (mirrors IMAP \Seen flag).
|
||||
if (uid && changeKey) {
|
||||
try {
|
||||
await this.ewsPost(config, markReadSoap(uid, changeKey), 'UpdateItem');
|
||||
} catch {
|
||||
// Non-fatal: message will reappear on next poll but IsRead filter will catch it
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return results;
|
||||
}
|
||||
|
||||
/** Split a GetItem response body into per-message XML blocks. */
|
||||
private splitMessageBlocks(xml: string): string[] {
|
||||
const blocks: string[] = [];
|
||||
const open = '<m:Items>';
|
||||
const close = '</m:Items>';
|
||||
let pos = 0;
|
||||
while (pos < xml.length) {
|
||||
const start = xml.indexOf(open, pos);
|
||||
if (start === -1) break;
|
||||
const end = xml.indexOf(close, start);
|
||||
if (end === -1) break;
|
||||
blocks.push(xml.slice(start + open.length, end));
|
||||
pos = end + close.length;
|
||||
}
|
||||
// If no <m:Items> blocks, treat whole response as one block
|
||||
return blocks.length > 0 ? blocks : [xml];
|
||||
}
|
||||
|
||||
/** POST a SOAP body to the EWS endpoint using NTLM authentication. */
|
||||
private async ewsPost(
|
||||
config: InboxConfig,
|
||||
soap: string,
|
||||
action: string,
|
||||
): Promise<{ statusCode: number; body: string }> {
|
||||
const opts: NtlmOptions = {
|
||||
url: config.host, // must be full EWS URL: https://server/EWS/Exchange.asmx
|
||||
username: config.username ?? '',
|
||||
password: config.password ?? '',
|
||||
domain: config.domain ?? '',
|
||||
workstation: '',
|
||||
body: soap,
|
||||
headers: {
|
||||
'Content-Type': 'text/xml; charset=utf-8',
|
||||
'SOAPAction': `"http://schemas.microsoft.com/exchange/services/2006/messages/${action}"`,
|
||||
},
|
||||
};
|
||||
return ntlmPost(opts);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,243 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import { ImapFlow, MessageStructureObject } from 'imapflow';
|
||||
import type { InboxAttachment, InboxConfig, InboxEmail } from './inbox-provider.interface';
|
||||
import type { InboxProvider } from './inbox-provider.interface';
|
||||
|
||||
/**
|
||||
* Max attachment size (bytes) accepted before buffering.
|
||||
* Prevents PDF-bomb DoS (T-07-05 — Research Security Domain).
|
||||
* 25 MB covers any realistic DKV invoice; larger files are skipped.
|
||||
*/
|
||||
const MAX_ATTACHMENT_BYTES = 25 * 1024 * 1024; // 25 MB
|
||||
|
||||
/**
|
||||
* Converts a Node.js Readable stream into a Buffer.
|
||||
* Accumulates chunks up to MAX_ATTACHMENT_BYTES; throws if limit exceeded.
|
||||
*/
|
||||
async function streamToBuffer(
|
||||
stream: NodeJS.ReadableStream,
|
||||
): Promise<Buffer> {
|
||||
return new Promise<Buffer>((resolve, reject) => {
|
||||
const chunks: Buffer[] = [];
|
||||
let total = 0;
|
||||
|
||||
stream.on('data', (chunk: Buffer) => {
|
||||
total += chunk.length;
|
||||
if (total > MAX_ATTACHMENT_BYTES) {
|
||||
// Destroy the stream to prevent further data emission
|
||||
(stream as any).destroy?.();
|
||||
reject(
|
||||
new Error(
|
||||
`Attachment exceeds maximum allowed size of ${MAX_ATTACHMENT_BYTES} bytes (T-07-05)`,
|
||||
),
|
||||
);
|
||||
return;
|
||||
}
|
||||
chunks.push(Buffer.from(chunk));
|
||||
});
|
||||
stream.on('end', () => resolve(Buffer.concat(chunks)));
|
||||
stream.on('error', reject);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Recursively collects all PDF body-part IDs from a message's MIME structure.
|
||||
*
|
||||
* imapflow represents the MIME tree as nested `MessageStructureObject` nodes.
|
||||
* Each node has a `type` (full MIME type, e.g. "application/pdf") and a `part`
|
||||
* ID that can be passed to `client.download()`.
|
||||
*
|
||||
* @param node Root or child MIME structure node
|
||||
* @param parts Accumulator (pass empty array on first call)
|
||||
*/
|
||||
function collectPdfParts(
|
||||
node: MessageStructureObject | undefined,
|
||||
parts: string[] = [],
|
||||
): string[] {
|
||||
if (!node) return parts;
|
||||
|
||||
const type = node.type?.toLowerCase() ?? '';
|
||||
// Some mail clients (e.g. Outlook) send PDFs as application/octet-stream.
|
||||
// Fall back to checking the filename from Content-Disposition or Content-Type parameters.
|
||||
const dispositionFilename =
|
||||
((node as any).disposition?.parameters?.filename as string | undefined)?.toLowerCase() ?? '';
|
||||
const typeFilename =
|
||||
((node as any).parameters?.name as string | undefined)?.toLowerCase() ?? '';
|
||||
const looksLikePdf =
|
||||
type === 'application/pdf' ||
|
||||
(type === 'application/octet-stream' &&
|
||||
(dispositionFilename.endsWith('.pdf') || typeFilename.endsWith('.pdf')));
|
||||
|
||||
if (looksLikePdf && node.part !== undefined) {
|
||||
parts.push(node.part);
|
||||
}
|
||||
|
||||
// Recurse into multipart children
|
||||
if (Array.isArray(node.childNodes)) {
|
||||
for (const child of node.childNodes) {
|
||||
collectPdfParts(child, parts);
|
||||
}
|
||||
}
|
||||
|
||||
return parts;
|
||||
}
|
||||
|
||||
/**
|
||||
* IMAP inbox provider using the imapflow library.
|
||||
*
|
||||
* Implements InboxProvider so it is interchangeable with ExchangeInboxProvider.
|
||||
*
|
||||
* Security:
|
||||
* - T-07-03: ImapFlow constructed with `logger: false` (no credential logging)
|
||||
* - T-07-03: Error messages are generic — credentials never appear in logs
|
||||
* - T-07-05: Max attachment size enforced before buffering (PDF-bomb mitigation)
|
||||
*
|
||||
* Pitfall avoidance:
|
||||
* - Pitfall 1: `fetchAll()` is called BEFORE any `download()` calls.
|
||||
* Never call `download()` inside a `client.fetch()` async iterator —
|
||||
* that deadlocks the IMAP connection.
|
||||
*/
|
||||
@Injectable()
|
||||
export class ImapProvider implements InboxProvider {
|
||||
private readonly logger = new Logger(ImapProvider.name);
|
||||
|
||||
/**
|
||||
* Connects to the IMAP server, locks the configured folder, searches for
|
||||
* emails from the configured sender, and downloads PDF attachments.
|
||||
*
|
||||
* @param config Decrypted inbox connection parameters
|
||||
* @returns Emails with PDF attachments as Buffers (empty array on error)
|
||||
*/
|
||||
async fetchPdfAttachments(config: InboxConfig): Promise<InboxEmail[]> {
|
||||
const client = this.buildClient(config);
|
||||
|
||||
await client.connect();
|
||||
// lock is declared outside try so the type is available in finally, but
|
||||
// getMailboxLock() itself is inside the try so that a lock failure still
|
||||
// triggers client.logout() — preventing a connection leak (CR-03).
|
||||
let lock: Awaited<ReturnType<typeof client.getMailboxLock>> | null = null;
|
||||
const results: InboxEmail[] = [];
|
||||
|
||||
try {
|
||||
lock = await client.getMailboxLock(config.folder || 'INBOX');
|
||||
// Search only UNSEEN emails to avoid reprocessing already-handled messages.
|
||||
// Combine with sender filter when configured.
|
||||
const searchQuery: Record<string, unknown> = { seen: false };
|
||||
if (config.senderFilter) {
|
||||
searchQuery.from = config.senderFilter;
|
||||
}
|
||||
const uids = await client.search(searchQuery, { uid: true });
|
||||
|
||||
if (!uids || uids.length === 0) {
|
||||
return results;
|
||||
}
|
||||
|
||||
// Fetch envelope + body structure for all matching UIDs
|
||||
// IMPORTANT: Must call fetchAll() and complete BEFORE any download() calls
|
||||
// Pitfall 1: calling download() inside a fetch() iterator deadlocks the connection
|
||||
const messages = await client.fetchAll(
|
||||
uids.join(','),
|
||||
{ envelope: true, bodyStructure: true },
|
||||
{ uid: true },
|
||||
);
|
||||
|
||||
for (const msg of messages) {
|
||||
// Find all PDF parts in the message MIME tree
|
||||
const pdfPartIds = collectPdfParts(msg.bodyStructure);
|
||||
const attachments: InboxAttachment[] = [];
|
||||
|
||||
for (const partId of pdfPartIds) {
|
||||
try {
|
||||
const { content } = await client.download(
|
||||
String(msg.uid),
|
||||
partId,
|
||||
{ uid: true },
|
||||
);
|
||||
|
||||
// streamToBuffer enforces MAX_ATTACHMENT_BYTES (T-07-05)
|
||||
const buf = await streamToBuffer(content);
|
||||
const filename =
|
||||
`attachment-${msg.uid}-${partId}.pdf`;
|
||||
|
||||
attachments.push({
|
||||
filename,
|
||||
contentType: 'application/pdf',
|
||||
buffer: buf,
|
||||
});
|
||||
} catch (partErr) {
|
||||
// Skip oversized or unreadable parts; log generic message (T-07-03)
|
||||
this.logger.warn(
|
||||
`Skipped IMAP attachment part ${partId} for UID ${String(msg.uid)}: ${(partErr as Error).message}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if (attachments.length > 0) {
|
||||
results.push({
|
||||
uid: msg.uid!,
|
||||
messageId: msg.envelope?.messageId ?? '',
|
||||
subject: msg.envelope?.subject ?? '',
|
||||
from: msg.envelope?.from?.[0]?.address ?? '',
|
||||
date: msg.envelope?.date ?? new Date(),
|
||||
attachments,
|
||||
});
|
||||
// Mark as read so subsequent polls skip this message (UNSEEN filter above).
|
||||
try {
|
||||
await client.messageFlagsAdd(String(msg.uid), ['\\Seen'], { uid: true });
|
||||
} catch {
|
||||
// Non-fatal: message will simply appear again on next poll
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
lock?.release();
|
||||
await client.logout();
|
||||
}
|
||||
|
||||
return results;
|
||||
}
|
||||
|
||||
/**
|
||||
* Tests whether the IMAP connection can be established.
|
||||
* Connects and immediately logs out without selecting any folder.
|
||||
*
|
||||
* @param config Decrypted inbox connection parameters
|
||||
* @returns true on success, false on any network/auth failure
|
||||
*/
|
||||
async testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }> {
|
||||
const client = this.buildClient(config);
|
||||
try {
|
||||
await client.connect();
|
||||
await client.logout();
|
||||
return { success: true };
|
||||
} catch (err) {
|
||||
const message = (err as Error).message;
|
||||
// T-07-03: log without credentials; return message to admin for diagnosis
|
||||
this.logger.error(`IMAP connection test failed: ${message}`);
|
||||
return { success: false, message };
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Constructs an ImapFlow client from InboxConfig.
|
||||
*
|
||||
* Security (T-07-03):
|
||||
* - `logger: false` suppresses imapflow verbose logging (which includes credentials)
|
||||
* - `auth` is only set when username is present (supports servers with no auth)
|
||||
*/
|
||||
private buildClient(config: InboxConfig): ImapFlow {
|
||||
return new ImapFlow({
|
||||
host: config.host,
|
||||
port: config.port,
|
||||
// ssl-tls = implicit TLS (port 993); starttls = STARTTLS upgrade (port 143)
|
||||
secure: config.encryption === 'ssl-tls',
|
||||
requireTLS: config.encryption === 'starttls',
|
||||
auth:
|
||||
config.username
|
||||
? { user: config.username, pass: config.password ?? '' }
|
||||
: undefined,
|
||||
// T-07-03: suppress imapflow verbose logs — they include auth credentials
|
||||
logger: false,
|
||||
} as any);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
import { InboxAttachment, InboxConfig, InboxEmail } from './inbox.types';
|
||||
|
||||
// Re-export types for downstream consumers that import from this module
|
||||
export type { InboxAttachment, InboxConfig, InboxEmail };
|
||||
|
||||
/**
|
||||
* Abstract inbox provider contract — implemented by ImapProvider and
|
||||
* ExchangeInboxProvider. Consumers (DkvService, DkvSchedulerService) depend
|
||||
* only on this interface, not on the concrete implementations.
|
||||
*
|
||||
* Security:
|
||||
* - T-07-03: Implementations MUST NOT log credential values
|
||||
* - InboxConfig carries decrypted credentials — never persist, never serialize
|
||||
*/
|
||||
export interface InboxProvider {
|
||||
/**
|
||||
* Connects to the configured inbox, searches for emails from the configured
|
||||
* sender filter, downloads PDF attachments, and returns all matching emails.
|
||||
*
|
||||
* @param config Decrypted inbox connection parameters
|
||||
* @returns Array of emails with their PDF attachments as Buffers
|
||||
*/
|
||||
fetchPdfAttachments(config: InboxConfig): Promise<InboxEmail[]>;
|
||||
|
||||
/**
|
||||
* Tests whether the inbox connection can be established.
|
||||
*
|
||||
* @param config Decrypted inbox connection parameters
|
||||
* @returns { success: true } on success, { success: false, message } on failure
|
||||
*/
|
||||
testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }>;
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { ExchangeInboxProvider } from './exchange-inbox.provider';
|
||||
import { ImapProvider } from './imap.provider';
|
||||
|
||||
/**
|
||||
* Shared inbox-connection module (IMAP/Exchange).
|
||||
*
|
||||
* Extracted out of DKV (D-01, 14-CONTEXT.md) so any module that needs to
|
||||
* poll a mailbox — e.g. the Phase-14 EmailAlertAdapter (Plan 14-03) — can
|
||||
* import InboxModule and inject ImapProvider / ExchangeInboxProvider without
|
||||
* depending on DKV. DKV imports this module unchanged in behavior; each
|
||||
* module keeps its own, independent mailbox configuration (D-03) — only the
|
||||
* connection code is shared.
|
||||
*/
|
||||
@Module({
|
||||
providers: [ImapProvider, ExchangeInboxProvider],
|
||||
exports: [ImapProvider, ExchangeInboxProvider],
|
||||
})
|
||||
export class InboxModule {}
|
||||
@@ -0,0 +1,73 @@
|
||||
/**
|
||||
* Shared types for the inbox module (IMAP/Exchange connection mechanics).
|
||||
*
|
||||
* Moved verbatim out of apps/api/src/dkv/dkv.types.ts as part of the
|
||||
* Phase 14 inbox-module extraction (D-01/D-02, 14-CONTEXT.md). DKV re-exports
|
||||
* these from ../inbox/inbox.types so existing DKV imports keep compiling.
|
||||
*/
|
||||
|
||||
/**
|
||||
* Inbox connection configuration (supports IMAP and Exchange).
|
||||
* Mirrors DkvModuleConfig from Prisma but with decrypted credentials.
|
||||
*/
|
||||
export interface InboxConfig {
|
||||
/** 'imap' | 'exchange' */
|
||||
protocol: string;
|
||||
host: string;
|
||||
port: number;
|
||||
/** Optional username (may be empty for anonymous or Exchange-with-integrated-auth) */
|
||||
username?: string;
|
||||
/** Decrypted password (never persisted in this form) */
|
||||
password?: string;
|
||||
/** 'none' | 'starttls' | 'ssl-tls' */
|
||||
encryption: string;
|
||||
/** IMAP folder to monitor (e.g., "INBOX") */
|
||||
folder: string;
|
||||
/** Sender filter: only process emails from this address */
|
||||
senderFilter?: string;
|
||||
/** Exchange only: Windows domain (e.g. "CONTOSO") */
|
||||
domain?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* A single email attachment, already downloaded as a Buffer.
|
||||
*/
|
||||
export interface InboxAttachment {
|
||||
filename: string;
|
||||
contentType: string;
|
||||
buffer: Buffer;
|
||||
}
|
||||
|
||||
/**
|
||||
* A fetched email message with its PDF attachments.
|
||||
*/
|
||||
export interface InboxEmail {
|
||||
/** IMAP UID or Exchange item ID */
|
||||
uid: number | string;
|
||||
/** RFC 2822 Message-ID header */
|
||||
messageId: string;
|
||||
subject: string;
|
||||
from: string;
|
||||
date: Date;
|
||||
attachments: InboxAttachment[];
|
||||
}
|
||||
|
||||
/**
|
||||
* A single fetched message with subject + HTML/text body (no attachments).
|
||||
*
|
||||
* Added additively (D-02) as the seam Plan 14-03's EmailAlertAdapter consumes
|
||||
* via InboxProvider.fetchMessages — DKV's fetchPdfAttachments is unaffected.
|
||||
*/
|
||||
export interface InboxMessage {
|
||||
/** IMAP UID or Exchange item ID */
|
||||
uid: number | string;
|
||||
/** RFC 2822 Message-ID header */
|
||||
messageId: string;
|
||||
subject: string;
|
||||
from: string;
|
||||
date: Date;
|
||||
/** HTML body part, or null when the message has no text/html part */
|
||||
bodyHtml: string | null;
|
||||
/** Plain-text body part; empty string when the message has no text/plain part */
|
||||
bodyText: string;
|
||||
}
|
||||
Reference in New Issue
Block a user