diff --git a/apps/api/src/tenders/tender-ingestion.service.ts b/apps/api/src/tenders/tender-ingestion.service.ts new file mode 100644 index 0000000..bbf8611 --- /dev/null +++ b/apps/api/src/tenders/tender-ingestion.service.ts @@ -0,0 +1,181 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { PrismaService } from '../prisma/prisma.service'; +import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter'; +import { TenderNormalizerService } from './tender-normalizer.service'; + +const DOE_SOURCE_TYPE = 'doe-opendata'; +const RETENTION_DAYS = 90; +/** Polite delay between successive catch-up day-fetches (RESEARCH Open Question 1). */ +const CATCH_UP_DELAY_MS = 1_500; +const MS_PER_DAY = 86_400_000; + +/** + * Day-cursor gate (RESEARCH.md "Day-cursor gate" snippet) — replaces a + * `since: Date` timestamp cursor, because the DÖE OpenData API has no + * incremental/since parameter; it is a daily batch-export API only. + * + * Returns the next `YYYY-MM-DD` to fetch, or `null` when nothing new is + * available yet (dayCursor is not strictly before Europe/Berlin "today" — + * D-01's "from now" semantics: the first eligible day for a never-polled + * config is today itself, which is gated out until the next calendar day). + */ +export function nextDayToFetch(lastIngestedDay: Date | null): string | null { + const berlinToday = new Date().toLocaleDateString('en-CA', { + timeZone: 'Europe/Berlin', + }); // YYYY-MM-DD + const next = lastIngestedDay + ? new Date(lastIngestedDay.getTime() + MS_PER_DAY).toISOString().slice(0, 10) + : new Date().toISOString().slice(0, 10); + return next < berlinToday ? next : null; // null = no-op this tick (Pitfall A) +} + +function addOneDay(day: string): string { + return new Date(new Date(`${day}T00:00:00Z`).getTime() + MS_PER_DAY) + .toISOString() + .slice(0, 10); +} + +/** + * TenderIngestionService — orchestrates the DÖE poll tick: day-cursor gate + * -> fetch -> normalize -> upsert (SCHEMA-02 change detection) -> advance + * cursor -> prune expired (D-05). + * + * Multi-tenant safety (D-03, T-10-09): uses the plain, non-tenant-scoped + * PrismaService injected as-is. Never wraps `Tender`/`TenderSourcePollConfig` + * queries in the tenant RLS extension — these are platform-global, + * RLS-exempt tables. Applying that extension here would silently filter + * out platform data for a 2nd tenant's session. + */ +@Injectable() +export class TenderIngestionService { + private readonly logger = new Logger(TenderIngestionService.name); + + /** Overridable (tests zero this out to avoid real sleeps). */ + protected politeDelayMs = CATCH_UP_DELAY_MS; + + constructor( + private readonly prisma: PrismaService, + private readonly doeAdapter: DoeOpenDataAdapter, + private readonly normalizer: TenderNormalizerService, + ) {} + + /** + * Single global poll tick — poll-once-fan-out-many (INGEST-06). Called by + * TenderSchedulerService's cron tick, on a shared, admin-configurable + * interval, regardless of tenant count. Never throws — catch-and-log per + * tick (DKV pattern), so a scheduler tick failure never crashes the process. + */ + async pollDueSources(): Promise { + try { + const config = await this.prisma.tenderSourcePollConfig.findUnique({ + where: { sourceType: DOE_SOURCE_TYPE }, + }); + if (!config?.isActive) return; + + let cursorDay = nextDayToFetch(config.lastIngestedDay); + if (!cursorDay) { + this.logger.debug( + 'Tender poll tick: no new DÖE day yet (day-cursor gate) — expected no-op, not a bug', + ); + return; + } + + const berlinToday = new Date().toLocaleDateString('en-CA', { + timeZone: 'Europe/Berlin', + }); + let fetchedAtLeastOneDay = false; + let isFirstFetch = true; + + while (cursorDay && cursorDay < berlinToday) { + if (!isFirstFetch) { + // Politeness delay between successive catch-up day-fetches only — + // never before the first fetch of a tick. + await this._delay(this.politeDelayMs); + } + isFirstFetch = false; + + const rawRecords = await this.doeAdapter.fetchTenders(cursorDay); + const normalized = rawRecords.map((r) => this.normalizer.normalize(r)); + + for (const tender of normalized) { + await this.prisma.tender.upsert({ + where: { dedupKey: tender.dedupKey }, + update: { + title: tender.title, + buyerName: tender.buyerName, + cpvCodes: tender.cpvCodes, + region: tender.region, + plz: tender.plz, + bundesland: tender.bundesland, + deadlineAt: tender.deadlineAt, + estimatedValue: tender.estimatedValue, + procedureType: tender.procedureType, + sourceUrl: tender.sourceUrl, + contentHash: tender.contentHash, + }, + create: { + sourcePortal: tender.sourcePortal, + sourceNoticeId: tender.sourceNoticeId, + ocid: tender.ocid, + dedupKey: tender.dedupKey, + title: tender.title, + buyerName: tender.buyerName, + cpvCodes: tender.cpvCodes, + region: tender.region, + plz: tender.plz, + bundesland: tender.bundesland, + deadlineAt: tender.deadlineAt, + estimatedValue: tender.estimatedValue, + procedureType: tender.procedureType, + status: tender.status, + sourceUrl: tender.sourceUrl, + contentHash: tender.contentHash, + publishedAt: tender.publishedAt, + }, + }); + } + + await this.prisma.tenderSourcePollConfig.update({ + where: { sourceType: DOE_SOURCE_TYPE }, + data: { lastIngestedDay: new Date(`${cursorDay}T00:00:00Z`) }, + }); + + fetchedAtLeastOneDay = true; + cursorDay = addOneDay(cursorDay); + } + + if (fetchedAtLeastOneDay) { + await this.pruneExpiredTenders(); + } + } catch (err) { + this.logger.error(`Tender poll tick failed: ${(err as Error).message}`); + } + } + + /** + * D-05 retention: mark past-deadline active rows 'expired', then delete + * 'expired' rows whose deadline is older than the 90-day retention window. + * Rows with `deadlineAt IS NULL` are excluded from BOTH operations — a + * wrongly-deleted no-deadline tender is unrecoverable (RESEARCH Pattern 4 + * recommendation (a)); Prisma's `lt` comparison against null is always + * false at the SQL level, so this exclusion holds structurally, not just + * by convention. + */ + async pruneExpiredTenders(): Promise { + const now = new Date(); + const retentionCutoff = new Date(now.getTime() - RETENTION_DAYS * MS_PER_DAY); + + await this.prisma.tender.updateMany({ + where: { status: 'active', deadlineAt: { lt: now } }, + data: { status: 'expired' }, + }); + + await this.prisma.tender.deleteMany({ + where: { status: 'expired', deadlineAt: { lt: retentionCutoff } }, + }); + } + + private _delay(ms: number): Promise { + return new Promise((resolve) => setTimeout(resolve, ms)); + } +} diff --git a/apps/api/src/tenders/tenders.module.ts b/apps/api/src/tenders/tenders.module.ts index bff0ead..7e50ac1 100644 --- a/apps/api/src/tenders/tenders.module.ts +++ b/apps/api/src/tenders/tenders.module.ts @@ -4,16 +4,19 @@ import { ModuleRegistryService } from '../module-registry/module-registry.servic import { PrismaService } from '../prisma/prisma.service'; import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter'; import { seedTendersModule } from './tenders.seed'; +import { TenderIngestionService } from './tender-ingestion.service'; import { TenderNormalizerService } from './tender-normalizer.service'; /** * NestJS module for the Ausschreibungs-Radar feature. * * Wave 2 registered the module in the marketplace and seeded the - * singleton DÖE poll config so the shared poll is admin-drivable. This - * plan (03) adds the DÖE source adapter + normalizer (parse+map core of - * INGEST-01/SCHEMA-01). Ingestion orchestration/scheduler/controller - * wiring is added in Plans 04-05. + * singleton DÖE poll config so the shared poll is admin-drivable. Plan 03 + * added the DÖE source adapter + normalizer (parse+map core of + * INGEST-01/SCHEMA-01). This plan (04) adds ingestion orchestration + * (TenderIngestionService: day-cursor gate, SCHEMA-02 change detection, + * D-05 retention) and the shared global scheduler. Controller/DTO wiring + * is added in Plan 05. * * PrismaModule is global (no explicit import needed). * @@ -23,7 +26,7 @@ import { TenderNormalizerService } from './tender-normalizer.service'; @Module({ imports: [ModuleRegistryModule], controllers: [], - providers: [DoeOpenDataAdapter, TenderNormalizerService], + providers: [DoeOpenDataAdapter, TenderNormalizerService, TenderIngestionService], }) export class TendersModule implements OnModuleInit { private readonly logger = new Logger(TendersModule.name);