feat(10-04): implement TenderIngestionService — day-cursor gate, upsert change-detect, D-05 retention
- pollDueSources(): singleton doe-opendata config via findUnique (fixed slug,
not findFirst); day-cursor gate (nextDayToFetch) no-ops when nothing new
(Pitfall A); catch-up loop from lastIngestedDay+1 to today-1 with a polite
1.5s delay between successive day-fetches
- prisma.tender.upsert({ where: { dedupKey } }) — SCHEMA-02 change-detection
seam: identical notice does not duplicate, changed contentHash updates in
place
- pruneExpiredTenders(): marks active+past-deadline rows 'expired', deletes
expired rows older than 90 days, never touches deadlineAt=null rows (D-05)
- Plain PrismaService throughout — no tenant RLS extension on the global
Tender/TenderSourcePollConfig tables (D-03, T-10-09)
- Registered in TendersModule.providers
This commit is contained in:
@@ -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<void> {
|
||||||
|
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<void> {
|
||||||
|
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<void> {
|
||||||
|
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -4,16 +4,19 @@ import { ModuleRegistryService } from '../module-registry/module-registry.servic
|
|||||||
import { PrismaService } from '../prisma/prisma.service';
|
import { PrismaService } from '../prisma/prisma.service';
|
||||||
import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter';
|
import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter';
|
||||||
import { seedTendersModule } from './tenders.seed';
|
import { seedTendersModule } from './tenders.seed';
|
||||||
|
import { TenderIngestionService } from './tender-ingestion.service';
|
||||||
import { TenderNormalizerService } from './tender-normalizer.service';
|
import { TenderNormalizerService } from './tender-normalizer.service';
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* NestJS module for the Ausschreibungs-Radar feature.
|
* NestJS module for the Ausschreibungs-Radar feature.
|
||||||
*
|
*
|
||||||
* Wave 2 registered the module in the marketplace and seeded the
|
* Wave 2 registered the module in the marketplace and seeded the
|
||||||
* singleton DÖE poll config so the shared poll is admin-drivable. This
|
* singleton DÖE poll config so the shared poll is admin-drivable. Plan 03
|
||||||
* plan (03) adds the DÖE source adapter + normalizer (parse+map core of
|
* added the DÖE source adapter + normalizer (parse+map core of
|
||||||
* INGEST-01/SCHEMA-01). Ingestion orchestration/scheduler/controller
|
* INGEST-01/SCHEMA-01). This plan (04) adds ingestion orchestration
|
||||||
* wiring is added in Plans 04-05.
|
* (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).
|
* PrismaModule is global (no explicit import needed).
|
||||||
*
|
*
|
||||||
@@ -23,7 +26,7 @@ import { TenderNormalizerService } from './tender-normalizer.service';
|
|||||||
@Module({
|
@Module({
|
||||||
imports: [ModuleRegistryModule],
|
imports: [ModuleRegistryModule],
|
||||||
controllers: [],
|
controllers: [],
|
||||||
providers: [DoeOpenDataAdapter, TenderNormalizerService],
|
providers: [DoeOpenDataAdapter, TenderNormalizerService, TenderIngestionService],
|
||||||
})
|
})
|
||||||
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