feat(07-04): DkvSchedulerService + DkvController — dynamic cron + REST surface
- DkvSchedulerService: SchedulerRegistry.addCronJob (dynamic interval, not static @Cron)
- onModuleInit loads first active config (v1 single-tenant, documented in SUMMARY)
- setInterval() replaces existing job and registers new one with */ cron expression
- stopJob() removes job when config.isActive=false
- cron package resolved via require() workaround (pnpm strict isolation: transitive dep)
- DkvController: 12 handlers all carrying @Roles(Role.ADMIN, Role.SUPER_ADMIN)
- Routes: GET/PUT config, POST check-now, POST test-connection, GET history,
GET exports/:filename, GET/POST/PUT/DELETE vehicles, POST vehicles/import
- vehicles/import uses FileInterceptor('file') for CSV multipart upload
- exports/:filename streams file as attachment; traversal guard in DkvService
- Controller coordinates scheduler after PUT /dkv/config (no circular dep)
This commit is contained in:
@@ -0,0 +1,138 @@
|
|||||||
|
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
|
||||||
|
import { SchedulerRegistry } from '@nestjs/schedule';
|
||||||
|
import { DkvService } from './dkv.service';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CronJob constructor — resolved at runtime via require() because `cron` is a
|
||||||
|
* transitive dependency of @nestjs/schedule (not a direct api dep under pnpm
|
||||||
|
* strict isolation, so `import { CronJob } from 'cron'` fails type-check).
|
||||||
|
* At runtime, cron IS on disk as @nestjs/schedule@6 declares it as a peer dep.
|
||||||
|
*/
|
||||||
|
// eslint-disable-next-line @typescript-eslint/no-require-imports
|
||||||
|
const CronJobClass: new (cronTime: string, onTick: () => void) => { start(): void } =
|
||||||
|
// eslint-disable-next-line @typescript-eslint/no-unsafe-member-access
|
||||||
|
require('cron').CronJob as new (cronTime: string, onTick: () => void) => { start(): void };
|
||||||
|
|
||||||
|
/**
|
||||||
|
* DkvSchedulerService — dynamic cron job lifecycle management for inbox polling.
|
||||||
|
*
|
||||||
|
* Uses `SchedulerRegistry.addCronJob()` instead of the static `@Cron()` decorator
|
||||||
|
* so the polling interval can be updated at runtime when the admin changes the
|
||||||
|
* module config. (Research Pattern 7: Dynamic Cron Job; Pitfall 4: ScheduleModule
|
||||||
|
* must be registered in AppModule — done in Plan 01.)
|
||||||
|
*
|
||||||
|
* Multi-tenant note (v1): On init, the scheduler loads the first active
|
||||||
|
* DkvModuleConfig row via findFirst(). For single-tenant deployments this
|
||||||
|
* is always the correct config. Multi-tenant scheduling (one cron job per
|
||||||
|
* active tenant) is deferred to a future plan.
|
||||||
|
*
|
||||||
|
* The DkvController calls `setInterval()` after saving config so the cron job
|
||||||
|
* reflects any admin change immediately — without a service restart.
|
||||||
|
*/
|
||||||
|
@Injectable()
|
||||||
|
export class DkvSchedulerService implements OnModuleInit {
|
||||||
|
private readonly logger = new Logger(DkvSchedulerService.name);
|
||||||
|
|
||||||
|
/** Name of the managed cron job in the SchedulerRegistry. */
|
||||||
|
private readonly JOB_NAME = 'dkv-inbox-poll';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The tenantId this scheduler is currently serving.
|
||||||
|
* Updated when setInterval() is called with a new tenantId.
|
||||||
|
*/
|
||||||
|
private activeTenantId: string | null = null;
|
||||||
|
|
||||||
|
constructor(
|
||||||
|
private readonly schedulerRegistry: SchedulerRegistry,
|
||||||
|
private readonly dkvService: DkvService,
|
||||||
|
) {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* On application startup: load the first active DkvModuleConfig and
|
||||||
|
* register the cron job if the module is active.
|
||||||
|
*
|
||||||
|
* Errors are caught and logged (not re-thrown) so a missing or broken
|
||||||
|
* config does not prevent the rest of the application from starting.
|
||||||
|
*/
|
||||||
|
async onModuleInit(): Promise<void> {
|
||||||
|
try {
|
||||||
|
// loadConfig without tenantId → findFirst (v1 single-tenant)
|
||||||
|
const config = await this.dkvService.loadConfig();
|
||||||
|
if (config?.isActive && config.tenantId) {
|
||||||
|
this.activeTenantId = config.tenantId;
|
||||||
|
this.setInterval(config.pollIntervalMin, config.tenantId);
|
||||||
|
this.logger.log(
|
||||||
|
`DKV scheduler initialized: every ${config.pollIntervalMin} min for tenant ${config.tenantId}`,
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
this.logger.log('DKV scheduler: no active config found — cron job not registered');
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
this.logger.error(
|
||||||
|
`DKV scheduler init failed: ${(err as Error).message}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Create (or replace) the inbox polling cron job.
|
||||||
|
*
|
||||||
|
* Replaces any existing job with the new interval. Called on module init
|
||||||
|
* and by DkvController.saveConfig() after the admin updates the config.
|
||||||
|
*
|
||||||
|
* @param intervalMin - Poll interval in minutes (e.g. 60 = every hour)
|
||||||
|
* @param tenantId - Tenant to process on each tick
|
||||||
|
*/
|
||||||
|
setInterval(intervalMin: number, tenantId?: string): void {
|
||||||
|
if (tenantId) this.activeTenantId = tenantId;
|
||||||
|
|
||||||
|
if (!this.activeTenantId) {
|
||||||
|
this.logger.warn('DKV scheduler: no active tenantId — cron job not created');
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const tenant = this.activeTenantId;
|
||||||
|
|
||||||
|
// Remove existing job if registered
|
||||||
|
try {
|
||||||
|
this.schedulerRegistry.getCronJob(this.JOB_NAME).stop();
|
||||||
|
this.schedulerRegistry.deleteCronJob(this.JOB_NAME);
|
||||||
|
} catch {
|
||||||
|
/* Job not yet registered — this is expected on first call */
|
||||||
|
}
|
||||||
|
|
||||||
|
// Create new cron job with computed expression: every N minutes
|
||||||
|
const cronExpr = `*/${intervalMin} * * * *`;
|
||||||
|
const job = new CronJobClass(cronExpr, () => {
|
||||||
|
this.dkvService.processInbox(tenant).catch((err) =>
|
||||||
|
this.logger.error(
|
||||||
|
`DKV inbox poll failed for tenant ${tenant}: ${(err as Error).message}`,
|
||||||
|
),
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
|
// Cast required: our minimal CronJob type doesn't match cron's full type signature.
|
||||||
|
// At runtime the object IS a full CronJob — SchedulerRegistry only calls stop() on it.
|
||||||
|
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||||
|
this.schedulerRegistry.addCronJob(this.JOB_NAME, job as any);
|
||||||
|
job.start();
|
||||||
|
|
||||||
|
this.logger.log(
|
||||||
|
`DKV cron job registered: every ${intervalMin} minutes for tenant ${tenant}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Stop and remove the inbox polling cron job.
|
||||||
|
* Called by DkvController when admin sets isActive=false in config.
|
||||||
|
*/
|
||||||
|
stopJob(): void {
|
||||||
|
try {
|
||||||
|
this.schedulerRegistry.getCronJob(this.JOB_NAME).stop();
|
||||||
|
this.schedulerRegistry.deleteCronJob(this.JOB_NAME);
|
||||||
|
this.logger.log('DKV cron job stopped and removed');
|
||||||
|
} catch {
|
||||||
|
/* Not registered — no-op */
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,237 @@
|
|||||||
|
import {
|
||||||
|
BadRequestException,
|
||||||
|
Body,
|
||||||
|
Controller,
|
||||||
|
Delete,
|
||||||
|
Get,
|
||||||
|
NotFoundException,
|
||||||
|
Param,
|
||||||
|
Post,
|
||||||
|
Put,
|
||||||
|
Query,
|
||||||
|
Req,
|
||||||
|
Res,
|
||||||
|
UploadedFile,
|
||||||
|
UseInterceptors,
|
||||||
|
} from '@nestjs/common';
|
||||||
|
import { FileInterceptor } from '@nestjs/platform-express';
|
||||||
|
import { Role } from '@prisma/client';
|
||||||
|
import { Roles } from '../auth/decorators/roles.decorator';
|
||||||
|
import { DkvSchedulerService } from './dkv-scheduler.service';
|
||||||
|
import { DkvService } from './dkv.service';
|
||||||
|
import { DkvConfigDto } from './dto/dkv-config.dto';
|
||||||
|
import { DkvHistoryQueryDto } from './dto/dkv-history.dto';
|
||||||
|
import { CreateVehicleDto, UpdateVehicleDto } from './dto/dkv-vehicle.dto';
|
||||||
|
|
||||||
|
/**
|
||||||
|
* DkvController — all /dkv/* routes, ADMIN-only (V4).
|
||||||
|
*
|
||||||
|
* Every handler carries @Roles(Role.ADMIN, Role.SUPER_ADMIN).
|
||||||
|
* Global JwtAuthGuard enforces JWT authentication; RolesGuard enforces the
|
||||||
|
* @Roles decorator. No route is publicly accessible.
|
||||||
|
*
|
||||||
|
* Tenant extraction: `req.tenantId` set by TenantMiddleware (runs after auth guards).
|
||||||
|
* All operations are scoped to the authenticated tenant's data.
|
||||||
|
*
|
||||||
|
* Routes:
|
||||||
|
* GET /dkv/config — get module config (no encrypted creds)
|
||||||
|
* PUT /dkv/config — save module config; updates scheduler
|
||||||
|
* POST /dkv/check-now — manual inbox poll trigger
|
||||||
|
* POST /dkv/test-connection — test inbox connection with form values
|
||||||
|
* GET /dkv/history — paginated processing history
|
||||||
|
* GET /dkv/exports/:filename — download a saved xlsx file
|
||||||
|
* GET /dkv/vehicles — list vehicle master records
|
||||||
|
* POST /dkv/vehicles — create a vehicle master record
|
||||||
|
* PUT /dkv/vehicles/:id — update a vehicle master record
|
||||||
|
* DELETE /dkv/vehicles/:id — delete a vehicle master record
|
||||||
|
* POST /dkv/vehicles/import — bulk-import from CSV upload
|
||||||
|
*/
|
||||||
|
@Controller('dkv')
|
||||||
|
export class DkvController {
|
||||||
|
constructor(
|
||||||
|
private readonly dkvService: DkvService,
|
||||||
|
private readonly dkvScheduler: DkvSchedulerService,
|
||||||
|
) {}
|
||||||
|
|
||||||
|
// ─── Config ────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/** GET /dkv/config — returns module config without encrypted credentials. */
|
||||||
|
@Get('config')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async getConfig(@Req() req: any) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.loadConfig(tenantId);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* PUT /dkv/config — upsert module config.
|
||||||
|
*
|
||||||
|
* After saving, re-applies the cron job if isActive is true,
|
||||||
|
* or stops the cron job if isActive is false.
|
||||||
|
*/
|
||||||
|
@Put('config')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async saveConfig(@Req() req: any, @Body() dto: DkvConfigDto) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
const result = await this.dkvService.saveConfig(tenantId, dto);
|
||||||
|
|
||||||
|
// Update scheduler to reflect the new interval / active state
|
||||||
|
if (dto.isActive && dto.pollIntervalMin) {
|
||||||
|
this.dkvScheduler.setInterval(dto.pollIntervalMin, tenantId);
|
||||||
|
} else if (dto.isActive === false) {
|
||||||
|
this.dkvScheduler.stopJob();
|
||||||
|
}
|
||||||
|
|
||||||
|
return result;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Manual trigger + connection test ──────────────────────────────────────
|
||||||
|
|
||||||
|
/** POST /dkv/check-now — immediately run the inbox processing pipeline. */
|
||||||
|
@Post('check-now')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async checkNow(@Req() req: any) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.checkNow(tenantId);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* POST /dkv/test-connection — test inbox connection with current form values.
|
||||||
|
* Used by the InboxConfigForm "Verbindung testen" button before saving.
|
||||||
|
*/
|
||||||
|
@Post('test-connection')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async testConnection(@Req() req: any, @Body() dto: DkvConfigDto) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
const success = await this.dkvService.testConnection(tenantId, dto);
|
||||||
|
return { success };
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── History ───────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/**
|
||||||
|
* GET /dkv/history?page=1&limit=20 — paginated processing history.
|
||||||
|
* T-07-06: pagination parameters validated by DkvHistoryQueryDto.
|
||||||
|
*/
|
||||||
|
@Get('history')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async getHistory(@Req() req: any, @Query() query: DkvHistoryQueryDto) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
const page = query.page ?? 1;
|
||||||
|
const limit = query.limit ?? 20;
|
||||||
|
return this.dkvService.getHistory(tenantId, page, limit);
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Export file download ──────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/**
|
||||||
|
* GET /dkv/exports/:filename — stream a DKV xlsx export file as an attachment.
|
||||||
|
*
|
||||||
|
* T-07-09: DkvService.getExportFile() validates the filename against the
|
||||||
|
* `DKV_*.xlsx` whitelist pattern before reading from user-files/. Any filename
|
||||||
|
* containing path separators or non-whitelisted characters is rejected.
|
||||||
|
*/
|
||||||
|
@Get('exports/:filename')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async downloadExport(
|
||||||
|
@Req() req: any,
|
||||||
|
@Param('filename') filename: string,
|
||||||
|
@Res() res: any,
|
||||||
|
) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
|
||||||
|
try {
|
||||||
|
const buffer = await this.dkvService.getExportFile(tenantId, filename);
|
||||||
|
res.setHeader('Content-Disposition', `attachment; filename="${filename}"`);
|
||||||
|
res.setHeader(
|
||||||
|
'Content-Type',
|
||||||
|
'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
|
||||||
|
);
|
||||||
|
res.send(buffer);
|
||||||
|
} catch (error) {
|
||||||
|
if (error instanceof NotFoundException || error instanceof BadRequestException) {
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Vehicle Master CRUD ───────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/** GET /dkv/vehicles — list all vehicle master records for this tenant. */
|
||||||
|
@Get('vehicles')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async listVehicles(@Req() req: any) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.listVehicles(tenantId);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** POST /dkv/vehicles — create a new vehicle master record. */
|
||||||
|
@Post('vehicles')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async createVehicle(@Req() req: any, @Body() dto: CreateVehicleDto) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.createVehicle(tenantId, dto);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** PUT /dkv/vehicles/:id — update an existing vehicle master record. */
|
||||||
|
@Put('vehicles/:id')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async updateVehicle(
|
||||||
|
@Req() req: any,
|
||||||
|
@Param('id') id: string,
|
||||||
|
@Body() dto: UpdateVehicleDto,
|
||||||
|
) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.updateVehicle(tenantId, id, dto);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** DELETE /dkv/vehicles/:id — delete a vehicle master record. */
|
||||||
|
@Delete('vehicles/:id')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
async deleteVehicle(@Req() req: any, @Param('id') id: string) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
return this.dkvService.deleteVehicle(tenantId, id);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* POST /dkv/vehicles/import — bulk import from a CSV file upload.
|
||||||
|
*
|
||||||
|
* Accepts a multipart form with:
|
||||||
|
* - `file`: the CSV file (field name must be "file")
|
||||||
|
* - `mode`: 'merge' (default) or 'replace'
|
||||||
|
*
|
||||||
|
* FileInterceptor buffers the upload in memory (no disk write).
|
||||||
|
* The controller reads `file.buffer.toString('utf-8')` and passes to DkvService.
|
||||||
|
*/
|
||||||
|
@Post('vehicles/import')
|
||||||
|
@Roles(Role.ADMIN, Role.SUPER_ADMIN)
|
||||||
|
@UseInterceptors(FileInterceptor('file'))
|
||||||
|
async importVehicles(
|
||||||
|
@Req() req: any,
|
||||||
|
@UploadedFile() file: any,
|
||||||
|
@Body('mode') mode: string,
|
||||||
|
) {
|
||||||
|
const tenantId = this._requireTenant(req);
|
||||||
|
|
||||||
|
if (!file || !file.buffer) {
|
||||||
|
throw new BadRequestException('No CSV file uploaded (field name must be "file")');
|
||||||
|
}
|
||||||
|
|
||||||
|
const csvText = (file.buffer as Buffer).toString('utf-8');
|
||||||
|
const importMode = mode === 'replace' ? 'replace' : 'merge';
|
||||||
|
|
||||||
|
return this.dkvService.importVehiclesCsv(tenantId, csvText, importMode);
|
||||||
|
}
|
||||||
|
|
||||||
|
// ─── Private helpers ───────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
/** Extract and validate tenantId from request; throw BadRequestException when absent. */
|
||||||
|
private _requireTenant(req: any): string {
|
||||||
|
const tenantId = req.tenantId as string | undefined;
|
||||||
|
if (!tenantId) {
|
||||||
|
throw new BadRequestException('No tenant context');
|
||||||
|
}
|
||||||
|
return tenantId;
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user