From c40a023321deb9b1a71c8798e63ee9c52bd8202d Mon Sep 17 00:00:00 2001 From: Schalli Date: Sat, 27 Jun 2026 00:24:47 +0200 Subject: [PATCH] =?UTF-8?q?feat(07-04):=20DkvService=20=E2=80=94=20pipelin?= =?UTF-8?q?e=20orchestration=20+=20vehicle/config/history=20logic?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Single-flight guard (processing flag) prevents concurrent inbox processing - processInbox: poll → 3-retry parse → driver-map → xlsx → 3-retry SMTP send → history - Parse failure (D-10): records Fehler history row with errorMessage - SMTP failure (D-16): exponential backoff 2s/4s, records Versand fehlgeschlagen - CONFIG_SAFE_SELECT excludes encryptedInboxCreds (T-07-12) - getExportFile rejects filenames with path separators or outside DKV_*.xlsx pattern (T-07-09) - CSV import: merge (upsert by tenantId+kennzeichen) and replace (deleteMany then createMany) modes - Invoice number extracted from email subject via regex; falls back to email-{uid} - Vehicle format string resolved via DkvExportService.resolveFahrzeug (D-19) --- apps/api/src/dkv/dkv.service.ts | 680 ++++++++++++++++++++++++++++++++ 1 file changed, 680 insertions(+) create mode 100644 apps/api/src/dkv/dkv.service.ts diff --git a/apps/api/src/dkv/dkv.service.ts b/apps/api/src/dkv/dkv.service.ts new file mode 100644 index 0000000..0af1aeb --- /dev/null +++ b/apps/api/src/dkv/dkv.service.ts @@ -0,0 +1,680 @@ +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 { imported: 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)); +}