feat(14-01): add fetchMessages() to InboxProvider + both implementations

- Add InboxMessage type (subject + html/text body) and fetchMessages() to
  the InboxProvider interface, ImapProvider, and ExchangeInboxProvider
- IMAP: findBodyParts() walks the MIME tree for first text/html + text/plain
  parts, reusing the connect/lock/search/fetchAll skeleton; marks \Seen
- EWS: new getItemBodySoap() requests item:Body, extracts BodyType via the
  existing extractAttr/extractAll helpers, marks IsRead via markReadSoap
- Net-new spec coverage (imap.provider.spec.ts, exchange-inbox.provider.spec.ts)
  mocking ImapFlow and httpntlm.post (via require.cache stub, since httpntlm
  is loaded with a raw require() that vi.mock cannot intercept)
- fetchPdfAttachments untouched in both providers (D-02); full API suite
  (301 tests) + tsc --noEmit stay green

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-23 13:09:45 +02:00
parent eb668fd5c8
commit c404954bee
5 changed files with 640 additions and 4 deletions
@@ -0,0 +1,207 @@
import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import type { ExchangeInboxProvider as ExchangeInboxProviderType } from './exchange-inbox.provider';
import type { InboxConfig } from './inbox.types';
/**
* ExchangeInboxProvider.fetchMessages spec (Plan 14-01, Task 2, TDD).
*
* Net-new coverage for the additive fetchMessages() method (D-02) — the seam
* the Plan 14-03 EmailAlertAdapter consumes. fetchPdfAttachments is NOT
* touched or retested here.
*
* httpntlm is loaded via a raw CommonJS `require('httpntlm')` inside
* exchange-inbox.provider.ts (not an ES import) — this bypasses Vitest's
* `vi.mock` ESM-graph interception entirely (verified: `vi.mock('httpntlm',
* ...)` has zero effect on a `require('httpntlm')` call, since Node's real
* `require` executes it directly against the module cache). Instead, this
* spec pre-seeds Node's own `require.cache` for the resolved httpntlm path
* with a stub `{ post }` BEFORE the provider module is first loaded (in
* `beforeAll`, via a dynamic `import()` — kept out of top-level module scope
* so this file stays compatible with the project's commonjs tsconfig, which
* rejects top-level `await`/`import.meta`), so the provider's own
* `require('httpntlm')` call resolves to the stub.
*/
const httpntlmPath = require.resolve('httpntlm');
const httpntlmPost = vi.fn(
(_opts: any, cb: (err: Error | null, res: any) => void) => {
cb(new Error('httpntlmPost not configured for this test'), null);
},
);
require.cache[httpntlmPath] = {
id: httpntlmPath,
filename: httpntlmPath,
loaded: true,
exports: { post: httpntlmPost },
} as any;
let ExchangeInboxProvider: typeof ExchangeInboxProviderType;
beforeAll(async () => {
({ ExchangeInboxProvider } = await import('./exchange-inbox.provider'));
});
const BASE_CONFIG: InboxConfig = {
protocol: 'exchange',
host: 'https://mail.example.com/EWS/Exchange.asmx',
port: 443,
username: 'alerts',
password: 'secret',
encryption: 'ssl-tls',
folder: 'INBOX',
domain: 'CONTOSO',
};
function soapResponse(statusCode: number, body: string) {
return { statusCode, body: Buffer.from(body, 'utf-8') };
}
const FIND_ITEM_RESPONSE = `<?xml version="1.0"?>
<s:Envelope><s:Body>
<m:FindItemResponseMessage ResponseClass="Success">
<m:RootFolder>
<t:Items>
<t:Message>
<t:ItemId Id="AAA1" ChangeKey="CK1"/>
<t:HasAttachments>false</t:HasAttachments>
</t:Message>
</t:Items>
</m:RootFolder>
</m:FindItemResponseMessage>
</s:Body></s:Envelope>`;
function getItemHtmlResponse(): string {
return `<?xml version="1.0"?>
<s:Envelope><s:Body>
<m:GetItemResponseMessage ResponseClass="Success">
<m:Items>
<t:Message>
<t:ItemId Id="AAA1" ChangeKey="CK1"/>
<t:Subject>Neue Ausschreibung verfuegbar</t:Subject>
<t:InternetMessageId>msg-1@example.com</t:InternetMessageId>
<t:From><t:Mailbox><t:EmailAddress>noreply@vergabeportal.de</t:EmailAddress></t:Mailbox></t:From>
<t:DateTimeReceived>2026-07-20T08:00:00Z</t:DateTimeReceived>
<t:Body BodyType="HTML">Hallo HTML Body</t:Body>
</t:Message>
</m:Items>
</m:GetItemResponseMessage>
</s:Body></s:Envelope>`;
}
function getItemTextResponse(): string {
return `<?xml version="1.0"?>
<s:Envelope><s:Body>
<m:GetItemResponseMessage ResponseClass="Success">
<m:Items>
<t:Message>
<t:ItemId Id="AAA1" ChangeKey="CK1"/>
<t:Subject>Neue Ausschreibung verfuegbar</t:Subject>
<t:InternetMessageId>msg-1@example.com</t:InternetMessageId>
<t:From><t:Mailbox><t:EmailAddress>noreply@vergabeportal.de</t:EmailAddress></t:Mailbox></t:From>
<t:DateTimeReceived>2026-07-20T08:00:00Z</t:DateTimeReceived>
<t:Body BodyType="Text">Hallo Text Body</t:Body>
</t:Message>
</m:Items>
</m:GetItemResponseMessage>
</s:Body></s:Envelope>`;
}
/** Routes the stubbed httpntlm.post by the SOAPAction header, mirroring ewsPost's per-action dispatch. */
function mockEwsFlow(getItemResponseBody: string) {
httpntlmPost.mockImplementation((opts: any, cb: (err: Error | null, res: any) => void) => {
const action = opts.headers.SOAPAction as string;
if (action.includes('FindItem')) {
cb(null, soapResponse(200, FIND_ITEM_RESPONSE));
} else if (action.includes('GetItem')) {
cb(null, soapResponse(200, getItemResponseBody));
} else if (action.includes('UpdateItem')) {
cb(null, soapResponse(200, '<s:Envelope><s:Body/></s:Envelope>'));
} else {
cb(null, soapResponse(200, '<s:Envelope><s:Body/></s:Envelope>'));
}
});
}
beforeEach(() => {
httpntlmPost.mockClear();
});
describe('ExchangeInboxProvider.fetchMessages', () => {
it('a BodyType="HTML" GetItem response yields bodyHtml populated and bodyText empty', async () => {
mockEwsFlow(getItemHtmlResponse());
const provider = new ExchangeInboxProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toHaveLength(1);
expect(messages[0]).toMatchObject({
uid: 'AAA1',
messageId: 'msg-1@example.com',
subject: 'Neue Ausschreibung verfuegbar',
from: 'noreply@vergabeportal.de',
bodyHtml: 'Hallo HTML Body',
bodyText: '',
});
});
it('a BodyType="Text" GetItem response yields bodyText populated and bodyHtml null', async () => {
mockEwsFlow(getItemTextResponse());
const provider = new ExchangeInboxProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toHaveLength(1);
expect(messages[0]!.bodyHtml).toBeNull();
expect(messages[0]!.bodyText).toBe('Hallo Text Body');
});
it('marks the item read (IsRead) like fetchPdfAttachments', async () => {
mockEwsFlow(getItemHtmlResponse());
const provider = new ExchangeInboxProvider();
await provider.fetchMessages(BASE_CONFIG);
const updateCall = httpntlmPost.mock.calls.find(
(call) => (call[0].headers.SOAPAction as string).includes('UpdateItem'),
);
expect(updateCall).toBeDefined();
expect(updateCall![0].body).toContain('IsRead');
expect(updateCall![0].body).toContain('AAA1');
});
it('returns [] on non-200 FindItem response, without throwing', async () => {
httpntlmPost.mockImplementation((_opts: any, cb: (err: Error | null, res: any) => void) => {
cb(null, soapResponse(500, 'Internal Server Error'));
});
const provider = new ExchangeInboxProvider();
await expect(provider.fetchMessages(BASE_CONFIG)).resolves.toEqual([]);
});
it('returns [] when httpntlm.post errors, without throwing', async () => {
httpntlmPost.mockImplementation((_opts: any, cb: (err: Error | null, res: any) => void) => {
cb(new Error('NTLM handshake failed'), null);
});
const provider = new ExchangeInboxProvider();
await expect(provider.fetchMessages(BASE_CONFIG)).resolves.toEqual([]);
});
it('returns [] when FindItem has no items', async () => {
httpntlmPost.mockImplementation((opts: any, cb: (err: Error | null, res: any) => void) => {
const action = opts.headers.SOAPAction as string;
if (action.includes('FindItem')) {
cb(null, soapResponse(200, '<s:Envelope><s:Body><m:FindItemResponseMessage/></s:Body></s:Envelope>'));
} else {
cb(null, soapResponse(200, '<s:Envelope><s:Body/></s:Envelope>'));
}
});
const provider = new ExchangeInboxProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toEqual([]);
});
});
+105 -1
View File
@@ -1,7 +1,7 @@
import { Injectable, Logger } from '@nestjs/common'; import { Injectable, Logger } from '@nestjs/common';
// eslint-disable-next-line @typescript-eslint/no-require-imports // 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 }; 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 { InboxAttachment, InboxConfig, InboxEmail, InboxMessage } from './inbox-provider.interface';
import type { InboxProvider } from './inbox-provider.interface'; import type { InboxProvider } from './inbox-provider.interface';
const MAX_ATTACHMENT_BYTES = 25 * 1024 * 1024; // 25 MB — T-07-05 const MAX_ATTACHMENT_BYTES = 25 * 1024 * 1024; // 25 MB — T-07-05
@@ -80,6 +80,29 @@ function getItemSoap(itemIds: string[]): string {
</m:GetItem>`); </m:GetItem>`);
} }
/**
* GetItem SOAP requesting item:Body instead of item:Attachments — used by
* fetchMessages (D-02) to retrieve the HTML/text body without touching the
* attachment-metadata path getItemSoap() serves for fetchPdfAttachments.
*/
function getItemBodySoap(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:Body"/>
</t:AdditionalProperties>
</m:ItemShape>
<m:ItemIds>${ids}</m:ItemIds>
</m:GetItem>`);
}
function markReadSoap(itemId: string, changeKey: string): string { function markReadSoap(itemId: string, changeKey: string): string {
return soapEnvelope(` return soapEnvelope(`
<m:UpdateItem MessageDisposition="SaveOnly" ConflictResolution="AlwaysOverwrite"> <m:UpdateItem MessageDisposition="SaveOnly" ConflictResolution="AlwaysOverwrite">
@@ -248,6 +271,26 @@ export class ExchangeInboxProvider implements InboxProvider {
} }
} }
/**
* Connects via EWS, finds unread items (optionally filtered by sender —
* same restriction as fetchPdfAttachments), and returns each message's
* subject + HTML/text body.
*
* Additive method (D-02) — does NOT touch fetchViaEws (fetchPdfAttachments'
* path). Marks each processed item IsRead so re-polls don't reprocess it.
*
* @param config Decrypted inbox connection parameters
* @returns Messages with subject + body (empty array on error)
*/
async fetchMessages(config: InboxConfig): Promise<InboxMessage[]> {
try {
return await this.fetchMessagesViaEws(config);
} catch (error) {
this.logger.error(`EWS fetchMessages failed: ${(error as Error).message}`);
return [];
}
}
async testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }> { async testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }> {
try { try {
const folderElement = await this.resolveFolderElement(config); const folderElement = await this.resolveFolderElement(config);
@@ -411,6 +454,67 @@ export class ExchangeInboxProvider implements InboxProvider {
return results; return results;
} }
private async fetchMessagesViaEws(config: InboxConfig): Promise<InboxMessage[]> {
const folderElement = await this.resolveFolderElement(config);
// 1. FindItem — get IDs of unread emails (no attachment filter — unlike fetchViaEws)
const findSoap = findItemSoap(folderElement, 50, config.senderFilter);
const findRes = await this.ewsPost(config, findSoap, 'FindItem');
if (findRes.statusCode !== 200) {
this.logger.warn(`EWS FindItem (fetchMessages) returned HTTP ${findRes.statusCode}`);
return [];
}
const itemIds = extractAttrs(findRes.body, 't:ItemId', 'Id');
if (itemIds.length === 0) return [];
const results: InboxMessage[] = [];
// 2. GetItem in batches of 10, requesting item:Body instead of item:Attachments
for (let i = 0; i < itemIds.length; i += 10) {
const batch = itemIds.slice(i, i + 10);
const getRes = await this.ewsPost(config, getItemBodySoap(batch), 'GetItem');
if (getRes.statusCode !== 200) continue;
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, mirrors fetchViaEws)
if (config.senderFilter && from &&
!from.toLowerCase().includes(config.senderFilter.toLowerCase())) {
continue;
}
const bodyType = extractAttr(block, 't:Body', 'BodyType');
const bodyContent = extractAll(block, 't:Body')[0] ?? '';
const bodyHtml = bodyType === 'HTML' ? bodyContent : null;
const bodyText = bodyType === 'Text' ? bodyContent : '';
results.push({ uid, messageId, subject, from, date, bodyHtml, bodyText });
// Mark as read so subsequent polls skip this message (mirrors fetchViaEws).
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. */ /** Split a GetItem response body into per-message XML blocks. */
private splitMessageBlocks(xml: string): string[] { private splitMessageBlocks(xml: string): string[] {
const blocks: string[] = []; const blocks: string[] = [];
+176
View File
@@ -0,0 +1,176 @@
import { beforeEach, describe, expect, it, vi } from 'vitest';
import { ImapFlow } from 'imapflow';
import { ImapProvider } from './imap.provider';
import type { InboxConfig } from './inbox.types';
/**
* ImapProvider.fetchMessages spec (Plan 14-01, Task 2, TDD).
*
* Net-new coverage for the additive fetchMessages() method (D-02) — the seam
* the Plan 14-03 EmailAlertAdapter consumes. fetchPdfAttachments is NOT
* touched or retested here (regression coverage lives in the DKV suite).
*
* ImapFlow is mocked entirely — no real IMAP connection. Mirrors the mock
* style used in tender-mail.service.spec.ts (module-level vi.mock +
* mockImplementation returning a stub client).
*/
vi.mock('imapflow', () => ({
ImapFlow: vi.fn(),
}));
function makeReadable(text: string): NodeJS.ReadableStream {
const { Readable } = require('stream') as typeof import('stream');
return Readable.from([Buffer.from(text, 'utf8')]);
}
const BASE_CONFIG: InboxConfig = {
protocol: 'imap',
host: 'imap.example.com',
port: 993,
username: 'alerts@example.com',
password: 'secret',
encryption: 'ssl-tls',
folder: 'INBOX',
};
/** Body structure with a text/html part (id '1') and a text/plain part (id '2'). */
const MULTIPART_BODY_STRUCTURE = {
type: 'multipart/alternative',
childNodes: [
{ type: 'text/plain', part: '2' },
{ type: 'text/html', part: '1' },
],
};
function makeMockClient(overrides: Partial<Record<string, unknown>> = {}) {
return {
connect: vi.fn().mockResolvedValue(undefined),
logout: vi.fn().mockResolvedValue(undefined),
getMailboxLock: vi.fn().mockResolvedValue({ release: vi.fn() }),
search: vi.fn().mockResolvedValue([42]),
fetchAll: vi.fn().mockResolvedValue([
{
uid: 42,
envelope: {
messageId: '<msg-42@example.com>',
subject: 'Neue Ausschreibung verfügbar',
from: [{ address: 'noreply@vergabeportal.de' }],
date: new Date('2026-07-20T08:00:00Z'),
},
bodyStructure: MULTIPART_BODY_STRUCTURE,
},
]),
download: vi.fn(async (_uid: string, partId: string) => {
if (partId === '1') {
return { content: makeReadable('<p>Hallo <b>Welt</b></p>') };
}
return { content: makeReadable('Hallo Welt (Text)') };
}),
messageFlagsAdd: vi.fn().mockResolvedValue(undefined),
...overrides,
};
}
beforeEach(() => {
vi.clearAllMocks();
});
describe('ImapProvider.fetchMessages', () => {
it('returns one InboxMessage with bodyHtml from the html part and bodyText from the plain part', async () => {
const client = makeMockClient();
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toHaveLength(1);
expect(messages[0]).toMatchObject({
uid: 42,
messageId: '<msg-42@example.com>',
subject: 'Neue Ausschreibung verfügbar',
from: 'noreply@vergabeportal.de',
bodyHtml: '<p>Hallo <b>Welt</b></p>',
bodyText: 'Hallo Welt (Text)',
});
});
it('marks each processed message \\Seen (idempotency for re-polls)', async () => {
const client = makeMockClient();
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
await provider.fetchMessages(BASE_CONFIG);
expect(client.messageFlagsAdd).toHaveBeenCalledWith('42', ['\\Seen'], { uid: true });
});
it('honors the same UNSEEN + optional senderFilter search as fetchPdfAttachments', async () => {
const client = makeMockClient();
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
await provider.fetchMessages({ ...BASE_CONFIG, senderFilter: 'vergabeportal.de' });
expect(client.search).toHaveBeenCalledWith(
{ seen: false, from: 'vergabeportal.de' },
{ uid: true },
);
});
it('returns [] on connect error without throwing', async () => {
const client = makeMockClient({
connect: vi.fn().mockRejectedValue(new Error('ECONNREFUSED')),
});
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
await expect(provider.fetchMessages(BASE_CONFIG)).resolves.toEqual([]);
});
it('returns [] on search error without throwing', async () => {
const client = makeMockClient({
search: vi.fn().mockRejectedValue(new Error('search boom')),
});
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
await expect(provider.fetchMessages(BASE_CONFIG)).resolves.toEqual([]);
expect(client.logout).toHaveBeenCalled();
});
it('returns [] when there are no unread messages', async () => {
const client = makeMockClient({ search: vi.fn().mockResolvedValue([]) });
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toEqual([]);
});
it('sets bodyHtml null and bodyText empty when the message has neither part', async () => {
const client = makeMockClient({
fetchAll: vi.fn().mockResolvedValue([
{
uid: 7,
envelope: {
messageId: '<msg-7@example.com>',
subject: 'Nur Betreff',
from: [{ address: 'a@b.de' }],
date: new Date('2026-07-20T08:00:00Z'),
},
bodyStructure: { type: 'text/x-unknown' },
},
]),
});
(ImapFlow as unknown as ReturnType<typeof vi.fn>).mockImplementation(() => client);
const provider = new ImapProvider();
const messages = await provider.fetchMessages(BASE_CONFIG);
expect(messages).toHaveLength(1);
expect(messages[0]!.bodyHtml).toBeNull();
expect(messages[0]!.bodyText).toBe('');
});
});
+137 -1
View File
@@ -1,6 +1,6 @@
import { Injectable, Logger } from '@nestjs/common'; import { Injectable, Logger } from '@nestjs/common';
import { ImapFlow, MessageStructureObject } from 'imapflow'; import { ImapFlow, MessageStructureObject } from 'imapflow';
import type { InboxAttachment, InboxConfig, InboxEmail } from './inbox-provider.interface'; import type { InboxAttachment, InboxConfig, InboxEmail, InboxMessage } from './inbox-provider.interface';
import type { InboxProvider } from './inbox-provider.interface'; import type { InboxProvider } from './inbox-provider.interface';
/** /**
@@ -82,6 +82,43 @@ function collectPdfParts(
return parts; return parts;
} }
/** First-html + first-text body part IDs found in a message's MIME tree. */
interface BodyParts {
htmlPart?: string;
textPart?: string;
}
/**
* Recursively walks a message's MIME structure collecting the first
* `text/html` and first `text/plain` part IDs (D-02 — additive sibling to
* collectPdfParts, does not affect the PDF-attachment path).
*
* @param node Root or child MIME structure node
* @param found Accumulator (pass empty object on first call)
*/
function findBodyParts(
node: MessageStructureObject | undefined,
found: BodyParts = {},
): BodyParts {
if (!node) return found;
const type = node.type?.toLowerCase() ?? '';
if (type === 'text/html' && found.htmlPart === undefined && node.part !== undefined) {
found.htmlPart = node.part;
} else if (type === 'text/plain' && found.textPart === undefined && node.part !== undefined) {
found.textPart = node.part;
}
if (Array.isArray(node.childNodes)) {
for (const child of node.childNodes) {
findBodyParts(child, found);
}
}
return found;
}
/** /**
* IMAP inbox provider using the imapflow library. * IMAP inbox provider using the imapflow library.
* *
@@ -197,6 +234,105 @@ export class ImapProvider implements InboxProvider {
return results; return results;
} }
/**
* Connects to the IMAP server, searches for unread emails (optionally
* filtered by sender — same query as fetchPdfAttachments), and returns
* each message's subject + HTML/text body.
*
* Additive method (D-02) — does NOT touch fetchPdfAttachments. Marks each
* processed message \Seen so re-polls don't reprocess it (idempotency).
*
* @param config Decrypted inbox connection parameters
* @returns Messages with subject + body (empty array on error)
*/
async fetchMessages(config: InboxConfig): Promise<InboxMessage[]> {
const client = this.buildClient(config);
const results: InboxMessage[] = [];
try {
await client.connect();
} catch (err) {
this.logger.error(`IMAP fetchMessages connect failed: ${(err as Error).message}`);
return results;
}
let lock: Awaited<ReturnType<typeof client.getMailboxLock>> | null = null;
try {
lock = await client.getMailboxLock(config.folder || 'INBOX');
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;
}
// Same fetchAll-before-download ordering as fetchPdfAttachments (Pitfall 1)
const messages = await client.fetchAll(
uids.join(','),
{ envelope: true, bodyStructure: true },
{ uid: true },
);
for (const msg of messages) {
const { htmlPart, textPart } = findBodyParts(msg.bodyStructure);
let bodyHtml: string | null = null;
let bodyText = '';
if (htmlPart !== undefined) {
try {
const { content } = await client.download(String(msg.uid), htmlPart, { uid: true });
bodyHtml = (await streamToBuffer(content)).toString('utf8');
} catch (partErr) {
this.logger.warn(
`Skipped IMAP html body part for UID ${String(msg.uid)}: ${(partErr as Error).message}`,
);
}
}
if (textPart !== undefined) {
try {
const { content } = await client.download(String(msg.uid), textPart, { uid: true });
bodyText = (await streamToBuffer(content)).toString('utf8');
} catch (partErr) {
this.logger.warn(
`Skipped IMAP text body part for UID ${String(msg.uid)}: ${(partErr as Error).message}`,
);
}
}
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(),
bodyHtml,
bodyText,
});
// Mark as read so subsequent polls skip this message (idempotency).
try {
await client.messageFlagsAdd(String(msg.uid), ['\\Seen'], { uid: true });
} catch {
// Non-fatal: message will simply appear again on next poll
}
}
} catch (err) {
this.logger.error(`IMAP fetchMessages failed: ${(err as Error).message}`);
} finally {
lock?.release();
await client.logout();
}
return results;
}
/** /**
* Tests whether the IMAP connection can be established. * Tests whether the IMAP connection can be established.
* Connects and immediately logs out without selecting any folder. * Connects and immediately logs out without selecting any folder.
+15 -2
View File
@@ -1,7 +1,7 @@
import { InboxAttachment, InboxConfig, InboxEmail } from './inbox.types'; import { InboxAttachment, InboxConfig, InboxEmail, InboxMessage } from './inbox.types';
// Re-export types for downstream consumers that import from this module // Re-export types for downstream consumers that import from this module
export type { InboxAttachment, InboxConfig, InboxEmail }; export type { InboxAttachment, InboxConfig, InboxEmail, InboxMessage };
/** /**
* Abstract inbox provider contract — implemented by ImapProvider and * Abstract inbox provider contract — implemented by ImapProvider and
@@ -29,4 +29,17 @@ export interface InboxProvider {
* @returns { success: true } on success, { success: false, message } on failure * @returns { success: true } on success, { success: false, message } on failure
*/ */
testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }>; testConnection(config: InboxConfig): Promise<{ success: boolean; message?: string }>;
/**
* Connects to the configured inbox, searches for unread emails (optionally
* filtered by sender), and returns each message's subject + HTML/text body.
*
* Additive method (D-02, Phase 14 inbox extraction) — the seam the
* EmailAlertAdapter (Plan 14-03) consumes. Does NOT affect
* fetchPdfAttachments; DKV is not switched to this method.
*
* @param config Decrypted inbox connection parameters
* @returns Array of messages with subject + body (empty array on error)
*/
fetchMessages(config: InboxConfig): Promise<InboxMessage[]>;
} }