feat(14-02): add TenderRssFeedSource CRUD, tick poll gate, and wire RssAdapter

Global admin-managed RSS feed list (TenderRssFeedSource, D-08/D-14) with
a save-time hostname/SSRF guard (TenderRssFeedSourceService) — RSS feed
URLs are runtime admin input, so the code-level SourceRegistry denylist
gate does not cover them; a separate check rejects DENYLISTED_PORTALS
hostnames, non-http(s) schemes, and private/loopback hosts.

Adds TenderSourcePollConfig.pollGranularity ('day' | 'tick', D-15):
pollDueSources() branches per source — 'day' sources keep the existing
lastIngestedDay gate byte-unchanged, 'tick' sources (rss) fetch on every
active scheduler tick regardless of lastIngestedDay, since the day-cursor
gate was built for a genuine daily batch-export API and would otherwise
silently cap RSS to one fetch per calendar day.

Wires RssAdapter.fetchTenders() to fan out over active feed rows (native
fetch + AbortController 15s + response-size ceiling, catch-per-feed),
registers it in tenders.module.ts, and seeds the 'rss' poll config
active with pollGranularity='tick' plus a default-active service.bund.de
feed row (subreport-elvis has no single canonical URL — zero rows seeded,
admin adds relevant municipality feeds).

