feat(13-03): pollDueSources fan-out over all active sources (SCHEMA-03)

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) <noreply@anthropic.com>
This commit is contained in:
2026-07-23 08:50:57 +02:00
parent 1cd2fcfa01
commit d453dbbd1f
2 changed files with 283 additions and 156 deletions
@@ -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<string, any> = {}) {
const tenders = new Map<string, any>();
const configs = new Map<string, any>();
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<string, any>) {
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<string, any>) {
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();