import { Injectable, Logger } from '@nestjs/common'; import { PrismaService } from '../prisma/prisma.service'; import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter'; import { TenderMatchingService } from './tender-matching.service'; 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. * * Plan 12-01 (NOTIFY-03): pollDueSources now collects the IDs of genuinely * NEW Tender rows created this tick (via an indexed pre-check against the * existing upsert) and calls `TenderMatchingService.matchDelta` with only * those IDs at the end of the tick — this is the delta-only matching * boundary (D-07) that structurally prevents a backfill flood. Rows that * merely changed (SCHEMA-02 contentHash update) are NOT included. */ @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, private readonly matching: TenderMatchingService, ) {} /** * 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; // Delta-only matching boundary (D-07, NOTIFY-03): only genuinely NEW // tender rows created THIS tick are collected here and handed to // matchDelta below — never changed/re-seen rows, never historical rows. const newTenderIds: string[] = []; 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) { // Indexed pre-check (dedupKey @unique) — the plain upsert below // does not report create-vs-update, so we check existence first // to determine whether this row is genuinely NEW this tick // (D-07 delta boundary). Negligible cost: a handful of records/ // day for the DÖE source. const existing = await this.prisma.tender.findUnique({ where: { dedupKey: tender.dedupKey }, select: { id: true }, }); const saved = await this.prisma.tender.upsert({ where: { dedupKey: tender.dedupKey }, update: { title: tender.title, buyerName: tender.buyerName, cpvCodes: tender.cpvCodes, cpvDivisions: tender.cpvDivisions, 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, cpvDivisions: tender.cpvDivisions, 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, }, }); if (!existing) newTenderIds.push(saved.id); // genuinely NEW row this tick (D-07) } 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(); } if (newTenderIds.length) { // Matching failures must never crash the ingestion tick — the // outer catch-and-log already covers this, but matchDelta itself // also catches per-profile (defense in depth). await this.matching.matchDelta(newTenderIds); } } 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)); } }