import { BadRequestException, Injectable, Logger, NotFoundException, } from '@nestjs/common'; import * as fs from 'fs'; import * as path from 'path'; import { CalendarCryptoService } from '../calendar/crypto.service'; import { PrismaService } from '../prisma/prisma.service'; import { DkvExportService } from './dkv-export.service'; import { DkvMailService } from './dkv-mail.service'; import { DkvParserService } from './dkv-parser.service'; import { DkvConfigDto } from './dto/dkv-config.dto'; import { CreateVehicleDto, UpdateVehicleDto } from './dto/dkv-vehicle.dto'; import { ExchangeInboxProvider } from './providers/exchange-inbox.provider'; import { ImapProvider } from './providers/imap.provider'; import type { DkvVehicleBlock, InboxConfig, InboxEmail } from './dkv.types'; /** * Prisma select for DkvModuleConfig — never includes encryptedInboxCreds. * T-07-12: Encrypted credential blob is excluded from all API responses. */ const CONFIG_SAFE_SELECT = { id: true, tenantId: true, protocol: true, host: true, port: true, encryption: true, folder: true, senderFilter: true, pollIntervalMin: true, isActive: true, exportRecipient: true, vehicleFormatString: true, // encryptedInboxCreds: NEVER included — T-07-12 createdAt: true, updatedAt: true, } as const; /** * DkvService — orchestrates the full DKV processing pipeline. * * Pipeline: * poll inbox → download PDF attachments → parse vehicle/transaction data * → map license plates to drivers → build xlsx → write to user-files/ * → send via SMTP → record history * * Security: * - T-05-13: Decrypted credentials never logged * - T-07-12: CONFIG_SAFE_SELECT excludes encryptedInboxCreds in all responses * - T-07-09: Export filename validated against safe pattern before reading (traversal guard) * - Single-flight guard: prevents concurrent inbox processing (Pitfall 7) * * Multi-tenant note (v1): The scheduler loads config via findFirst(). * Each processInbox(tenantId) call is per-tenant. The Controller scopes all * operations to req.tenantId. Full per-tenant scheduling (one cron per active * tenant) is deferred to a future plan — v1 covers single-tenant deployments. */ @Injectable() export class DkvService { private readonly logger = new Logger(DkvService.name); /** * Single-flight guard: if processing is already in progress, any concurrent * call to processInbox() returns early without starting a second pipeline * run (Pitfall 7 — prevents the prune race condition and duplicate records). */ private processing = false; /** Resolved path to user-files/ directory (monorepo root). */ private readonly userFilesDir: string; constructor( private readonly prisma: PrismaService, private readonly crypto: CalendarCryptoService, private readonly parser: DkvParserService, private readonly exporter: DkvExportService, private readonly mailer: DkvMailService, private readonly imapProvider: ImapProvider, private readonly exchangeProvider: ExchangeInboxProvider, ) { // __dirname at runtime = apps/api/dist/dkv/ — go up 4 levels to monorepo root this.userFilesDir = path.resolve(__dirname, '..', '..', '..', '..', 'user-files'); } // ─── Config ────────────────────────────────────────────────────────────────── /** * Load DKV module config for a tenant (safe — no encrypted creds). * When tenantId is omitted, returns the first row (used by scheduler on init). */ async loadConfig(tenantId?: string) { if (!tenantId) { return this.prisma.dkvModuleConfig.findFirst({ select: CONFIG_SAFE_SELECT }); } return this.prisma.dkvModuleConfig.findUnique({ where: { tenantId }, select: CONFIG_SAFE_SELECT, }); } /** * Upsert DKV module config for a tenant. * * Credential handling: * - If dto.password is non-empty: re-encrypt {username, password} together. * username is taken from dto.username if provided, else preserved from DB. * - If dto.username is non-empty but dto.password is empty: preserve existing * password; re-encrypt with new username. * - If both are empty/undefined: preserve existing encryptedInboxCreds entirely. * * T-07-12: Returns safe select (no encryptedInboxCreds). * T-05-13: Never logs decrypted credentials. */ async saveConfig(tenantId: string, dto: DkvConfigDto) { let encryptedInboxCreds: string | undefined; const credChanged = (dto.password && dto.password.length > 0) || (dto.username !== undefined && dto.username !== null); if (credChanged) { let username: string = dto.username ?? ''; let password: string = dto.password ?? ''; // Preserve the field that was left empty from the existing stored value if (!dto.password || !dto.username) { try { const existing = await this.prisma.dkvModuleConfig.findUnique({ where: { tenantId } }); if (existing?.encryptedInboxCreds) { // T-05-13: decrypt only to preserve — never log the result const stored = JSON.parse(this.crypto.decrypt(existing.encryptedInboxCreds)) as { username?: string; password?: string; }; if (!dto.username) username = stored.username ?? ''; if (!dto.password) password = stored.password ?? ''; } } catch { // Ignore decrypt errors — will overwrite with whatever was provided } } encryptedInboxCreds = this.crypto.encrypt(JSON.stringify({ username, password })); } const data: Record = { protocol: dto.protocol, encryption: dto.encryption, ...(dto.host !== undefined && { host: dto.host }), ...(dto.port !== undefined && { port: dto.port }), ...(dto.folder !== undefined && { folder: dto.folder }), ...(dto.senderFilter !== undefined && { senderFilter: dto.senderFilter }), ...(dto.pollIntervalMin !== undefined && { pollIntervalMin: dto.pollIntervalMin }), ...(dto.isActive !== undefined && { isActive: dto.isActive }), ...(dto.exportRecipient !== undefined && { exportRecipient: dto.exportRecipient }), ...(dto.vehicleFormatString !== undefined && { vehicleFormatString: dto.vehicleFormatString }), ...(encryptedInboxCreds !== undefined && { encryptedInboxCreds }), }; return this.prisma.dkvModuleConfig.upsert({ where: { tenantId }, create: { tenantId, ...data }, update: data, select: CONFIG_SAFE_SELECT, }); } /** * Test the inbox connection using credentials from the form DTO. * When dto.password is empty, falls back to the stored encrypted password. * * T-05-13: Decrypted password used only within this method scope — never logged. */ async testConnection(tenantId: string, dto: DkvConfigDto): Promise { let password: string | undefined = dto.password; // If no password in DTO, fall back to the stored one if (!password) { try { const existing = await this.prisma.dkvModuleConfig.findUnique({ where: { tenantId } }); if (existing?.encryptedInboxCreds) { const stored = JSON.parse(this.crypto.decrypt(existing.encryptedInboxCreds)) as { password?: string; }; password = stored.password; } } catch { // T-05-13: generic — no credential details this.logger.error(`DKV testConnection: failed to load stored credentials for tenant ${tenantId}`); } } const inboxConfig: InboxConfig = { protocol: dto.protocol, host: dto.host ?? '', port: dto.port ?? 993, username: dto.username, password, encryption: dto.encryption, folder: dto.folder ?? 'INBOX', senderFilter: dto.senderFilter, }; const provider = dto.protocol === 'exchange' ? this.exchangeProvider : this.imapProvider; return provider.testConnection(inboxConfig); } // ─── Pipeline ───────────────────────────────────────────────────────────────── /** * Trigger an inbox check for a tenant. * Returns a summary object: status, message, checkedAt. * * Used by POST /dkv/check-now and indirectly by the scheduler. */ async checkNow(tenantId: string): Promise<{ status: string; message: string; checkedAt: string; }> { const checkedAt = new Date().toISOString(); try { await this.processInbox(tenantId); return { status: 'ok', message: 'Posteingang geprüft', checkedAt }; } catch (err) { return { status: 'error', message: (err as Error).message, checkedAt, }; } } /** * Run the full DKV inbox processing pipeline for a tenant. * * Single-flight guard: returns early if already processing (Pitfall 7). * Pipeline: load config → decrypt creds → fetch PDFs → parse → map drivers * → build xlsx → write to user-files/ → send SMTP → record history. */ async processInbox(tenantId: string): Promise { if (this.processing) { this.logger.warn(`DKV inbox already processing for tenant ${tenantId} — skipping`); return; } this.processing = true; try { await this._runPipeline(tenantId); } finally { this.processing = false; } } private async _runPipeline(tenantId: string): Promise { // Load raw config (need encryptedInboxCreds for decryption) const config = await this.prisma.dkvModuleConfig.findUnique({ where: { tenantId } }); if (!config) { this.logger.warn(`DKV processInbox: no config for tenant ${tenantId}`); return; } if (!config.isActive) { this.logger.log(`DKV inbox polling is inactive for tenant ${tenantId}`); return; } // Decrypt inbox credentials — T-05-13: never log result let inboxConfig: InboxConfig; try { const creds = config.encryptedInboxCreds ? (JSON.parse(this.crypto.decrypt(config.encryptedInboxCreds)) as { username?: string; password?: string; }) : {}; inboxConfig = { protocol: config.protocol, host: config.host ?? '', port: config.port ?? 993, username: creds.username, password: creds.password, encryption: config.encryption, folder: config.folder, senderFilter: config.senderFilter ?? undefined, }; } catch { this.logger.error(`DKV inbox config decryption failed for tenant ${tenantId}`); return; } // Select provider by protocol const provider = config.protocol === 'exchange' ? this.exchangeProvider : this.imapProvider; // Fetch PDF attachments from inbox let emails: InboxEmail[]; try { emails = await provider.fetchPdfAttachments(inboxConfig); } catch (err) { this.logger.error( `DKV inbox fetch failed for tenant ${tenantId}: ${(err as Error).message}`, ); return; } if (emails.length === 0) { this.logger.log(`DKV inbox check: no matching emails for tenant ${tenantId}`); return; } const vehicleFormatString = config.vehicleFormatString ?? '{Marke}/{Modell}/{Kennzeichen}'; for (const email of emails) { for (const attachment of email.attachments) { await this._processAttachment( tenantId, attachment.buffer, email, config.exportRecipient ?? undefined, vehicleFormatString, ); } } } /** * Process a single PDF attachment through the full pipeline. * D-10: Up to 3 parse retries; on final failure record Fehler history row. * D-16: Up to 3 SMTP send retries with exponential backoff; on final failure * record Versand fehlgeschlagen (file stays locally available). */ private async _processAttachment( tenantId: string, pdfBuffer: Buffer, email: InboxEmail, recipient: string | undefined, vehicleFormatString: string, ): Promise { // D-10: Up to 3 parse retries let vehicles: DkvVehicleBlock[] | null = null; let parseError: string | null = null; for (let attempt = 1; attempt <= 3; attempt++) { try { vehicles = await this.parser.parsePdf(pdfBuffer); parseError = null; break; } catch (err) { parseError = (err as Error).message; this.logger.warn( `DKV PDF parse attempt ${attempt}/3 for tenant ${tenantId}: ${parseError}`, ); } } if (!vehicles) { // Record parse failure in history (D-10) await this.prisma.dkvInvoiceHistory.create({ data: { tenantId, rechnungsnummer: this._extractInvoiceNumber(email.subject, email.uid), anzahlFahrzeuge: 0, anzahlTransaktionen: 0, status: 'Fehler', errorMessage: parseError ?? 'PDF parse failed', }, }); return; } // Extract invoice metadata for filename + history record const rechnungsnummer = this._extractInvoiceNumber(email.subject, email.uid); const invoiceMonth = this._extractInvoiceMonth(vehicles); const anzahlTransaktionen = vehicles.reduce((s, v) => s + v.transactions.length, 0); // Build export rows — look up drivers from vehicle master const exportRows = await this._buildExportRows(tenantId, vehicles, vehicleFormatString); // Generate xlsx buffer and write to user-files/ (DkvExportService) const xlsxBuffer = this.exporter.buildExcelBuffer(exportRows); const exportFilename = this.exporter.writeAndPrune(xlsxBuffer, rechnungsnummer, invoiceMonth); // D-16: SMTP send with 3-retry exponential backoff let smtpStatus: 'Verarbeitet' | 'Versand fehlgeschlagen' = 'Verarbeitet'; let smtpError: string | undefined; if (recipient) { for (let attempt = 1; attempt <= 3; attempt++) { try { await this.mailer.sendExportEmail(tenantId, recipient, xlsxBuffer, exportFilename); smtpStatus = 'Verarbeitet'; smtpError = undefined; break; } catch (err) { smtpError = (err as Error).message; this.logger.warn(`DKV SMTP send attempt ${attempt}/3 failed: ${smtpError}`); if (attempt < 3) { await _delay(Math.pow(2, attempt) * 1000); // 2s, 4s } else { smtpStatus = 'Versand fehlgeschlagen'; // D-16: file stays locally available for manual download } } } } else { // No recipient configured — mark as processed without sending this.logger.warn(`DKV: no exportRecipient configured for tenant ${tenantId} — skipping SMTP send`); } // Record history row (D-20) await this.prisma.dkvInvoiceHistory.create({ data: { tenantId, rechnungsnummer, anzahlFahrzeuge: vehicles.length, anzahlTransaktionen, status: smtpStatus, ...(smtpError && { errorMessage: smtpError }), exportFilename, }, }); this.logger.log( `DKV processed: ${rechnungsnummer} — ${vehicles.length} vehicles, ${anzahlTransaktionen} transactions, status: ${smtpStatus}`, ); } // ─── Vehicle CRUD ──────────────────────────────────────────────────────────── async listVehicles(tenantId: string) { return this.prisma.dkvVehicleMaster.findMany({ where: { tenantId }, orderBy: { kennzeichen: 'asc' }, }); } async createVehicle(tenantId: string, dto: CreateVehicleDto) { return this.prisma.dkvVehicleMaster.create({ data: { tenantId, ...dto }, }); } async updateVehicle(tenantId: string, id: string, dto: UpdateVehicleDto) { const existing = await this.prisma.dkvVehicleMaster.findFirst({ where: { id, tenantId }, }); if (!existing) throw new NotFoundException('Vehicle not found'); return this.prisma.dkvVehicleMaster.update({ where: { id }, data: dto }); } async deleteVehicle(tenantId: string, id: string): Promise<{ deleted: boolean }> { const existing = await this.prisma.dkvVehicleMaster.findFirst({ where: { id, tenantId }, }); if (!existing) throw new NotFoundException('Vehicle not found'); await this.prisma.dkvVehicleMaster.delete({ where: { id } }); return { deleted: true }; } /** * Bulk-import vehicles from semicolon-delimited CSV text. * * CSV header (case-insensitive): Kennzeichen;Marke;Modell;Fahrer * * mode='merge': upsert each vehicle (add new, update existing by Kennzeichen). * mode='replace': delete all existing vehicles for this tenant, then insert. * * Research pattern: CSV Vehicle Import Pattern (RESEARCH.md Code Examples). */ async importVehiclesCsv( tenantId: string, csvText: string, mode: 'merge' | 'replace', ): Promise<{ imported: number; mode: string }> { const vehicles = _parseVehicleCsv(csvText); if (vehicles.length === 0) { throw new BadRequestException( 'CSV enthält keine gültigen Fahrzeuge. Erwarteter Header: Kennzeichen;Marke;Modell;Fahrer', ); } if (mode === 'replace') { await this.prisma.dkvVehicleMaster.deleteMany({ where: { tenantId } }); await this.prisma.dkvVehicleMaster.createMany({ data: vehicles.map((v) => ({ tenantId, ...v })), }); } else { // Merge: upsert by (tenantId, kennzeichen) compound unique key for (const v of vehicles) { await this.prisma.dkvVehicleMaster.upsert({ where: { tenantId_kennzeichen: { tenantId, kennzeichen: v.kennzeichen } }, create: { tenantId, ...v }, update: { marke: v.marke, modell: v.modell, fahrer: v.fahrer }, }); } } this.logger.log(`DKV CSV import: ${vehicles.length} vehicles, mode=${mode}, tenant=${tenantId}`); return { count: vehicles.length, mode }; } // ─── History ───────────────────────────────────────────────────────────────── /** * Get paginated processing history for a tenant. * Ordered by datumZeit descending (most recent first). * T-07-06: pagination prevents unbounded result-set DoS. */ async getHistory( tenantId: string, page = 1, limit = 20, ): Promise<{ items: unknown[]; total: number; page: number; limit: number; }> { const skip = (page - 1) * limit; const [items, total] = await Promise.all([ this.prisma.dkvInvoiceHistory.findMany({ where: { tenantId }, orderBy: { datumZeit: 'desc' }, skip, take: limit, }), this.prisma.dkvInvoiceHistory.count({ where: { tenantId } }), ]); return { items, total, page, limit }; } // ─── Export file download ──────────────────────────────────────────────────── /** * Read a DKV export file from user-files/ and return its Buffer. * * T-07-09 traversal guard: validates filename against the server-generated * pattern `DKV_*.xlsx` before reading. Rejects any filename containing path * separators, `..`, or characters outside the expected character set. * * @throws BadRequestException when filename fails validation * @throws NotFoundException when the file does not exist */ async getExportFile(tenantId: string, filename: string): Promise { // Traversal guard: whitelist-validate the filename before reading if ( filename.includes('/') || filename.includes('\\') || filename.includes('..') || !/^DKV_[\w\-]+\.xlsx$/.test(filename) ) { throw new BadRequestException('Invalid export filename'); } const filePath = path.join(this.userFilesDir, filename); if (!fs.existsSync(filePath)) { throw new NotFoundException(`Export file not found: ${filename}`); } return fs.readFileSync(filePath); } // ─── Private helpers ────────────────────────────────────────────────────────── /** * Build ExportRow[] for the xlsx generator by looking up each vehicle's driver. * * Vehicles with no matching DkvVehicleMaster entry still appear in the export * with an empty Fahrer field (D-13). The Fahrzeug column is resolved via the * format string with empty Marke/Modell/Fahrer placeholders for unknown plates. */ private async _buildExportRows( tenantId: string, vehicles: DkvVehicleBlock[], vehicleFormatString: string, ): Promise<{ lieferdatum: string; fahrzeug: string; fahrer: string; ort: string; kilometerstand: number }[]> { // Batch load vehicle master to avoid N+1 queries const masters = await this.prisma.dkvVehicleMaster.findMany({ where: { tenantId } }); const masterMap = new Map(masters.map((m) => [m.kennzeichen, m])); const rows: { lieferdatum: string; fahrzeug: string; fahrer: string; ort: string; kilometerstand: number }[] = []; for (const vehicle of vehicles) { const master = masterMap.get(vehicle.kennzeichen); // For unknown plates: use empty string values for unknown fields const vehicleForFormat = { kennzeichen: vehicle.kennzeichen, marke: master?.marke ?? '', modell: master?.modell ?? '', fahrer: master?.fahrer ?? '', }; const fahrzeug = this.exporter.resolveFahrzeug(vehicleForFormat, vehicleFormatString); const fahrer = master?.fahrer ?? ''; for (const tx of vehicle.transactions) { rows.push({ lieferdatum: tx.lieferdatum, fahrzeug, fahrer, ort: tx.ort, kilometerstand: tx.kilometerstand, }); } } return rows; } /** * Extract a DKV invoice number from the email subject line. * Pattern: DD-DDDDDDDDD-DDD (e.g. "26-651566449-001") * Falls back to "email-{uid}" when no match is found. */ private _extractInvoiceNumber(subject: string, uid: number | string): string { const match = subject?.match(/(\d{2}-\d{9}-\d{3})/); if (match?.[1]) return match[1]; // Fallback when invoice number cannot be parsed from subject return `email-${String(uid)}`; } /** * Derive the invoice month (YYYY-MM) from the first transaction date. * DKV dates are "DD.MM.YYYY" — this converts to "YYYY-MM" for the filename. * Falls back to the current month when no transaction date is available. */ private _extractInvoiceMonth(vehicles: DkvVehicleBlock[]): string { const firstDate = vehicles[0]?.transactions[0]?.lieferdatum; if (firstDate) { const parts = firstDate.split('.'); if (parts.length === 3) return `${parts[2]}-${parts[1]}`; } return new Date().toISOString().slice(0, 7); } } // ─── Module-level helpers ──────────────────────────────────────────────────── /** * Parse semicolon-delimited CSV into vehicle master records. * Expected header: Kennzeichen;Marke;Modell;Fahrer (case-insensitive). * Skips rows with empty Kennzeichen. * * Research: CSV Vehicle Import Pattern (RESEARCH.md Code Examples). */ function _parseVehicleCsv( csvText: string, ): { kennzeichen: string; marke: string; modell: string; fahrer: string }[] { const lines = csvText .trim() .split('\n') .map((l) => l.trim()) .filter(Boolean); if (lines.length < 2) return []; const headers = lines[0].split(';').map((h) => h.trim().toLowerCase()); return lines .slice(1) .map((line) => { const cols = line.split(';').map((c) => c.trim()); return { kennzeichen: cols[headers.indexOf('kennzeichen')] ?? '', marke: cols[headers.indexOf('marke')] ?? '', modell: cols[headers.indexOf('modell')] ?? '', fahrer: cols[headers.indexOf('fahrer')] ?? '', }; }) .filter((v) => v.kennzeichen.length > 0); } /** Simple promise-based delay for retry backoff. */ function _delay(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); }