Migration applied locally per project convention (host -> container IP).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-23 13:24:36 +02:00
parent 3a96cbbbe6
commit e812738c3a
10 changed files with 809 additions and 21 deletions
+139 -13
View File
@@ -1,6 +1,6 @@
import { readFileSync } from 'fs';
import { join } from 'path';
import { describe, expect, it } from 'vitest';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { RssAdapter } from './rss.adapter';
/**
@@ -20,16 +20,46 @@ const SUBREPORT_ELVIS_FIXTURE = readFileSync(
'utf8',
);
/** Minimal prisma-shaped fake — only `tenderRssFeedSource.findMany` is used by RssAdapter. */
function makeFakePrisma(feeds: any[] = []) {
return {
tenderRssFeedSource: {
findMany: vi.fn(async () => feeds),
},
} as any;
}
/** Stubs global `fetch` to resolve with `xml` for every call (fan-out tests). */
function stubFetchResolvingXml(xmlByUrl: Record<string, string>): void {
vi.stubGlobal(
'fetch',
vi.fn(async (url: string) => {
const xml = xmlByUrl[url];
if (xml === undefined) {
return { ok: false, status: 404 } as Response;
}
const bytes = Buffer.from(xml, 'utf8');
return {
ok: true,
status: 200,
headers: new Headers({ 'content-length': String(bytes.byteLength) }),
arrayBuffer: async () =>
bytes.buffer.slice(bytes.byteOffset, bytes.byteOffset + bytes.byteLength),
} as unknown as Response;
}),
);
}
describe('RssAdapter', () => {
it('declares sourceType rss and the symbolic rss portal', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
expect(adapter.sourceType).toBe('rss');
expect(adapter.portals).toEqual(['rss']);
});
describe('parseFeed (pure, fixture-driven)', () => {
it('parses service.bund.de items: sourceType/sourcePortal set, pubDate present, numeric HTML entities decoded in title', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const records = adapter.parseFeed(SERVICE_BUND_FIXTURE, 'service-bund');
@@ -54,7 +84,7 @@ describe('RssAdapter', () => {
});
it('uses the item guid as sourceNoticeId for service.bund.de (plain-text guid, no isPermaLink attribute)', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const records = adapter.parseFeed(SERVICE_BUND_FIXTURE, 'service-bund');
const ids = records.map((r) => r.sourceNoticeId);
@@ -63,7 +93,7 @@ describe('RssAdapter', () => {
});
it('parses subreport-elvis items: publishedAt null (no item pubDate), title from CDATA, guid isPermaLink=false extracted as plain text', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const records = adapter.parseFeed(
SUBREPORT_ELVIS_FIXTURE,
@@ -92,7 +122,7 @@ describe('RssAdapter', () => {
});
it('never carries description HTML into the record (V5 — no stored-XSS surface)', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const records = adapter.parseFeed(
SUBREPORT_ELVIS_FIXTURE,
'subreport-neuss',
@@ -106,7 +136,7 @@ describe('RssAdapter', () => {
});
it('bare-minimum bag: buyerName/procedureType/deadlineAt always null (baseline mapping only, D-04/D-05)', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const records = adapter.parseFeed(SERVICE_BUND_FIXTURE, 'service-bund');
for (const record of records) {
@@ -124,26 +154,26 @@ describe('RssAdapter', () => {
});
it('returns [] for empty XML instead of throwing', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
expect(adapter.parseFeed('', 'empty-feed')).toEqual([]);
});
it('returns [] for malformed XML instead of throwing', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
expect(
adapter.parseFeed('<rss><channel><item><title>', 'broken-feed'),
).toEqual([]);
});
it('returns [] for well-formed XML with no <item> at all', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const xml =
'<?xml version="1.0"?><rss version="2.0"><channel><title>Empty</title></channel></rss>';
expect(adapter.parseFeed(xml, 'no-items-feed')).toEqual([]);
});
it('skips a single item missing <link> without aborting the rest', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const xml = `<?xml version="1.0"?>
<rss version="2.0"><channel>
<item><title>Ohne Link</title><guid>no-link-guid</guid></item>
@@ -157,7 +187,7 @@ describe('RssAdapter', () => {
});
it('falls back to sha256(link) for sourceNoticeId when guid is absent', () => {
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
const xml = `<?xml version="1.0"?>
<rss version="2.0"><channel>
<item><title>No Guid</title><link>https://example.invalid/no-guid</link></item>
@@ -177,11 +207,107 @@ describe('RssAdapter', () => {
// only proves RssAdapter's OWN output shape (parseFeed), not the
// normalizer dispatch itself (kept in the normalizer's own spec file
// to avoid duplicating TenderNormalizerService test infrastructure).
const adapter = new RssAdapter();
const adapter = new RssAdapter(makeFakePrisma());
expect(adapter.sourceType).toBe('rss');
});
});
describe('fetchTenders (internal fan-out over active TenderRssFeedSource rows, Plan 14-02 Task 2)', () => {
afterEach(() => {
vi.unstubAllGlobals();
});
it('fetches and parses every active feed, fanning out over findMany({isActive:true})', async () => {
const feeds = [
{ id: 'f1', url: 'https://service-bund.invalid/rss.xml', label: 'service-bund', isActive: true },
{ id: 'f2', url: 'https://subreport.invalid/rss.xml', label: 'subreport-neuss', isActive: true },
];
const prisma = makeFakePrisma(feeds);
stubFetchResolvingXml({
'https://service-bund.invalid/rss.xml': SERVICE_BUND_FIXTURE,
'https://subreport.invalid/rss.xml': SUBREPORT_ELVIS_FIXTURE,
});
const adapter = new RssAdapter(prisma);
const records = await adapter.fetchTenders('2026-07-23');
expect(prisma.tenderRssFeedSource.findMany).toHaveBeenCalledWith({
where: { isActive: true },
});
const portals = new Set(records.map((r) => r.sourcePortal));
expect(portals).toEqual(new Set(['service-bund', 'subreport-neuss']));
expect(records.length).toBeGreaterThan(1);
});
it('catch-per-feed: a feed that fails to fetch is skipped, the other feed still resolves (D-01)', async () => {
const feeds = [
{ id: 'f1', url: 'https://broken.invalid/rss.xml', label: 'broken', isActive: true },
{ id: 'f2', url: 'https://ok.invalid/rss.xml', label: 'ok-feed', isActive: true },
];
const prisma = makeFakePrisma(feeds);
stubFetchResolvingXml({
'https://ok.invalid/rss.xml': SERVICE_BUND_FIXTURE,
// 'broken.invalid' deliberately absent -> stub resolves 404 -> throws
});
const adapter = new RssAdapter(prisma);
const records = await adapter.fetchTenders('2026-07-23');
expect(records.every((r) => r.sourcePortal === 'ok-feed')).toBe(true);
expect(records.length).toBeGreaterThanOrEqual(1);
});
it('returns [] when there are no active feeds', async () => {
const prisma = makeFakePrisma([]);
const adapter = new RssAdapter(prisma);
const records = await adapter.fetchTenders('2026-07-23');
expect(records).toEqual([]);
});
it('returns [] without throwing when a feed fetch throws (network error)', async () => {
const feeds = [
{ id: 'f1', url: 'https://unreachable.invalid/rss.xml', label: 'unreachable', isActive: true },
];
const prisma = makeFakePrisma(feeds);
vi.stubGlobal(
'fetch',
vi.fn(async () => {
throw new Error('network unreachable');
}),
);
const adapter = new RssAdapter(prisma);
const records = await adapter.fetchTenders('2026-07-23');
expect(records).toEqual([]);
});
it('rejects a feed response exceeding the size ceiling (DoS mitigation, T-14-02-02)', async () => {
const feeds = [
{ id: 'f1', url: 'https://huge.invalid/rss.xml', label: 'huge', isActive: true },
];
const prisma = makeFakePrisma(feeds);
vi.stubGlobal(
'fetch',
vi.fn(async () => {
return {
ok: true,
status: 200,
headers: new Headers({ 'content-length': String(11 * 1024 * 1024) }),
arrayBuffer: async () => new ArrayBuffer(0),
} as unknown as Response;
}),
);
const adapter = new RssAdapter(prisma);
const records = await adapter.fetchTenders('2026-07-23');
expect(records).toEqual([]); // oversized feed skipped, catch-per-feed (D-01)
});
});
it('never imports or uses axios (native fetch is the sole HTTP client convention)', () => {
const source = readFileSync(join(__dirname, 'rss.adapter.ts'), 'utf8');
expect(source).not.toMatch(/from ['"]axios['"]/);
+84 -5
View File
@@ -1,6 +1,7 @@
import { Injectable, Logger } from '@nestjs/common';
import { createHash } from 'crypto';
import { XMLParser } from 'fast-xml-parser';
import { PrismaService } from '../../prisma/prisma.service';
import type { RawTenderRecord, SourceType } from '../tender.types';
import type { TenderSourceAdapter } from './tender-source-adapter.interface';
@@ -39,7 +40,21 @@ import type { TenderSourceAdapter } from './tender-source-adapter.interface';
* read by this adapter, so no raw HTML fragment is ever carried into a
* RawTenderRecord (V5 — no stored-XSS surface, same text-only discipline
* as NetServerAdapter/CosinexAdapter).
*
* `fetchTenders()` (Plan 14-02 Task 2) is an internal-fan-out adapter
* (14-RESEARCH.md Pattern 1, same shape as `NetServerAdapter`'s
* multi-portal loop): reads every `isActive` `TenderRssFeedSource` row and
* fetches+parses each feed independently, catch-per-feed (D-01 discipline)
* — one broken/unreachable feed never blocks the others in the same tick.
* `_dayCursor` is accepted for interface conformance but NOT used as a
* filter — RSS has no day-batch concept; it is polled every scheduler
* tick via `pollGranularity: 'tick'` (D-15, wired in
* `tender-ingestion.service.ts`), not gated by the day-cursor at all.
*/
const RSS_FETCH_TIMEOUT_MS = 15_000;
/** RSS feeds are normally small; an admin-supplied URL is less trusted than a hardcoded portal (RESEARCH.md Security Domain, DoS mitigation). */
const RSS_RESPONSE_SIZE_CEILING_BYTES = 10 * 1024 * 1024;
@Injectable()
export class RssAdapter implements TenderSourceAdapter {
readonly sourceType: SourceType = 'rss';
@@ -53,14 +68,78 @@ export class RssAdapter implements TenderSourceAdapter {
attributeNamePrefix: '@_',
});
constructor(private readonly prisma: PrismaService) {}
/**
* Placeholder — wired to the real `TenderRssFeedSource` internal fan-out
* (native fetch + AbortController per admin-added feed URL) in Plan
* 14-02 Task 2. Kept here only so the class satisfies
* `TenderSourceAdapter` for this task's fixture-driven `parseFeed` proof.
* Internal fan-out over every active admin-managed feed (D-14/D-08,
* global — no tenantId). Catch-per-feed (D-01): a single feed that fails
* to fetch/parse is skipped (logger.warn), never aborting the rest.
*/
async fetchTenders(_dayCursor: string): Promise<RawTenderRecord[]> {
return [];
const feeds = await this.prisma.tenderRssFeedSource.findMany({
where: { isActive: true },
});
const records: RawTenderRecord[] = [];
for (const feed of feeds) {
try {
const xml = await this.fetchFeedXml(feed.url);
records.push(...this.parseFeed(xml, feed.label));
} catch (error) {
this.logger.warn(
`RSS feed '${feed.label}' (${feed.url}) fetch failed, skipping this feed for this tick: ${(error as Error).message}`,
);
}
}
return records;
}
/**
* Native fetch + AbortController 15s timeout (project-wide convention,
* DoeOpenDataAdapter/NetServerAdapter/CosinexAdapter idiom) plus a
* response-size ceiling (T-14-02-02, DoS mitigation) — checked both via
* the `Content-Length` header (fast-path, may be absent/wrong) and the
* actually-received byte count (authoritative).
*/
private async fetchFeedXml(url: string): Promise<string> {
const controller = new AbortController();
const timeout = setTimeout(
() => controller.abort(),
RSS_FETCH_TIMEOUT_MS,
);
try {
const res = await fetch(url, {
signal: controller.signal,
redirect: 'follow',
});
if (!res.ok) {
throw new Error(`RSS feed fetch failed: HTTP ${res.status}`);
}
const contentLength = res.headers.get('content-length');
if (
contentLength &&
Number(contentLength) > RSS_RESPONSE_SIZE_CEILING_BYTES
) {
throw new Error(
`RSS feed exceeds size ceiling (Content-Length: ${contentLength} bytes)`,
);
}
const buffer = await res.arrayBuffer();
if (buffer.byteLength > RSS_RESPONSE_SIZE_CEILING_BYTES) {
throw new Error(
`RSS feed exceeds size ceiling (${buffer.byteLength} bytes received)`,
);
}
return new TextDecoder('utf-8').decode(buffer);
} finally {
clearTimeout(timeout);
}
}
/**