import { readFileSync } from 'fs'; import { join } from 'path'; import { describe, expect, it, vi } from 'vitest'; import { TenderIngestionService } from './tender-ingestion.service'; /** * 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 * hand-rolled prisma-shaped mocks rather than a real DB connection (see * 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 { return new Date().toLocaleDateString('en-CA', { timeZone: 'Europe/Berlin' }); } /** Day string, offset from Berlin-today by `offsetDays` (negative = past). */ function dayString(offsetDays: number): string { const base = new Date(`${berlinTodayStr()}T00:00:00Z`); base.setUTCDate(base.getUTCDate() + offsetDays); return base.toISOString().slice(0, 10); } /** UTC-midnight Date for the given day offset — day-granularity, no time-of-day drift. */ function dayDate(offsetDays: number): Date { return new Date(`${dayString(offsetDays)}T00:00:00Z`); } function makeFakePrisma(configSeeds: Record = {}) { const tenders = new Map(); const configs = new Map(); 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: { 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 }; configs.set(where.sourceType, updated); return updated; }), }, tender: { updateMany: vi.fn(async ({ where, data }: any) => { let count = 0; for (const [key, row] of tenders) { if ( row.status === where.status && row.deadlineAt instanceof Date && where.deadlineAt?.lt && row.deadlineAt < where.deadlineAt.lt ) { tenders.set(key, { ...row, ...data }); count++; } } return { count }; }), deleteMany: vi.fn(async ({ where }: any) => { let count = 0; for (const [key, row] of tenders) { if ( row.status === where.status && row.deadlineAt instanceof Date && where.deadlineAt?.lt && row.deadlineAt < where.deadlineAt.lt ) { tenders.delete(key); count++; } } return { count }; }), }, __store: { tenders, configs }, }; return prisma; } const BASE_RECORD = { sourcePortal: 'doe-opendata', sourceNoticeId: 'notice-1', ocid: 'ocds-1', dedupKey: 'ocds-1', 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, deadlineAt: null as Date | null, estimatedValue: null as number | null, procedureType: null as string | null, status: 'active', sourceUrl: null as string | null, contentHash: 'hash-1', 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, registry: any, normalizer: any, matching: any = { matchDelta: vi.fn() }, dedup: any = makeFakeDedup(prisma.__store.tenders), ) { 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 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, registry, normalizer); await service.pollDueSources(); expect(doeAdapter.fetchTenders).not.toHaveBeenCalled(); }); 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, registry, normalizer); await service.pollDueSources(); expect(doeAdapter.fetchTenders).not.toHaveBeenCalled(); }); }); 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({ '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, registry, normalizer); await service.pollDueSources(); // Catch-up loop covers 2 days (today-2, today-1) — same record both days. expect(doeAdapter.fetchTenders).toHaveBeenCalledTimes(2); expect(prisma.__store.tenders.size).toBe(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) — SCHEMA-02 preserved through resolve()', async () => { const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-3) } }); const day1Record = { ...BASE_RECORD, contentHash: 'hash-1', deadlineAt: new Date('2026-08-01T10:00:00Z'), }; const day2Record = { ...BASE_RECORD, contentHash: 'hash-2', deadlineAt: new Date('2026-08-15T10:00:00Z'), }; 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, registry, normalizer); await service.pollDueSources(); expect(prisma.__store.tenders.size).toBe(1); // updated in place, no duplicate row const row = prisma.__store.tenders.get('ocds-1'); expect(row.contentHash).toBe('hash-2'); expect(row.deadlineAt).toEqual(new Date('2026-08-15T10:00:00Z')); }); }); 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 — pollGranularity 'tick' gate (D-15, Phase 14 Plan 02)", () => { it("fetches a pollGranularity='tick' source on every tick, while a 'day' source with today's next-day gated-out lastIngestedDay is skipped in the SAME tick", async () => { const prisma = makeFakePrisma({ 'doe-opendata': { lastIngestedDay: dayDate(-1) }, // -> next day is today, gated out ('day' path unchanged) rss: { pollGranularity: 'tick', lastIngestedDay: null }, // 'tick' — no day-cursor gate at all }); const doeAdapter = { fetchTenders: vi.fn() }; const rssAdapter = { fetchTenders: vi.fn().mockResolvedValue([ { ...BASE_RECORD, sourcePortal: 'rss', sourceNoticeId: 'rss-1', ocid: null, dedupKey: 'rss:rss-1' }, ]), }; const registry = makeFakeRegistry({ 'doe-opendata': doeAdapter, rss: rssAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); expect(doeAdapter.fetchTenders).not.toHaveBeenCalled(); // 'day' gate unchanged expect(rssAdapter.fetchTenders).toHaveBeenCalledTimes(1); // 'tick' — fetched regardless expect(prisma.__store.tenders.size).toBe(1); }); it("does NOT advance or read lastIngestedDay for a 'tick' source — it is fetched again next tick even with lastIngestedDay already set to today", async () => { const prisma = makeFakePrisma({ rss: { pollGranularity: 'tick', lastIngestedDay: dayDate(0) }, // today — would gate a 'day' source out entirely }); const rssAdapter = { fetchTenders: vi.fn().mockResolvedValue([]) }; const registry = makeFakeRegistry({ rss: rssAdapter }); const normalizer = { normalize: vi.fn() }; const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); await service.pollDueSources(); expect(rssAdapter.fetchTenders).toHaveBeenCalledTimes(2); // lastIngestedDay is never touched by the 'tick' path. expect(prisma.__store.configs.get('rss').lastIngestedDay).toEqual(dayDate(0)); }); it("passes berlinToday (not a day-cursor loop) as the single argument to a 'tick' adapter's fetchTenders", async () => { const prisma = makeFakePrisma({ rss: { pollGranularity: 'tick' } }); const rssAdapter = { fetchTenders: vi.fn().mockResolvedValue([]) }; const registry = makeFakeRegistry({ rss: rssAdapter }); const normalizer = { normalize: vi.fn() }; const service = makeService(prisma, registry, normalizer); await service.pollDueSources(); expect(rssAdapter.fetchTenders).toHaveBeenCalledWith(berlinTodayStr()); }); }); 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({ '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', { id: 'id-existing', dedupKey: 'existing-notice', contentHash: 'old-hash', status: 'active', }); 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 registry = makeFakeRegistry({ 'doe-opendata': doeAdapter }); const normalizer = { normalize: vi.fn((raw: any) => raw) }; const matching = { matchDelta: vi.fn() }; const service = makeService(prisma, registry, 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({ 'doe-opendata': { lastIngestedDay: dayDate(-2) } }); prisma.__store.tenders.set('existing-notice', { id: 'id-existing', dedupKey: 'existing-notice', contentHash: 'old-hash', status: 'active', }); 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, registry, 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({ '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, registry, normalizer); await service.pollDueSources(); const calledDays = doeAdapter.fetchTenders.mock.calls.map((c: any[]) => c[0]); expect(calledDays).toEqual([dayString(-2), dayString(-1)]); const finalConfig = prisma.__store.configs.get('doe-opendata'); expect(finalConfig.lastIngestedDay.toISOString().slice(0, 10)).toBe(dayString(-1)); }); }); describe('TenderIngestionService.pruneExpiredTenders — D-05 retention', () => { it('marks past-deadline active rows expired, deletes expired rows older than 90 days, and never touches null-deadline rows', async () => { const prisma = makeFakePrisma(); prisma.__store.tenders.set('t-active-past-deadline', { dedupKey: 't-active-past-deadline', status: 'active', deadlineAt: dayDate(-1), }); prisma.__store.tenders.set('t-expired-old', { dedupKey: 't-expired-old', status: 'expired', deadlineAt: dayDate(-91), }); prisma.__store.tenders.set('t-expired-recent', { dedupKey: 't-expired-recent', status: 'expired', deadlineAt: dayDate(-30), }); prisma.__store.tenders.set('t-null-deadline', { dedupKey: 't-null-deadline', status: 'active', deadlineAt: null, }); prisma.__store.tenders.set('t-active-future', { dedupKey: 't-active-future', status: 'active', deadlineAt: dayDate(1), }); const registry = makeFakeRegistry({}); const service = makeService(prisma, registry, { normalize: vi.fn() }); await service.pruneExpiredTenders(); expect(prisma.__store.tenders.get('t-active-past-deadline').status).toBe('expired'); expect(prisma.__store.tenders.has('t-expired-old')).toBe(false); expect(prisma.__store.tenders.get('t-expired-recent').status).toBe('expired'); expect(prisma.__store.tenders.get('t-null-deadline').status).toBe('active'); expect(prisma.__store.tenders.get('t-active-future').status).toBe('active'); }); }); describe('TenderIngestionService — multi-tenant safety invariant (D-03)', () => { it('never calls forTenant() — uses the plain global PrismaService on Tender/TenderSourcePollConfig', () => { const source = readFileSync(join(__dirname, 'tender-ingestion.service.ts'), 'utf8'); expect(source).not.toMatch(/forTenant/); }); });