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:
@@ -4,15 +4,22 @@ import { describe, expect, it, vi } from 'vitest';
|
|||||||
import { TenderIngestionService } from './tender-ingestion.service';
|
import { TenderIngestionService } from './tender-ingestion.service';
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* TenderIngestionService.spec — day-cursor gate, SCHEMA-02 change detection,
|
* TenderIngestionService.spec — day-cursor gate, fan-out over active
|
||||||
* and D-05 retention (Plan 10-04, Task 1).
|
* 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
|
* A minimal in-memory fake PrismaService (Maps for tenderSourcePollConfig
|
||||||
* tender) is used, matching this repo's established test convention of
|
* and tender) is used, matching this repo's established test convention of
|
||||||
* hand-rolled prisma-shaped mocks rather than a real DB connection (see
|
* hand-rolled prisma-shaped mocks rather than a real DB connection (see
|
||||||
* ldap.service.spec.ts). doeAdapter/normalizer are also plain fakes — the
|
* ldap.service.spec.ts). The registry/adapter/normalizer/dedup collaborators
|
||||||
* normalizer is stubbed as an identity function so test fixtures can be
|
* are also plain fakes.
|
||||||
* pre-shaped as NormalizedTenderFields directly.
|
*
|
||||||
|
* 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 {
|
function berlinTodayStr(): string {
|
||||||
@@ -31,20 +38,27 @@ function dayDate(offsetDays: number): Date {
|
|||||||
return new Date(`${dayString(offsetDays)}T00:00:00Z`);
|
return new Date(`${dayString(offsetDays)}T00:00:00Z`);
|
||||||
}
|
}
|
||||||
|
|
||||||
function makeFakePrisma() {
|
function makeFakePrisma(configSeeds: Record<string, any> = {}) {
|
||||||
const tenders = new Map<string, any>();
|
const tenders = new Map<string, any>();
|
||||||
const configs = new Map<string, any>();
|
const configs = new Map<string, any>();
|
||||||
configs.set('doe-opendata', {
|
for (const [sourceType, seed] of Object.entries(configSeeds)) {
|
||||||
id: 'cfg1',
|
configs.set(sourceType, {
|
||||||
sourceType: 'doe-opendata',
|
id: `cfg-${sourceType}`,
|
||||||
|
sourceType,
|
||||||
pollIntervalMin: 60,
|
pollIntervalMin: 60,
|
||||||
isActive: true,
|
isActive: true,
|
||||||
lastIngestedDay: null as Date | null,
|
lastIngestedDay: null as Date | null,
|
||||||
|
...seed,
|
||||||
});
|
});
|
||||||
|
}
|
||||||
|
|
||||||
const prisma = {
|
const prisma = {
|
||||||
tenderSourcePollConfig: {
|
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) => {
|
update: vi.fn(async ({ where, data }: any) => {
|
||||||
const existing = configs.get(where.sourceType);
|
const existing = configs.get(where.sourceType);
|
||||||
const updated = { ...existing, ...data };
|
const updated = { ...existing, ...data };
|
||||||
@@ -53,18 +67,6 @@ function makeFakePrisma() {
|
|||||||
}),
|
}),
|
||||||
},
|
},
|
||||||
tender: {
|
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) => {
|
updateMany: vi.fn(async ({ where, data }: any) => {
|
||||||
let count = 0;
|
let count = 0;
|
||||||
for (const [key, row] of tenders) {
|
for (const [key, row] of tenders) {
|
||||||
@@ -110,6 +112,7 @@ const BASE_RECORD = {
|
|||||||
title: 'Test-Ausschreibung',
|
title: 'Test-Ausschreibung',
|
||||||
buyerName: 'Stadt Testhausen',
|
buyerName: 'Stadt Testhausen',
|
||||||
cpvCodes: [] as string[],
|
cpvCodes: [] as string[],
|
||||||
|
cpvDivisions: [] as string[],
|
||||||
region: null as string | null,
|
region: null as string | null,
|
||||||
plz: null as string | null,
|
plz: null as string | null,
|
||||||
bundesland: null as string | null,
|
bundesland: null as string | null,
|
||||||
@@ -122,38 +125,68 @@ const BASE_RECORD = {
|
|||||||
publishedAt: new Date(),
|
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(
|
function makeService(
|
||||||
prisma: any,
|
prisma: any,
|
||||||
doeAdapter: any,
|
registry: any,
|
||||||
normalizer: any,
|
normalizer: any,
|
||||||
matching: any = { matchDelta: vi.fn() },
|
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.
|
// 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;
|
||||||
}
|
}
|
||||||
|
|
||||||
describe('TenderIngestionService.pollDueSources — day-cursor gate (Pitfall A)', () => {
|
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 () => {
|
it('makes NO adapter call when the day-cursor is not strictly before Berlin-today (no-op tick, expected)', async () => {
|
||||||
const prisma = makeFakePrisma();
|
const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-1) } }); // -> next day is today, gated out
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-1); // -> next day is today, gated out
|
|
||||||
const doeAdapter = { fetchTenders: vi.fn() };
|
const doeAdapter = { fetchTenders: vi.fn() };
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn() };
|
const normalizer = { normalize: vi.fn() };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer);
|
const service = makeService(prisma, registry, normalizer);
|
||||||
|
|
||||||
await service.pollDueSources();
|
await service.pollDueSources();
|
||||||
|
|
||||||
expect(doeAdapter.fetchTenders).not.toHaveBeenCalled();
|
expect(doeAdapter.fetchTenders).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
it('does nothing when the singleton doe-opendata config is not active', async () => {
|
it('does nothing when there are no active configs at all', async () => {
|
||||||
const prisma = makeFakePrisma();
|
const prisma = makeFakePrisma({ 'doe-opendata': { isActive: false, lastIngestedDay: dayDate(-3) } });
|
||||||
prisma.__store.configs.get('doe-opendata').isActive = false;
|
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3);
|
|
||||||
const doeAdapter = { fetchTenders: vi.fn() };
|
const doeAdapter = { fetchTenders: vi.fn() };
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn() };
|
const normalizer = { normalize: vi.fn() };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer);
|
const service = makeService(prisma, registry, normalizer);
|
||||||
|
|
||||||
await service.pollDueSources();
|
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 () => {
|
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();
|
const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } });
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3);
|
|
||||||
const record = { ...BASE_RECORD };
|
const record = { ...BASE_RECORD };
|
||||||
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([record]) };
|
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([record]) };
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer);
|
const service = makeService(prisma, registry, normalizer);
|
||||||
|
|
||||||
await service.pollDueSources();
|
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');
|
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 () => {
|
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();
|
const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } });
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3);
|
|
||||||
const day1Record = {
|
const day1Record = {
|
||||||
...BASE_RECORD,
|
...BASE_RECORD,
|
||||||
contentHash: 'hash-1',
|
contentHash: 'hash-1',
|
||||||
@@ -194,8 +226,9 @@ describe('TenderIngestionService.pollDueSources — SCHEMA-02 change detection',
|
|||||||
const doeAdapter = {
|
const doeAdapter = {
|
||||||
fetchTenders: vi.fn().mockResolvedValueOnce([day1Record]).mockResolvedValueOnce([day2Record]),
|
fetchTenders: vi.fn().mockResolvedValueOnce([day1Record]).mockResolvedValueOnce([day2Record]),
|
||||||
};
|
};
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer);
|
const service = makeService(prisma, registry, normalizer);
|
||||||
|
|
||||||
await service.pollDueSources();
|
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)', () => {
|
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 () => {
|
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
|
// Pre-seed one tender that already exists BEFORE this tick — simulates
|
||||||
// a previously-ingested row that reappears in this tick's fetch.
|
// a previously-ingested row that reappears in this tick's fetch.
|
||||||
prisma.__store.tenders.set('existing-notice', {
|
prisma.__store.tenders.set('existing-notice', {
|
||||||
@@ -217,7 +346,6 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12
|
|||||||
contentHash: 'old-hash',
|
contentHash: 'old-hash',
|
||||||
status: 'active',
|
status: 'active',
|
||||||
});
|
});
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-2); // single catch-up day
|
|
||||||
|
|
||||||
const existingRecord = {
|
const existingRecord = {
|
||||||
...BASE_RECORD,
|
...BASE_RECORD,
|
||||||
@@ -230,9 +358,10 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12
|
|||||||
const doeAdapter = {
|
const doeAdapter = {
|
||||||
fetchTenders: vi.fn().mockResolvedValue([existingRecord, newRecord1, newRecord2]),
|
fetchTenders: vi.fn().mockResolvedValue([existingRecord, newRecord1, newRecord2]),
|
||||||
};
|
};
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
||||||
const matching = { matchDelta: vi.fn() };
|
const matching = { matchDelta: vi.fn() };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer, matching);
|
const service = makeService(prisma, registry, normalizer, matching);
|
||||||
|
|
||||||
await service.pollDueSources();
|
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 () => {
|
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', {
|
prisma.__store.tenders.set('existing-notice', {
|
||||||
id: 'id-existing',
|
id: 'id-existing',
|
||||||
dedupKey: 'existing-notice',
|
dedupKey: 'existing-notice',
|
||||||
contentHash: 'old-hash',
|
contentHash: 'old-hash',
|
||||||
status: 'active',
|
status: 'active',
|
||||||
});
|
});
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-2);
|
|
||||||
|
|
||||||
const existingRecord = { ...BASE_RECORD, dedupKey: 'existing-notice' };
|
const existingRecord = { ...BASE_RECORD, dedupKey: 'existing-notice' };
|
||||||
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([existingRecord]) };
|
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([existingRecord]) };
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
const normalizer = { normalize: vi.fn((raw: any) => raw) };
|
||||||
const matching = { matchDelta: vi.fn() };
|
const matching = { matchDelta: vi.fn() };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer, matching);
|
const service = makeService(prisma, registry, normalizer, matching);
|
||||||
|
|
||||||
await service.pollDueSources();
|
await service.pollDueSources();
|
||||||
|
|
||||||
@@ -270,11 +399,11 @@ describe('TenderIngestionService.pollDueSources — delta-only matching (Plan 12
|
|||||||
|
|
||||||
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({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } });
|
||||||
prisma.__store.configs.get('doe-opendata').lastIngestedDay = dayDate(-3);
|
|
||||||
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([]) };
|
const doeAdapter = { fetchTenders: vi.fn().mockResolvedValue([]) };
|
||||||
|
const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter });
|
||||||
const normalizer = { normalize: vi.fn() };
|
const normalizer = { normalize: vi.fn() };
|
||||||
const service = makeService(prisma, doeAdapter, normalizer);
|
const service = makeService(prisma, registry, normalizer);
|
||||||
|
|
||||||
await service.pollDueSources();
|
await service.pollDueSources();
|
||||||
|
|
||||||
@@ -315,7 +444,8 @@ describe('TenderIngestionService.pruneExpiredTenders — D-05 retention', () =>
|
|||||||
deadlineAt: dayDate(1),
|
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();
|
await service.pruneExpiredTenders();
|
||||||
|
|
||||||
|
|||||||
@@ -1,10 +1,11 @@
|
|||||||
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 { SourceRegistry } from './source-registry';
|
||||||
|
import { TenderDedupService } from './tender-dedup.service';
|
||||||
import { TenderMatchingService } from './tender-matching.service';
|
import { TenderMatchingService } from './tender-matching.service';
|
||||||
import { TenderNormalizerService } from './tender-normalizer.service';
|
import { TenderNormalizerService } from './tender-normalizer.service';
|
||||||
|
import type { SourceType } from './tender.types';
|
||||||
|
|
||||||
const DOE_SOURCE_TYPE = 'doe-opendata';
|
|
||||||
const RETENTION_DAYS = 90;
|
const RETENTION_DAYS = 90;
|
||||||
/** Polite delay between successive catch-up day-fetches (RESEARCH Open Question 1). */
|
/** Polite delay between successive catch-up day-fetches (RESEARCH Open Question 1). */
|
||||||
const CATCH_UP_DELAY_MS = 1_500;
|
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
|
* TenderIngestionService — orchestrates the poll tick: fan out over every
|
||||||
* -> fetch -> normalize -> upsert (SCHEMA-02 change detection) -> advance
|
* active `TenderSourcePollConfig` (poll-once-fan-out-many, NOT a single
|
||||||
* cursor -> prune expired (D-05).
|
* 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
|
* Multi-tenant safety (D-03, T-10-09): uses the plain, non-tenant-scoped
|
||||||
* PrismaService injected as-is. Never wraps `Tender`/`TenderSourcePollConfig`
|
* 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
|
* 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
|
* Phase 13, Plan 03 (SCHEMA-03): `pollDueSources` now loads ALL active
|
||||||
* NEW Tender rows created this tick (via an indexed pre-check against the
|
* configs via `findMany` and resolves each config's `sourceType` through
|
||||||
* existing upsert) and calls `TenderMatchingService.matchDelta` with only
|
* `SourceRegistry.get()` — replacing the Phase-10 hardwired single
|
||||||
* those IDs at the end of the tick — this is the delta-only matching
|
* `DoeOpenDataAdapter` injection. Each config's day-cursor-gate/fetch/
|
||||||
* boundary (D-07) that structurally prevents a backfill flood. Rows that
|
* normalize/dedup block runs inside its own try/catch (catch-per-source,
|
||||||
* merely changed (SCHEMA-02 contentHash update) are NOT included.
|
* 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()
|
@Injectable()
|
||||||
export class TenderIngestionService {
|
export class TenderIngestionService {
|
||||||
@@ -63,115 +80,95 @@ export class TenderIngestionService {
|
|||||||
|
|
||||||
constructor(
|
constructor(
|
||||||
private readonly prisma: PrismaService,
|
private readonly prisma: PrismaService,
|
||||||
private readonly doeAdapter: DoeOpenDataAdapter,
|
private readonly registry: SourceRegistry,
|
||||||
private readonly normalizer: TenderNormalizerService,
|
private readonly normalizer: TenderNormalizerService,
|
||||||
private readonly matching: TenderMatchingService,
|
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
|
* TenderSchedulerService's cron tick, on a shared, admin-configurable
|
||||||
* interval, regardless of tenant count. Never throws — catch-and-log per
|
* 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<void> {
|
async pollDueSources(): Promise<void> {
|
||||||
try {
|
try {
|
||||||
const config = await this.prisma.tenderSourcePollConfig.findUnique({
|
const configs = await this.prisma.tenderSourcePollConfig.findMany({
|
||||||
where: { sourceType: DOE_SOURCE_TYPE },
|
where: { isActive: true }, // fan-out, NOT findFirst/findUnique
|
||||||
});
|
});
|
||||||
if (!config?.isActive) return;
|
const activePortalCount = configs.length;
|
||||||
|
const dedupActive = activePortalCount >= 2; // D-05 dedup gate
|
||||||
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', {
|
const berlinToday = new Date().toLocaleDateString('en-CA', {
|
||||||
timeZone: 'Europe/Berlin',
|
timeZone: 'Europe/Berlin',
|
||||||
});
|
});
|
||||||
let fetchedAtLeastOneDay = false;
|
|
||||||
let isFirstFetch = true;
|
|
||||||
// Delta-only matching boundary (D-07, NOTIFY-03): only genuinely NEW
|
// Delta-only matching boundary (D-07, NOTIFY-03): only genuinely NEW
|
||||||
// tender rows created THIS tick are collected here and handed to
|
// tender rows created THIS tick (across ALL sources) are collected
|
||||||
// matchDelta below — never changed/re-seen rows, never historical rows.
|
// here and handed to matchDelta below — never changed/re-seen rows,
|
||||||
|
// never merged-as-additional-source rows, never historical rows.
|
||||||
const newTenderIds: string[] = [];
|
const newTenderIds: string[] = [];
|
||||||
|
let anyDayFetched = false;
|
||||||
|
|
||||||
|
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) {
|
while (cursorDay && cursorDay < berlinToday) {
|
||||||
if (!isFirstFetch) {
|
if (!isFirstFetch) {
|
||||||
// Politeness delay between successive catch-up day-fetches only —
|
// Politeness delay between successive catch-up day-fetches
|
||||||
// never before the first fetch of a tick.
|
// only — never before the first fetch of a source this tick.
|
||||||
await this._delay(this.politeDelayMs);
|
await this._delay(this.politeDelayMs);
|
||||||
}
|
}
|
||||||
isFirstFetch = false;
|
isFirstFetch = false;
|
||||||
|
|
||||||
const rawRecords = await this.doeAdapter.fetchTenders(cursorDay);
|
const rawRecords = await adapter.fetchTenders(cursorDay);
|
||||||
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) {
|
||||||
// Indexed pre-check (dedupKey @unique) — the plain upsert below
|
const { tenderId, created } = await this.dedup.resolve(tender, {
|
||||||
// does not report create-vs-update, so we check existence first
|
dedupActive,
|
||||||
// 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 },
|
|
||||||
});
|
});
|
||||||
|
if (created) newTenderIds.push(tenderId); // genuinely NEW row this tick (D-07)
|
||||||
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({
|
await this.prisma.tenderSourcePollConfig.update({
|
||||||
where: { sourceType: DOE_SOURCE_TYPE },
|
where: { sourceType: config.sourceType },
|
||||||
data: { lastIngestedDay: new Date(`${cursorDay}T00:00:00Z`) },
|
data: { lastIngestedDay: new Date(`${cursorDay}T00:00:00Z`) },
|
||||||
});
|
});
|
||||||
|
|
||||||
fetchedAtLeastOneDay = true;
|
anyDayFetched = true;
|
||||||
cursorDay = addOneDay(cursorDay);
|
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}`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (fetchedAtLeastOneDay) {
|
if (anyDayFetched) {
|
||||||
await this.pruneExpiredTenders();
|
await this.pruneExpiredTenders();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user