feat(10-04): implement TenderSchedulerService — single global cron (poll-once-fan-out-many)
- One named cron job 'tender-doe-poll' for the whole platform; setInterval() takes no tenant argument (INGEST-06) — reuses DkvSchedulerService's CronJob require()-resolution + SchedulerRegistry mechanics, drops the per-tenant activeTenantId framing entirely - onModuleInit() loads the singleton doe-opendata config via findUnique on the fixed sourceType slug, never findFirst (Pitfall D) - Day-cursor gate stays inside TenderIngestionService.pollDueSources() — this scheduler only controls cron-tick frequency (Pitfall A separation) - Registered in TendersModule.providers; ScheduleModule already global via AppModule, no re-registration needed
This commit is contained in:
@@ -0,0 +1,146 @@
|
|||||||
|
import { Injectable, Logger, OnModuleInit } from '@nestjs/common';
|
||||||
|
import { SchedulerRegistry } from '@nestjs/schedule';
|
||||||
|
import { PrismaService } from '../prisma/prisma.service';
|
||||||
|
import { TenderIngestionService } from './tender-ingestion.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. Reuses the exact DkvSchedulerService resolution workaround verbatim.
|
||||||
|
*/
|
||||||
|
// 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 };
|
||||||
|
|
||||||
|
/**
|
||||||
|
* TenderSchedulerService — a SINGLE global cron job driving the shared DÖE
|
||||||
|
* poll (INGEST-06, D-04). This is the poll-once-fan-out-many heart of the
|
||||||
|
* phase: unlike `DkvSchedulerService`, there is NO tenant dimension at all.
|
||||||
|
*
|
||||||
|
* Deliberate deviations from the DKV scheduler this reuses mechanics from
|
||||||
|
* (RESEARCH.md Pitfall D / PATTERNS.md ANTI-PATTERN note):
|
||||||
|
* - No per-tenant "active tenant id" instance field at all — DÖE config is
|
||||||
|
* a genuine platform-wide singleton, not a per-tenant config to track.
|
||||||
|
* - `setInterval()` takes NO tenant argument — there is exactly one cron
|
||||||
|
* job (`tender-doe-poll`) for the whole platform, regardless of how many
|
||||||
|
* tenants activate the `tender-radar` module.
|
||||||
|
* - `onModuleInit()` loads the singleton config via `findUnique` on the
|
||||||
|
* fixed `sourceType: 'doe-opendata'` slug — a fixed-slug lookup, not an
|
||||||
|
* unfiltered/ordered "first match" query — to make the "one config, no
|
||||||
|
* tenant iteration" intent explicit and self-documenting against future
|
||||||
|
* copy-paste into a per-tenant source.
|
||||||
|
* - The day-cursor gate (whether an HTTP call actually happens) is NOT
|
||||||
|
* here — it lives inside `TenderIngestionService.pollDueSources()`. This
|
||||||
|
* scheduler only controls cron-tick frequency (Pitfall A separation).
|
||||||
|
*
|
||||||
|
* The two-tenant safety property (Success Criteria 4 & 5) is proven by
|
||||||
|
* `tender-scheduler.service.spec.ts` (Plan 04, Task 3): activating the
|
||||||
|
* module for a 2nd tenant must trigger zero additional DÖE calls, zero
|
||||||
|
* additional cron jobs, and zero additional Tender rows.
|
||||||
|
*/
|
||||||
|
@Injectable()
|
||||||
|
export class TenderSchedulerService implements OnModuleInit {
|
||||||
|
private readonly logger = new Logger(TenderSchedulerService.name);
|
||||||
|
|
||||||
|
/** Name of the single, platform-global managed cron job. */
|
||||||
|
private readonly JOB_NAME = 'tender-doe-poll';
|
||||||
|
|
||||||
|
constructor(
|
||||||
|
private readonly schedulerRegistry: SchedulerRegistry,
|
||||||
|
private readonly tenderIngestionService: TenderIngestionService,
|
||||||
|
private readonly prisma: PrismaService,
|
||||||
|
) {}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* On application startup: load the singleton doe-opendata poll config and
|
||||||
|
* register the single global cron job if active. Errors are caught and
|
||||||
|
* logged (never re-thrown) so a missing/broken config does not prevent
|
||||||
|
* the rest of the application from starting.
|
||||||
|
*/
|
||||||
|
async onModuleInit(): Promise<void> {
|
||||||
|
try {
|
||||||
|
const config = await this.prisma.tenderSourcePollConfig.findUnique({
|
||||||
|
where: { sourceType: 'doe-opendata' },
|
||||||
|
});
|
||||||
|
|
||||||
|
if (config?.isActive) {
|
||||||
|
this.setInterval(config.pollIntervalMin);
|
||||||
|
this.logger.log(
|
||||||
|
`Tender scheduler initialized: every ${config.pollIntervalMin} min (single global job — platform-wide, no tenant dimension)`,
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
this.logger.log(
|
||||||
|
'Tender scheduler: doe-opendata config inactive — cron job not registered',
|
||||||
|
);
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
this.logger.error(
|
||||||
|
`Tender scheduler init failed: ${(err as Error).message}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Create (or replace) the single global DÖE-poll cron job.
|
||||||
|
*
|
||||||
|
* Called on module init and by the (Plan 05) admin source-config route
|
||||||
|
* after the platform-wide config is updated, so the cron job reflects any
|
||||||
|
* admin change immediately — without a service restart.
|
||||||
|
*
|
||||||
|
* @param intervalMin - Cron-tick frequency in minutes (D-04 default 60).
|
||||||
|
* This is NOT the DÖE fetch frequency — the actual upstream HTTP call
|
||||||
|
* only happens when TenderIngestionService's day-cursor gate allows it
|
||||||
|
* (Pitfall A: cron tick frequency != actual-fetch frequency for this
|
||||||
|
* day-granular source).
|
||||||
|
*/
|
||||||
|
setInterval(intervalMin: number): void {
|
||||||
|
// 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 */
|
||||||
|
}
|
||||||
|
|
||||||
|
// Standard cron minute field only accepts 0-59; for longer intervals use the hours field.
|
||||||
|
const cronExpr =
|
||||||
|
intervalMin < 60
|
||||||
|
? `*/${intervalMin} * * * *`
|
||||||
|
: `0 */${Math.floor(intervalMin / 60)} * * *`;
|
||||||
|
|
||||||
|
const job = new CronJobClass(cronExpr, () => {
|
||||||
|
this.tenderIngestionService.pollDueSources().catch((err) =>
|
||||||
|
this.logger.error(
|
||||||
|
`Tender poll tick failed: ${(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(
|
||||||
|
`Tender cron job registered: every ${intervalMin} minutes (single global job, no tenant parameter)`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Stop and remove the DÖE-poll cron job.
|
||||||
|
* Called by the (Plan 05) admin route when isActive=false is saved.
|
||||||
|
*/
|
||||||
|
stopJob(): void {
|
||||||
|
try {
|
||||||
|
this.schedulerRegistry.getCronJob(this.JOB_NAME).stop();
|
||||||
|
this.schedulerRegistry.deleteCronJob(this.JOB_NAME);
|
||||||
|
this.logger.log('Tender cron job stopped and removed');
|
||||||
|
} catch {
|
||||||
|
/* Not registered — no-op */
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -6,6 +6,7 @@ import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter';
|
|||||||
import { seedTendersModule } from './tenders.seed';
|
import { seedTendersModule } from './tenders.seed';
|
||||||
import { TenderIngestionService } from './tender-ingestion.service';
|
import { TenderIngestionService } from './tender-ingestion.service';
|
||||||
import { TenderNormalizerService } from './tender-normalizer.service';
|
import { TenderNormalizerService } from './tender-normalizer.service';
|
||||||
|
import { TenderSchedulerService } from './tender-scheduler.service';
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* NestJS module for the Ausschreibungs-Radar feature.
|
* NestJS module for the Ausschreibungs-Radar feature.
|
||||||
@@ -26,7 +27,12 @@ import { TenderNormalizerService } from './tender-normalizer.service';
|
|||||||
@Module({
|
@Module({
|
||||||
imports: [ModuleRegistryModule],
|
imports: [ModuleRegistryModule],
|
||||||
controllers: [],
|
controllers: [],
|
||||||
providers: [DoeOpenDataAdapter, TenderNormalizerService, TenderIngestionService],
|
providers: [
|
||||||
|
DoeOpenDataAdapter,
|
||||||
|
TenderNormalizerService,
|
||||||
|
TenderIngestionService,
|
||||||
|
TenderSchedulerService,
|
||||||
|
],
|
||||||
})
|
})
|
||||||
export class TendersModule implements OnModuleInit {
|
export class TendersModule implements OnModuleInit {
|
||||||
private readonly logger = new Logger(TendersModule.name);
|
private readonly logger = new Logger(TendersModule.name);
|
||||||
|
|||||||
Reference in New Issue
Block a user