import { Injectable, Logger, Optional } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { parse } from 'csv-parse/sync'; import * as fs from 'fs/promises'; import * as path from 'path'; import { CsvRateLoaderPort } from '@domain/ports/out/csv-rate-loader.port'; import { CsvRate, FreightPricing, FobCharges, DgSurchargeInfo, DgSurchargeValue, HandlingUnit, FrequencyType, RateDirection, } from '@domain/entities/csv-rate.entity'; import { PortCode } from '@domain/value-objects/port-code.vo'; import { ContainerType } from '@domain/value-objects/container-type.vo'; import { DateRange } from '@domain/value-objects/date-range.vo'; import { S3StorageAdapter } from '@infrastructure/storage/s3-storage.adapter'; import { TypeOrmCsvRateConfigRepository } from '@infrastructure/persistence/typeorm/repositories/typeorm-csv-rate-config.repository'; /** * Standardized 33-column CSV row. * All suppliers share this exact schema. */ interface CsvRow { // Supplier identity companyName: string; companyEmail: string; // Route geography originCFS: string; originCode: string; portOfLoading: string; routing: string; destinationCFS: string; destinationCode: string; destinationCountry: string; // Container containerType: string; // Freight freightCurrency: string; freightRatePerCBM: string; freightMinimum: string; // FOB charges fobCurrency: string; fobDocumentation: string; fobISPS: string; fobHandling: string; fobHandlingUnit: string; fobHandlingMinimum: string; fobSolas: string; fobCustoms: string; fobAMS_ACI: string; fobISF5: string; fobDGAdmin: string; // DG surcharge dgSurchargeCurrency: string; dgSurchargeRate: string; dgSurchargeUnit: string; dgSurchargeMin: string; // Metadata remarks: string; frequency: string; transitDays: string; validFrom: string; validUntil: string; } const REQUIRED_COLUMNS = [ 'companyName', 'companyEmail', 'originCFS', 'originCode', 'portOfLoading', 'routing', 'destinationCFS', 'destinationCode', 'destinationCountry', 'containerType', 'freightCurrency', 'freightRatePerCBM', 'freightMinimum', 'fobCurrency', 'fobDocumentation', 'fobISPS', 'fobHandling', 'fobHandlingUnit', 'fobHandlingMinimum', 'fobSolas', 'fobCustoms', 'fobAMS_ACI', 'fobISF5', 'fobDGAdmin', 'dgSurchargeCurrency', 'dgSurchargeRate', 'dgSurchargeUnit', 'dgSurchargeMin', 'remarks', 'frequency', 'transitDays', 'validFrom', 'validUntil', ]; @Injectable() export class CsvRateLoaderAdapter implements CsvRateLoaderPort { private readonly logger = new Logger(CsvRateLoaderAdapter.name); private readonly csvDirectory: string; // Legacy fallback mapping (export grids). Files are suffixed by direction: // -export.csv / -import.csv. private readonly companyFileMapping: Map = new Map([ ['SSC Consolidation', 'ssc-consolidation-export.csv'], ['ECU Worldwide', 'ecu-worldwide-export.csv'], ['TCC Logistics', 'tcc-logistics-export.csv'], ['NVO Consolidation', 'nvo-consolidation-export.csv'], ]); constructor( @Optional() private readonly s3Storage?: S3StorageAdapter, @Optional() private readonly configService?: ConfigService, @Optional() private readonly csvConfigRepository?: TypeOrmCsvRateConfigRepository ) { this.csvDirectory = path.join( process.cwd(), 'src', 'infrastructure', 'storage', 'csv-storage', 'rates' ); this.logger.log(`CSV directory initialized: ${this.csvDirectory}`); } async loadRatesFromCsv( filePath: string, companyEmail: string, companyNameOverride?: string, direction?: RateDirection ): Promise { this.logger.log(`Loading rates from CSV: ${filePath}`); try { let fileContent: string; let resolvedDirection: RateDirection | undefined = direction; if (this.s3Storage && this.configService && this.csvConfigRepository && companyNameOverride) { try { const config = await this.csvConfigRepository.findByCompanyName( companyNameOverride, direction ); resolvedDirection = resolvedDirection ?? config?.direction ?? undefined; const minioObjectKey = config?.metadata?.minioObjectKey as string | undefined; if (minioObjectKey) { const bucket = this.configService.get('AWS_S3_BUCKET', 'xpeditis-csv-rates'); const buffer = await this.s3Storage.download({ bucket, key: minioObjectKey }); fileContent = buffer.toString('utf-8'); } else { throw new Error('No MinIO object key'); } } catch (minioError: any) { this.logger.warn(`MinIO unavailable: ${minioError.message}. Using local file.`); const fullPath = this.resolvePath(filePath); fileContent = await fs.readFile(fullPath, 'utf-8'); } } else { const fullPath = this.resolvePath(filePath); fileContent = await fs.readFile(fullPath, 'utf-8'); } // Last-resort direction: filename suffix (-import.csv), default EXPORT resolvedDirection = resolvedDirection ?? (/-import\.csv$/i.test(filePath) ? 'IMPORT' : 'EXPORT'); const records: CsvRow[] = parse(fileContent, { columns: true, skip_empty_lines: true, trim: true, }); this.logger.log(`Parsed ${records.length} rows from ${filePath}`); this.validateCsvStructure(records); const rates = records.map((record, index) => { try { return this.mapToCsvRate(record, companyEmail, companyNameOverride, resolvedDirection); } catch (error) { const msg = error instanceof Error ? error.message : String(error); throw new Error(`Row ${index + 1} in ${filePath}: ${msg}`); } }); this.logger.log(`Loaded ${rates.length} rates from ${filePath}`); return rates; } catch (error) { const msg = error instanceof Error ? error.message : String(error); this.logger.error(`Failed to load ${filePath}: ${msg}`); throw new Error(`CSV loading failed for ${filePath}: ${msg}`); } } async loadRatesByCompany(companyName: string): Promise { const fileName = this.companyFileMapping.get(companyName); if (!fileName) { this.logger.warn(`No CSV file for company: ${companyName}`); return []; } const email = `info@${companyName.toLowerCase().replace(/\s+/g, '-')}.com`; return this.loadRatesFromCsv(fileName, email); } async validateCsvFile( filePath: string ): Promise<{ valid: boolean; errors: string[]; rowCount?: number }> { const errors: string[] = []; try { const fullPath = this.resolvePath(filePath); try { await fs.access(fullPath); } catch { return { valid: false, errors: [`File not found: ${filePath}`] }; } const fileContent = await fs.readFile(fullPath, 'utf-8'); const records: CsvRow[] = parse(fileContent, { columns: true, skip_empty_lines: true, trim: true, }); if (records.length === 0) { return { valid: false, errors: ['CSV file is empty'], rowCount: 0 }; } try { this.validateCsvStructure(records); } catch (e) { errors.push(e instanceof Error ? e.message : String(e)); } records.forEach((record, index) => { try { this.mapToCsvRate(record, 'validation@example.com'); } catch (e) { errors.push(`Row ${index + 1}: ${e instanceof Error ? e.message : String(e)}`); } }); return { valid: errors.length === 0, errors, rowCount: records.length }; } catch (e) { return { valid: false, errors: [`Validation failed: ${e instanceof Error ? e.message : String(e)}`], }; } } async getAvailableCsvFiles(): Promise { try { if (this.s3Storage && this.csvConfigRepository) { try { const configs = await this.csvConfigRepository.findAll(); const minioFiles = configs .filter(c => c.metadata?.minioObjectKey) .map(c => c.metadata?.minioObjectKey as string); if (minioFiles.length > 0) return minioFiles; } catch { // fall through to local } } try { await fs.access(this.csvDirectory); } catch { return []; } const files = await fs.readdir(this.csvDirectory); return files.filter(f => f.endsWith('.csv')); } catch { return []; } } private resolvePath(filePath: string): string { return path.isAbsolute(filePath) ? filePath : path.join(this.csvDirectory, filePath); } private validateCsvStructure(records: CsvRow[]): void { if (records.length === 0) throw new Error('CSV file is empty'); const firstRecord = records[0]; const missing = REQUIRED_COLUMNS.filter(col => !(col in firstRecord)); if (missing.length > 0) { throw new Error(`Missing required columns: ${missing.join(', ')}`); } } private mapToCsvRate( r: CsvRow, companyEmail: string, companyNameOverride?: string, direction: RateDirection = 'EXPORT' ): CsvRate { const companyName = companyNameOverride || r.companyName.trim(); // Admin-configured email always takes priority over the value in the CSV row const email = companyEmail?.trim() || r.companyEmail?.trim(); const freight: FreightPricing = { freightCurrency: r.freightCurrency.toUpperCase(), freightRatePerCBM: parseFloat(r.freightRatePerCBM) || 0, freightMinimum: parseFloat(r.freightMinimum) || 0, }; const fob: FobCharges = { fobCurrency: r.fobCurrency.toUpperCase(), fobDocumentation: parseInt(r.fobDocumentation, 10) || 0, fobISPS: parseInt(r.fobISPS, 10) || 0, fobHandling: parseInt(r.fobHandling, 10) || 0, fobHandlingUnit: (r.fobHandlingUnit?.toUpperCase() === 'W' ? 'W' : 'UP') as HandlingUnit, fobHandlingMinimum: parseInt(r.fobHandlingMinimum, 10) || 0, fobSolas: parseInt(r.fobSolas, 10) || 0, fobCustoms: parseInt(r.fobCustoms, 10) || 0, fobAMS_ACI: parseFloat(r.fobAMS_ACI) || 0, fobISF5: parseFloat(r.fobISF5) || 0, fobDGAdmin: parseInt(r.fobDGAdmin, 10) || 0, }; const dgSurcharge: DgSurchargeInfo = { dgSurchargeCurrency: (r.dgSurchargeCurrency || r.fobCurrency).toUpperCase(), dgSurchargeRate: parseDgValue(r.dgSurchargeRate), dgSurchargeUnit: (['UP', 'LS', '%'].includes(r.dgSurchargeUnit?.toUpperCase()) ? r.dgSurchargeUnit.toUpperCase() : 'LS') as 'UP' | 'LS' | '%', dgSurchargeMin: parseDgValue(r.dgSurchargeMin), }; const validFrom = new Date(r.validFrom); const validUntil = new Date(r.validUntil); const validity = DateRange.create(validFrom, validUntil, true); const frequency = parseFrequency(r.frequency); return new CsvRate( companyName, email, r.originCFS.trim(), PortCode.create(r.originCode.trim()), r.portOfLoading.trim(), r.routing.trim(), r.destinationCFS.trim(), PortCode.create(r.destinationCode.trim()), r.destinationCountry.trim(), ContainerType.create(r.containerType.trim()), freight, fob, dgSurcharge, r.remarks?.trim() || '', frequency, parseInt(r.transitDays, 10), validity, direction ); } } function parseDgValue(raw: string): DgSurchargeValue { if (!raw || raw.trim() === '') return 0; const upper = raw.trim().toUpperCase(); if (upper === 'ON REQUEST') return 'ON REQUEST'; if (upper === 'NOT ACCEPTED') return 'NOT ACCEPTED'; const num = parseFloat(raw); return isNaN(num) ? 0 : num; } function parseFrequency(raw: string): FrequencyType { switch (raw?.trim()) { case 'Weekly': return 'Weekly'; case 'Bi-Weekly': return 'Bi-Weekly'; case 'Bi-Monthly': return 'Bi-Monthly'; case 'Monthly': return 'Monthly'; default: return 'Weekly'; } }