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