feat(12-01): wire matchDelta into pollDueSources; register TenderMatchingService

pollDueSources now collects genuinely-new tender IDs via an indexed
dedupKey pre-check (existing upsert doesn't report create-vs-update),
and calls TenderMatchingService.matchDelta(newTenderIds) once at the
end of the tick — the delta-only matching boundary (D-07). Changed/
re-seen rows are excluded, only genuinely new rows trigger matching.

TenderMatchingService registered as a provider in TendersModule and
injected into TenderIngestionService. Ingestion spec extended to
assert matchDelta receives only the new IDs, and is not called when
no new tenders were ingested this tick.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-07-22 09:05:41 +02:00
parent 92d4c969de
commit b45047f0b3
3 changed files with 111 additions and 3 deletions
@@ -53,6 +53,7 @@ function makeFakePrisma() {
}), }),
}, },
tender: { tender: {
findUnique: vi.fn(async ({ where }: any) => tenders.get(where.dedupKey) ?? null),
upsert: vi.fn(async ({ where, update, create }: any) => { upsert: vi.fn(async ({ where, update, create }: any) => {
const existing = tenders.get(where.dedupKey); const existing = tenders.get(where.dedupKey);
if (existing) { if (existing) {
@@ -121,8 +122,13 @@ const BASE_RECORD = {
publishedAt: new Date(), publishedAt: new Date(),
}; };
function makeService(prisma: any, doeAdapter: any, normalizer: any) { function makeService(
const service = new TenderIngestionService(prisma, doeAdapter, normalizer); 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. // Override the polite catch-up delay so tests don't sleep for real.
(service as any).politeDelayMs = 0; (service as any).politeDelayMs = 0;
return service; 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', () => { 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 () => { it('advances lastIngestedDay by one day per successful fetch, looping from lastIngestedDay+1 up to today-1', async () => {
const prisma = makeFakePrisma(); const prisma = makeFakePrisma();
@@ -1,6 +1,7 @@
import { Injectable, Logger } from '@nestjs/common'; import { Injectable, Logger } from '@nestjs/common';
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 { TenderMatchingService } from './tender-matching.service';
import { TenderNormalizerService } from './tender-normalizer.service'; import { TenderNormalizerService } from './tender-normalizer.service';
const DOE_SOURCE_TYPE = 'doe-opendata'; 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, * queries in the tenant RLS extension — these are platform-global,
* RLS-exempt tables. Applying that extension here would silently filter * RLS-exempt tables. Applying that extension here would silently filter
* out platform data for a 2nd tenant's session. * 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() @Injectable()
export class TenderIngestionService { export class TenderIngestionService {
@@ -57,6 +65,7 @@ export class TenderIngestionService {
private readonly prisma: PrismaService, private readonly prisma: PrismaService,
private readonly doeAdapter: DoeOpenDataAdapter, private readonly doeAdapter: DoeOpenDataAdapter,
private readonly normalizer: TenderNormalizerService, private readonly normalizer: TenderNormalizerService,
private readonly matching: TenderMatchingService,
) {} ) {}
/** /**
@@ -85,6 +94,10 @@ export class TenderIngestionService {
}); });
let fetchedAtLeastOneDay = false; let fetchedAtLeastOneDay = false;
let isFirstFetch = true; 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) { while (cursorDay && cursorDay < berlinToday) {
if (!isFirstFetch) { if (!isFirstFetch) {
@@ -98,7 +111,17 @@ export class TenderIngestionService {
const normalized = rawRecords.map((r) => this.normalizer.normalize(r)); const normalized = rawRecords.map((r) => this.normalizer.normalize(r));
for (const tender of normalized) { 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 }, where: { dedupKey: tender.dedupKey },
update: { update: {
title: tender.title, title: tender.title,
@@ -135,6 +158,8 @@ export class TenderIngestionService {
publishedAt: tender.publishedAt, publishedAt: tender.publishedAt,
}, },
}); });
if (!existing) newTenderIds.push(saved.id); // genuinely NEW row this tick (D-07)
} }
await this.prisma.tenderSourcePollConfig.update({ await this.prisma.tenderSourcePollConfig.update({
@@ -149,6 +174,13 @@ export class TenderIngestionService {
if (fetchedAtLeastOneDay) { if (fetchedAtLeastOneDay) {
await this.pruneExpiredTenders(); 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) { } catch (err) {
this.logger.error(`Tender poll tick failed: ${(err as Error).message}`); this.logger.error(`Tender poll tick failed: ${(err as Error).message}`);
} }
+8
View File
@@ -5,6 +5,7 @@ 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 { TenderIngestionService } from './tender-ingestion.service';
import { TenderMatchingService } from './tender-matching.service';
import { TenderNormalizerService } from './tender-normalizer.service'; import { TenderNormalizerService } from './tender-normalizer.service';
import { TenderSavedSearchService } from './tender-saved-search.service'; import { TenderSavedSearchService } from './tender-saved-search.service';
import { TenderSchedulerService } from './tender-scheduler.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 * as a further provider — same userId-scoping convention as
* TenderTriageService, no forTenant()/RLS (Pitfall 4). * 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 * Seeds itself into the module registry on application startup via
* OnModuleInit lifecycle hook — same pattern as DkvModule. * OnModuleInit lifecycle hook — same pattern as DkvModule.
*/ */
@@ -44,6 +51,7 @@ import { TendersController } from './tenders.controller';
TenderSchedulerService, TenderSchedulerService,
TenderTriageService, TenderTriageService,
TenderSavedSearchService, TenderSavedSearchService,
TenderMatchingService,
], ],
}) })
export class TendersModule implements OnModuleInit { export class TendersModule implements OnModuleInit {