feat(07-04): DkvService — pipeline orchestration + vehicle/config/history logic
- 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)
This commit is contained in:
@@ -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<string, unknown> = {
|
||||||
|
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<boolean> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
// 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<void> {
|
||||||
|
// 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<Buffer> {
|
||||||
|
// 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<void> {
|
||||||
|
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user