From d453dbbd1f48597079a839bb2d4d84649bab1484 Mon Sep 17 00:00:00 2001 From: Schalli Date: Thu, 23 Jul 2026 08:50:57 +0200 Subject: [PATCH] feat(13-03): pollDueSources fan-out over all active sources (SCHEMA-03) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace the DÖE-only findUnique with findMany({isActive:true}) fan-out (poll-once-fan-out-many, D-01). Each active TenderSourcePollConfig is resolved through SourceRegistry.get(sourceType) and processed inside its own try/catch (catch-per-source, D-01) — one broken/blocking source no longer aborts the tick for the others. dedupActive = activePortalCount >= 2 (D-05) is computed once per tick and passed to TenderDedupService.resolve(), which now replaces the direct tender.upsert call. Delta-only matchDelta boundary (D-07) preserved: only genuinely-created tender IDs across all sources are collected. Extended tender-ingestion.service.spec.ts: multi-config fan-out, catch-per-source isolation, dedupActive gate assertion, adapter-missing skip, plus the existing SCHEMA-02/D-07/retention/day-cursor suites updated to the new registry+dedup constructor shape (all green). Co-Authored-By: Claude Opus 4.8 (1M context) --- .../tenders/tender-ingestion.service.spec.ts | 244 ++++++++++++++---- .../src/tenders/tender-ingestion.service.ts | 195 +++++++------- 2 files changed, 283 insertions(+), 156 deletions(-) diff --git a/apps/api/src/tenders/tender-ingestion.service.spec.ts b/apps/api/src/tenders/tender-ingestion.service.spec.ts index e05a1ac..868b268 100644 --- a/apps/api/src/tenders/tender-ingestion.service.spec.ts +++ b/apps/api/src/tenders/tender-ingestion.service.spec.ts @@ -4,15 +4,22 @@ import { describe, expect, it, vi } from 'vitest'; import { TenderIngestionService } from './tender-ingestion.service'; /** - * TenderIngestionService.spec — day-cursor gate, SCHEMA-02 change detection, - * and D-05 retention (Plan 10-04, Task 1). + * TenderIngestionService.spec — day-cursor gate, fan-out over active + * sources (SCHEMA-03, D-01), dedupActive gate (D-05), and D-07 delta-only + * matching (Plan 10-04 Task 1, extended Plan 13-03 Task 2). * - * A minimal in-memory fake PrismaService (Maps for tenderSourcePollConfig and - * tender) is used, matching this repo's established test convention of + * A minimal in-memory fake PrismaService (Maps for tenderSourcePollConfig + * and tender) is used, matching this repo's established test convention of * hand-rolled prisma-shaped mocks rather than a real DB connection (see - * ldap.service.spec.ts). doeAdapter/normalizer are also plain fakes — the - * normalizer is stubbed as an identity function so test fixtures can be - * pre-shaped as NormalizedTenderFields directly. + * ldap.service.spec.ts). The registry/adapter/normalizer/dedup collaborators + * are also plain fakes. + * + * SCHEMA-02 change-detection is now exercised through a fake `dedup.resolve` + * that mimics TenderDedupService's real behavior (upsert-by-dedupKey, report + * created-vs-updated) — this proves pollDueSources still wires contentHash- + * driven updates through correctly after the resolve() hand-off, without + * re-testing TenderDedupService's own tier logic (covered by + * tender-dedup.service.spec.ts). */ function berlinTodayStr(): string { @@ -31,20 +38,27 @@ function dayDate(offsetDays: number): Date { return new Date(`${dayString(offsetDays)}T00:00:00Z`); } -function makeFakePrisma() { +function makeFakePrisma(configSeeds: Record = {}) { const tenders = new Map(); const configs = new Map(); - configs.set('doe-opendata', { - id: 'cfg1', - sourceType: 'doe-opendata', - pollIntervalMin: 60, - isActive: true, - lastIngestedDay: null as Date | null, - }); + for (const [sourceType, seed] of Object.entries(configSeeds)) { + configs.set(sourceType, { + id: `cfg-${sourceType}`, + sourceType, + pollIntervalMin: 60, + isActive: true, + lastIngestedDay: null as Date | null, + ...seed, + }); + } const prisma = { tenderSourcePollConfig: { - findUnique: vi.fn(async ({ where }: any) => configs.get(where.sourceType) ?? null), + findMany: vi.fn(async ({ where }: any) => { + const all = [...configs.values()]; + if (where?.isActive === undefined) return all; + return all.filter((c) => c.isActive === where.isActive); + }), update: vi.fn(async ({ where, data }: any) => { const existing = configs.get(where.sourceType); const updated = { ...existing, ...data }; @@ -53,18 +67,6 @@ function makeFakePrisma() { }), }, tender: { - findUnique: vi.fn(async ({ where }: any) => tenders.get(where.dedupKey) ?? null), - upsert: vi.fn(async ({ where, update, create }: any) => { - const existing = tenders.get(where.dedupKey); - if (existing) { - const updated = { ...existing, ...update }; - tenders.set(where.dedupKey, updated); - return updated; - } - const created = { id: `id-${tenders.size + 1}`, ...create }; - tenders.set(where.dedupKey, created); - return created; - }), updateMany: vi.fn(async ({ where, data }: any) => { let count = 0; for (const [key, row] of tenders) { @@ -110,6 +112,7 @@ const BASE_RECORD = { title: 'Test-Ausschreibung', buyerName: 'Stadt Testhausen', cpvCodes: [] as string[], + cpvDivisions: [] as string[], region: null as string | null, plz: null as string | null, bundesland: null as string | null, @@ -122,38 +125,68 @@ const BASE_RECORD = { publishedAt: new Date(), }; +/** + * Fake TenderDedupService.resolve() — upserts by dedupKey into the shared + * `tenders` store and reports created-vs-updated, mirroring the real + * service's contract without re-implementing its three-tier match logic. + */ +function makeFakeDedup(tenders: Map) { + return { + resolve: vi.fn(async (n: any, _opts: { dedupActive: boolean }) => { + const existing = tenders.get(n.dedupKey); + if (existing) { + Object.assign(existing, { + title: n.title, + deadlineAt: n.deadlineAt, + contentHash: n.contentHash, + }); + return { tenderId: existing.id, created: false }; + } + const created = { id: `id-${tenders.size + 1}`, ...n }; + tenders.set(n.dedupKey, created); + return { tenderId: created.id, created: true }; + }), + }; +} + +function makeFakeRegistry(adapters: Record) { + return { + get: vi.fn((sourceType: string) => adapters[sourceType]), + }; +} + function makeService( prisma: any, - doeAdapter: any, + registry: any, normalizer: any, matching: any = { matchDelta: vi.fn() }, + dedup: any = makeFakeDedup(prisma.__store.tenders), ) { - const service = new TenderIngestionService(prisma, doeAdapter, normalizer, matching); + const service = new TenderIngestionService(prisma, registry, normalizer, matching, dedup); // Override the polite catch-up delay so tests don't sleep for real. (service as any).politeDelayMs = 0; return service; } describe('TenderIngestionService.pollDueSources — day-cursor gate (Pitfall A)', () => { - it('makes NO adapter call and returns when the day-cursor is not strictly before Berlin-today (no-op tick, expected)', async () => { - const prisma = makeFakePrisma(); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-1); // -> next day is today, gated out + it('makes NO adapter call when the day-cursor is not strictly before Berlin-today (no-op tick, expected)', async () => { + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-1) } }); // -> next day is today, gated out const doeAdapter = { fetchTenders: vi.fn() }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn() }; - const service = makeService(prisma, doeAdapter, normalizer); + const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); expect(doeAdapter.fetchTenders).not.toHaveBeenCalled(); }); - it('does nothing when the singleton doe-opendata config is not active', async () => { - const prisma = makeFakePrisma(); - prisma.__store.configs.get('doe-opendata').isActive = false; - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3); + it('does nothing when there are no active configs at all', async () => { + const prisma = makeFakePrisma({ 'doe-opendata': { isActive: false, lastIngestedDay: dayDate(-3) } }); const doeAdapter = { fetchTenders: vi.fn() }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn() }; - const service = makeService(prisma, doeAdapter, normalizer); + const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); @@ -161,14 +194,14 @@ describe('TenderIngestionService.pollDueSources — day-cursor gate (Pitfall A)' }); }); -describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection', () => { +describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection (via resolve())', () => { it('inserts a fresh notice, and does NOT duplicate when the identical notice (same dedupKey/contentHash) reappears the next catch-up day', async () => { - const prisma = makeFakePrisma(); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3); + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } }); const record = { ...BASE_RECORD }; const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([record]) }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; - const service = makeService(prisma, doeAdapter, normalizer); + const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); @@ -178,9 +211,8 @@ describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection', expect(prisma.__store.tenders.get('ocds-1').contentHash).toBe('hash-1'); }); - it('updates the existing row in place when the same dedupKey reappears with a changed contentHash (extended deadline)', async () => { - const prisma = makeFakePrisma(); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3); + it('updates the existing row in place when the same dedupKey reappears with a changed contentHash (extended deadline) — SCHEMA-02 preserved through resolve()', async () => { + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } }); const day1Record = { ...BASE_RECORD, contentHash: 'hash-1', @@ -194,8 +226,9 @@ describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection', const doeAdapter = { fetchTenders: vi.fn().mockResolvedValueOnce([day1Record]).mockResolvedValueOnce([day2Record]), }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; - const service = makeService(prisma, doeAdapter, normalizer); + const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); @@ -206,9 +239,105 @@ describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection', }); }); +describe('TenderIngestionService.pollDueSources — fan-out over multiple active sources (SCHEMA-03)', () => { + it('polls every active config, resolving each source through its own registered adapter', async () => { + const prisma = makeFakePrisma({ + 'doe-opendata': { lastIngestedDay: dayDate(-2) }, + 'ai-netserver': { lastIngestedDay: dayDate(-2) }, + }); + const doeAdapter = { + fetchTenders: vi.fn().mockResolvedValue([{ ...BASE_RECORD, dedupKey: 'doe-1', ocid: 'doe-1' }]), + }; + const netserverAdapter = { + fetchTenders: vi.fn().mockResolvedValue([ + { ...BASE_RECORD, sourcePortal: 'ai-netserver', sourceNoticeId: 'ns-1', ocid: null, dedupKey: 'ai-netserver:ns-1' }, + ]), + }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter, 'ai-netserver': netserverAdapter }); + const normalizer = { normalize: vi.fn((raw: any) => raw) }; + const service = makeService(prisma, registry, normalizer); + + await service.pollDueSources(); + + expect(doeAdapter.fetchTenders).toHaveBeenCalled(); + expect(netserverAdapter.fetchTenders).toHaveBeenCalled(); + expect(prisma.__store.tenders.size).toBe(2); + }); + + it('skips a config whose sourceType has no registered adapter, without throwing', async () => { + const prisma = makeFakePrisma({ 'cosinex-dtvp': { lastIngestedDay: dayDate(-2) } }); + const registry = makeFakeRegistry({}); // nothing registered + const normalizer = { normalize: vi.fn() }; + const service = makeService(prisma, registry, normalizer); + + await expect(service.pollDueSources()).resolves.toBeUndefined(); + expect(prisma.__store.tenders.size).toBe(0); + }); +}); + +describe('TenderIngestionService.pollDueSources — catch-per-source error isolation (D-01)', () => { + it('a throwing source does not prevent another active source from being polled, and the tick itself does not throw', async () => { + const prisma = makeFakePrisma({ + 'doe-opendata': { lastIngestedDay: dayDate(-2) }, + 'ai-netserver': { lastIngestedDay: dayDate(-2) }, + }); + const brokenAdapter = { fetchTenders: vi.fn().mockRejectedValue(new Error('portal down')) }; + const healthyAdapter = { + fetchTenders: vi.fn().mockResolvedValue([ + { ...BASE_RECORD, sourcePortal: 'ai-netserver', sourceNoticeId: 'ns-1', ocid: null, dedupKey: 'ai-netserver:ns-1' }, + ]), + }; + const registry = makeFakeRegistry({ 'doe-opendata': brokenAdapter, 'ai-netserver': healthyAdapter }); + const normalizer = { normalize: vi.fn((raw: any) => raw) }; + const service = makeService(prisma, registry, normalizer); + + await expect(service.pollDueSources()).resolves.toBeUndefined(); + + expect(brokenAdapter.fetchTenders).toHaveBeenCalled(); + expect(healthyAdapter.fetchTenders).toHaveBeenCalled(); + expect(prisma.__store.tenders.size).toBe(1); // only the healthy source's tender landed + }); +}); + +describe('TenderIngestionService.pollDueSources — dedupActive gate (D-05)', () => { + it('passes dedupActive=false to resolve() when exactly one config is active', async () => { + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-2) } }); + const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([{ ...BASE_RECORD }]) }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); + const normalizer = { normalize: vi.fn((raw: any) => raw) }; + const dedup = makeFakeDedup(prisma.__store.tenders); + const service = makeService(prisma, registry, normalizer, undefined, dedup); + + await service.pollDueSources(); + + expect(dedup.resolve).toHaveBeenCalledWith(expect.anything(), { dedupActive: false }); + }); + + it('passes dedupActive=true to resolve() when two or more configs are active', async () => { + const prisma = makeFakePrisma({ + 'doe-opendata': { lastIngestedDay: dayDate(-2) }, + 'ai-netserver': { lastIngestedDay: dayDate(-2) }, + }); + const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([{ ...BASE_RECORD, dedupKey: 'doe-1', ocid: 'doe-1' }]) }; + const netserverAdapter = { + fetchTenders: vi.fn().mockResolvedValue([ + { ...BASE_RECORD, sourcePortal: 'ai-netserver', sourceNoticeId: 'ns-1', ocid: null, dedupKey: 'ai-netserver:ns-1' }, + ]), + }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter, 'ai-netserver': netserverAdapter }); + const normalizer = { normalize: vi.fn((raw: any) => raw) }; + const dedup = makeFakeDedup(prisma.__store.tenders); + const service = makeService(prisma, registry, normalizer, undefined, dedup); + + await service.pollDueSources(); + + expect(dedup.resolve).toHaveBeenCalledWith(expect.anything(), { dedupActive: true }); + }); +}); + describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12-01, D-07)', () => { it('calls matching.matchDelta with ONLY the genuinely-new tender IDs of this tick — a pre-existing record reappearing (even with a changed contentHash) is excluded', async () => { - const prisma = makeFakePrisma(); + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-2) } }); // Pre-seed one tender that already exists BEFORE this tick — simulates // a previously-ingested row that reappears in this tick's fetch. prisma.__store.tenders.set('existing-notice', { @@ -217,7 +346,6 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12 contentHash: 'old-hash', status: 'active', }); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-2); // single catch-up day const existingRecord = { ...BASE_RECORD, @@ -230,9 +358,10 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12 const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([existingRecord, newRecord1, newRecord2]), }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; const matching = { matchDelta: vi.fn() }; - const service = makeService(prisma, doeAdapter, normalizer, matching); + const service = makeService(prisma, registry, normalizer, matching); await service.pollDueSources(); @@ -247,20 +376,20 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12 }); it('does NOT call matchDelta when the tick yields zero genuinely-new tenders', async () => { - const prisma = makeFakePrisma(); + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-2) } }); prisma.__store.tenders.set('existing-notice', { id: 'id-existing', dedupKey: 'existing-notice', contentHash: 'old-hash', status: 'active', }); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-2); const existingRecord = { ...BASE_RECORD, dedupKey: 'existing-notice' }; const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([existingRecord]) }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; const matching = { matchDelta: vi.fn() }; - const service = makeService(prisma, doeAdapter, normalizer, matching); + const service = makeService(prisma, registry, normalizer, matching); await service.pollDueSources(); @@ -270,11 +399,11 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12 describe('TenderIngestionService.pollDueSources — catch-up cursor advance', () => { it('advances lastIngestedDay by one day per successful fetch, looping from lastIngestedDay+1 up to today-1', async () => { - const prisma = makeFakePrisma(); - prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3); + const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } }); const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([]) }; + const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn() }; - const service = makeService(prisma, doeAdapter, normalizer); + const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); @@ -315,7 +444,8 @@ describe('TenderIngestionService.pruneExpiredTenders — D-05 retention', () => deadlineAt: dayDate(1), }); - const service = makeService(prisma, { fetchTenders: vi.fn() }, { normalize: vi.fn() }); + const registry = makeFakeRegistry({}); + const service = makeService(prisma, registry, { normalize: vi.fn() }); await service.pruneExpiredTenders(); diff --git a/apps/api/src/tenders/tender-ingestion.service.ts b/apps/api/src/tenders/tender-ingestion.service.ts index 5643d04..aa3ce53 100644 --- a/apps/api/src/tenders/tender-ingestion.service.ts +++ b/apps/api/src/tenders/tender-ingestion.service.ts @@ -1,10 +1,11 @@ import { Injectable, Logger } from '@nestjs/common'; import { PrismaService } from '../prisma/prisma.service'; -import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter'; +import { SourceRegistry } from './source-registry'; +import { TenderDedupService } from './tender-dedup.service'; import { TenderMatchingService } from './tender-matching.service'; import { TenderNormalizerService } from './tender-normalizer.service'; +import type { SourceType } from './tender.types'; -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; @@ -37,9 +38,12 @@ function addOneDay(day: string): string { } /** - * TenderIngestionService — orchestrates the DÖE poll tick: day-cursor gate - * -> fetch -> normalize -> upsert (SCHEMA-02 change detection) -> advance - * cursor -> prune expired (D-05). + * TenderIngestionService — orchestrates the poll tick: fan out over every + * active `TenderSourcePollConfig` (poll-once-fan-out-many, NOT a single + * findUnique — 13-RESEARCH Pattern 3), and per source: day-cursor gate -> + * fetch -> normalize -> `TenderDedupService.resolve()` (SCHEMA-03 + * three-tier dedup, replacing the old direct DÖE-only tender.upsert) -> + * 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` @@ -47,12 +51,25 @@ function addOneDay(day: string): string { * 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. + * Phase 13, Plan 03 (SCHEMA-03): `pollDueSources` now loads ALL active + * configs via `findMany` and resolves each config's `sourceType` through + * `SourceRegistry.get()` — replacing the Phase-10 hardwired single + * `DoeOpenDataAdapter` injection. Each config's day-cursor-gate/fetch/ + * normalize/dedup block runs inside its own try/catch (catch-per-source, + * D-01): one blocking/broken source can never abort the other sources' + * processing within the same tick. `dedupActive = activePortalCount >= 2` + * (D-05) is computed once per tick from `configs.length` and passed to + * every `TenderDedupService.resolve()` call — with only one active source + * it is always false, so the fingerprint dedup tier structurally never + * runs (Erfolgskriterium 3, inert). + * + * Plan 12-01 (NOTIFY-03): pollDueSources collects the IDs of genuinely + * NEW Tender rows created this tick (via `resolve()`'s `created` flag) 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) or were merged as an additional source + * (`created: false`) are NOT included. */ @Injectable() export class TenderIngestionService { @@ -63,115 +80,95 @@ export class TenderIngestionService { constructor( private readonly prisma: PrismaService, - private readonly doeAdapter: DoeOpenDataAdapter, + private readonly registry: SourceRegistry, private readonly normalizer: TenderNormalizerService, private readonly matching: TenderMatchingService, + private readonly dedup: TenderDedupService, ) {} /** - * Single global poll tick — poll-once-fan-out-many (INGEST-06). Called by + * Single global poll tick — poll-once-fan-out-many (INGEST-06), fanned + * out across every active source (SCHEMA-03). 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. + * tick (DKV pattern) AND catch-per-source (D-01) inside the fan-out, so + * neither a scheduler tick failure nor a single broken source crashes the + * process or blocks the other sources. */ async pollDueSources(): Promise { try { - const config = await this.prisma.tenderSourcePollConfig.findUnique({ - where: { sourceType: DOE_SOURCE_TYPE }, + const configs = await this.prisma.tenderSourcePollConfig.findMany({ + where: { isActive: true }, // fan-out, NOT findFirst/findUnique }); - 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 activePortalCount = configs.length; + const dedupActive = activePortalCount >= 2; // D-05 dedup gate 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. + // tender rows created THIS tick (across ALL sources) are collected + // here and handed to matchDelta below — never changed/re-seen rows, + // never merged-as-additional-source rows, never historical rows. const newTenderIds: string[] = []; + let anyDayFetched = false; - 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); + for (const config of configs) { + try { + const adapter = this.registry.get(config.sourceType as SourceType); + if (!adapter) { + this.logger.warn( + `No adapter registered for sourceType '${config.sourceType}' — skipping this tick`, + ); + continue; + } + + let cursorDay = nextDayToFetch(config.lastIngestedDay); + if (!cursorDay) { + this.logger.debug( + `Tender poll tick: no new day yet for '${config.sourceType}' (day-cursor gate) — expected no-op, not a bug`, + ); + continue; + } + + 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 source this tick. + await this._delay(this.politeDelayMs); + } + isFirstFetch = false; + + const rawRecords = await adapter.fetchTenders(cursorDay); + const normalized = rawRecords.map((r) => this.normalizer.normalize(r)); + + for (const tender of normalized) { + const { tenderId, created } = await this.dedup.resolve(tender, { + dedupActive, + }); + if (created) newTenderIds.push(tenderId); // genuinely NEW row this tick (D-07) + } + + await this.prisma.tenderSourcePollConfig.update({ + where: { sourceType: config.sourceType }, + data: { lastIngestedDay: new Date(`${cursorDay}T00:00:00Z`) }, + }); + + anyDayFetched = true; + cursorDay = addOneDay(cursorDay); + } + } catch (err) { + // catch-per-source (D-01): a blocking/broken source must never + // abort the fan-out for the remaining sources. + this.logger.error( + `Tender poll tick failed for source '${config.sourceType}': ${(err as Error).message}`, + ); } - 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) { + if (anyDayFetched) { await this.pruneExpiredTenders(); }