diff --git a/apps/api/src/tenders/tender-ingestion.service.spec.ts b/apps/api/src/tenders/tender-ingestion.service.spec.ts index 7e2a31b..e05a1ac 100644 --- a/apps/api/src/tenders/tender-ingestion.service.spec.ts +++ b/apps/api/src/tenders/tender-ingestion.service.spec.ts @@ -53,6 +53,7 @@ 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) { @@ -121,8 +122,13 @@ const BASE_RECORD = { publishedAt: new Date(), }; -function makeService(prisma: any, doeAdapter: any, normalizer: any) { - const service = new TenderIngestionService(prisma, doeAdapter, normalizer); +function makeService( + prisma: any, + doeAdapter: any, + normalizer: any, + matching: any = { matchDelta: vi.fn() }, +) { + const service = new TenderIngestionService(prisma, doeAdapter, normalizer, matching); // Override the polite catch-up delay so tests don't sleep for real. (service as any).politeDelayMs = 0; return service; @@ -200,6 +206,68 @@ describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection', }); }); +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(); + // 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', { + id: 'id-existing', + dedupKey: 'existing-notice', + contentHash: 'old-hash', + status: 'active', + }); + prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-2); // single catch-up day + + const existingRecord = { + ...BASE_RECORD, + dedupKey: 'existing-notice', + contentHash: 'updated-hash', // changed, but NOT genuinely new + }; + const newRecord1 = { ...BASE_RECORD, dedupKey: 'new-notice-1' }; + const newRecord2 = { ...BASE_RECORD, dedupKey: 'new-notice-2' }; + + const doeAdapter = { + fetchTenders: vi.fn().mockResolvedValue([existingRecord, newRecord1, newRecord2]), + }; + const normalizer = { normalize: vi.fn((raw: any) => raw) }; + const matching = { matchDelta: vi.fn() }; + const service = makeService(prisma, doeAdapter, normalizer, matching); + + await service.pollDueSources(); + + expect(matching.matchDelta).toHaveBeenCalledTimes(1); + const calledIds: string[] = matching.matchDelta.mock.calls[0][0]; + const newIds = [ + prisma.__store.tenders.get('new-notice-1').id, + prisma.__store.tenders.get('new-notice-2').id, + ]; + expect(calledIds.sort()).toEqual(newIds.sort()); + expect(calledIds).not.toContain('id-existing'); + }); + + it('does NOT call matchDelta when the tick yields zero genuinely-new tenders', async () => { + const prisma = makeFakePrisma(); + 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 normalizer = { normalize: vi.fn((raw: any) => raw) }; + const matching = { matchDelta: vi.fn() }; + const service = makeService(prisma, doeAdapter, normalizer, matching); + + await service.pollDueSources(); + + expect(matching.matchDelta).not.toHaveBeenCalled(); + }); +}); + 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(); diff --git a/apps/api/src/tenders/tender-ingestion.service.ts b/apps/api/src/tenders/tender-ingestion.service.ts index fce97bc..5643d04 100644 --- a/apps/api/src/tenders/tender-ingestion.service.ts +++ b/apps/api/src/tenders/tender-ingestion.service.ts @@ -1,6 +1,7 @@ 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'; @@ -45,6 +46,13 @@ function addOneDay(day: string): string { * 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 { @@ -57,6 +65,7 @@ export class TenderIngestionService { private readonly prisma: PrismaService, private readonly doeAdapter: DoeOpenDataAdapter, private readonly normalizer: TenderNormalizerService, + private readonly matching: TenderMatchingService, ) {} /** @@ -85,6 +94,10 @@ export class TenderIngestionService { }); 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) { @@ -98,7 +111,17 @@ export class TenderIngestionService { const normalized = rawRecords.map((r) => this.normalizer.normalize(r)); for (const tender of normalized) { - await this.prisma.tender.upsert({ + // 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, @@ -135,6 +158,8 @@ export class TenderIngestionService { publishedAt: tender.publishedAt, }, }); + + if (!existing) newTenderIds.push(saved.id); // genuinely NEW row this tick (D-07) } await this.prisma.tenderSourcePollConfig.update({ @@ -149,6 +174,13 @@ export class TenderIngestionService { 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}`); } diff --git a/apps/api/src/tenders/tenders.module.ts b/apps/api/src/tenders/tenders.module.ts index a1bc9d1..20ac873 100644 --- a/apps/api/src/tenders/tenders.module.ts +++ b/apps/api/src/tenders/tenders.module.ts @@ -5,6 +5,7 @@ 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 { TenderMatchingService } from './tender-matching.service'; import { TenderNormalizerService } from './tender-normalizer.service'; import { TenderSavedSearchService } from './tender-saved-search.service'; import { TenderSchedulerService } from './tender-scheduler.service'; @@ -31,6 +32,12 @@ import { TendersController } from './tenders.controller'; * as a further provider — same userId-scoping convention as * TenderTriageService, no forTenant()/RLS (Pitfall 4). * + * Phase 12, Plan 01 (NOTIFY-03) adds TenderMatchingService: injected into + * TenderIngestionService and called at the end of pollDueSources with only + * the genuinely-new tender IDs of the current tick (delta-only matching, + * D-07) — reuses buildTenderWhere against every active TenderSavedSearch + * profile and upserts TenderMatch rows (matched-vs-notified state, D-06). + * * Seeds itself into the module registry on application startup via * OnModuleInit lifecycle hook — same pattern as DkvModule. */ @@ -44,6 +51,7 @@ import { TendersController } from './tenders.controller'; TenderSchedulerService, TenderTriageService, TenderSavedSearchService, + TenderMatchingService, ], }) export class TendersModule implements OnModuleInit {