feat(10-03): implement DoeOpenDataAdapter — fetch, extract, D-02 filter
- Native fetch + AbortController 15s timeout (icon-discovery idiom, no
axios); URL host hardcoded, only the internally-computed dayCursor is
interpolated (T-10-06)
- HTTP 400 treated as an expected no-op (pubDay today/future) -> []
- adm-zip extraction with a pre-extraction decompression-bomb ceiling
check (sum entry.header.size vs ~50MB) before any entry buffer is read
(T-10-07); entries are never written to disk
- D-02 open-tender filter: positive tag.includes('tender') match only —
award/planning/untagged-with-awards excluded (Pitfall C)
- Register DoeOpenDataAdapter in TendersModule.providers
- All doe-opendata.adapter.spec.ts tests green (6/6)
This commit is contained in:
@@ -0,0 +1,184 @@
|
||||
import { Injectable, Logger } from '@nestjs/common';
|
||||
import AdmZip from 'adm-zip';
|
||||
import { XMLParser } from 'fast-xml-parser';
|
||||
import type { RawTenderRecord, SourceType } from '../tender.types';
|
||||
import type { TenderSourceAdapter } from './tender-source-adapter.interface';
|
||||
|
||||
/**
|
||||
* DoeOpenDataAdapter — fetches, extracts, and parses the DÖE OpenData
|
||||
* day-batch export (eForms-DE XML + OCDS JSON) into filtered, open-tender
|
||||
* RawTenderRecord[] (INGEST-01).
|
||||
*
|
||||
* Security notes (threat model T-10-06/T-10-07):
|
||||
* - URL host is a hardcoded constant; only the internally-computed
|
||||
* `dayCursor` is interpolated — never user/admin input (T-10-06).
|
||||
* - Native fetch + AbortController 15s timeout (icon-discovery.service.ts
|
||||
* idiom) — no axios anywhere in this codebase.
|
||||
* - HTTP 400 is an EXPECTED no-op signal (pubDay is today/future per the
|
||||
* live-verified DÖE boundary rule) — returns [], never throws. Any other
|
||||
* `!res.ok` throws.
|
||||
* - Zip entries are summed via `entry.header.size` (uncompressed size
|
||||
* claimed in the archive's own header) and rejected above a ~50MB
|
||||
* ceiling BEFORE any entry buffer is read — decompression-bomb guard
|
||||
* (T-10-07). Entries are read fully in-memory and NEVER written to disk,
|
||||
* so zip-slip path traversal is structurally absent, not just mitigated.
|
||||
*/
|
||||
const DOE_BASE_URL = 'https://oeffentlichevergabe.de';
|
||||
const DOE_FETCH_TIMEOUT_MS = 15_000;
|
||||
const ZIP_SIZE_CEILING_BYTES = 50 * 1024 * 1024;
|
||||
|
||||
type DoeFormat = 'eforms.zip' | 'ocds.zip';
|
||||
|
||||
interface OcdsRelease {
|
||||
id: string;
|
||||
ocid?: string;
|
||||
tag?: string[];
|
||||
awards?: unknown[];
|
||||
contracts?: unknown[];
|
||||
}
|
||||
|
||||
interface OcdsDocument {
|
||||
uri?: string;
|
||||
publishedDate?: string;
|
||||
releases?: OcdsRelease[];
|
||||
}
|
||||
|
||||
/**
|
||||
* D-02 open-tender filter (RESEARCH.md Pattern 2 / Pitfall C): positive
|
||||
* tag-presence match only. Missing/other tags (award, planning, or no tag
|
||||
* at all — even with populated awards/contracts) are conservatively
|
||||
* excluded, never treated as "open" by default.
|
||||
*/
|
||||
export function isOpenTenderNotice(tag: string[] | undefined): boolean {
|
||||
return Array.isArray(tag) && tag.includes('tender');
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class DoeOpenDataAdapter implements TenderSourceAdapter {
|
||||
readonly sourceType: SourceType = 'doe-opendata';
|
||||
|
||||
private readonly logger = new Logger(DoeOpenDataAdapter.name);
|
||||
private readonly xmlParser = new XMLParser({
|
||||
removeNSPrefix: true,
|
||||
ignoreAttributes: false,
|
||||
attributeNamePrefix: '@_',
|
||||
});
|
||||
|
||||
async fetchTenders(dayCursor: string): Promise<RawTenderRecord[]> {
|
||||
const eformsBuffer = await this.fetchDoeZip(dayCursor, 'eforms.zip');
|
||||
if (eformsBuffer === null) return []; // 400 no-op — nothing to fetch yet this cycle
|
||||
|
||||
const eformsZip = new AdmZip(eformsBuffer);
|
||||
this.assertWithinSizeCeiling(eformsZip);
|
||||
const eformsEntries = eformsZip.getEntries();
|
||||
|
||||
const ocdsBuffer = await this.fetchDoeZip(dayCursor, 'ocds.zip');
|
||||
if (ocdsBuffer === null) return [];
|
||||
|
||||
const ocdsZip = new AdmZip(ocdsBuffer);
|
||||
this.assertWithinSizeCeiling(ocdsZip);
|
||||
const ocdsEntries = ocdsZip.getEntries();
|
||||
|
||||
const eformsByBasename = new Map<string, unknown>();
|
||||
for (const entry of eformsEntries) {
|
||||
const basename = entry.entryName.replace(/\.xml$/i, '');
|
||||
try {
|
||||
eformsByBasename.set(
|
||||
basename,
|
||||
this.xmlParser.parse(entry.getData().toString('utf8')),
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.warn(
|
||||
`Failed to parse eForms-DE XML entry ${entry.entryName}: ${(error as Error).message}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const records: RawTenderRecord[] = [];
|
||||
const fetchedAt = new Date();
|
||||
|
||||
for (const entry of ocdsEntries) {
|
||||
const basename = entry.entryName.replace(/\.json$/i, '');
|
||||
|
||||
let ocdsDocument: OcdsDocument;
|
||||
try {
|
||||
ocdsDocument = JSON.parse(entry.getData().toString('utf8'));
|
||||
} catch (error) {
|
||||
this.logger.warn(
|
||||
`Failed to parse OCDS JSON entry ${entry.entryName}: ${(error as Error).message}`,
|
||||
);
|
||||
continue;
|
||||
}
|
||||
|
||||
const release = ocdsDocument.releases?.[0];
|
||||
if (!release) continue;
|
||||
|
||||
// D-02: whole-Germany ingest, no region/CPV pre-filtering (D-03) —
|
||||
// the ONLY filter applied here is the open-tender tag check.
|
||||
if (!isOpenTenderNotice(release.tag)) continue;
|
||||
|
||||
records.push({
|
||||
sourceType: this.sourceType,
|
||||
sourcePortal: 'doe-opendata',
|
||||
sourceNoticeId: release.id,
|
||||
ocid: release.ocid,
|
||||
sourceUrl: ocdsDocument.uri,
|
||||
fetchedAt,
|
||||
publishedAt: ocdsDocument.publishedDate
|
||||
? new Date(ocdsDocument.publishedDate)
|
||||
: null,
|
||||
eformsPayload: eformsByBasename.get(basename) ?? null,
|
||||
ocdsPayload: release,
|
||||
});
|
||||
}
|
||||
|
||||
return records;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sum each entry's declared uncompressed size (`header.size`) BEFORE any
|
||||
* entry buffer is read/decompressed, and reject if the total exceeds the
|
||||
* decompression-bomb ceiling (T-10-07). `adm-zip`'s `getEntries()` reads
|
||||
* only the central directory — this check runs before `getData()` is
|
||||
* ever called on any entry.
|
||||
*/
|
||||
private assertWithinSizeCeiling(zip: AdmZip): void {
|
||||
const totalUncompressed = zip
|
||||
.getEntries()
|
||||
.reduce((sum, entry) => sum + entry.header.size, 0);
|
||||
|
||||
if (totalUncompressed > ZIP_SIZE_CEILING_BYTES) {
|
||||
throw new Error(
|
||||
`DÖE export exceeds decompression-bomb ceiling (${totalUncompressed} bytes > ${ZIP_SIZE_CEILING_BYTES} bytes)`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private async fetchDoeZip(
|
||||
dayCursor: string,
|
||||
format: DoeFormat,
|
||||
): Promise<Buffer | null> {
|
||||
const controller = new AbortController();
|
||||
const timeout = setTimeout(
|
||||
() => controller.abort(),
|
||||
DOE_FETCH_TIMEOUT_MS,
|
||||
);
|
||||
|
||||
try {
|
||||
const res = await fetch(
|
||||
`${DOE_BASE_URL}/api/notice-exports?pubDay=${dayCursor}&format=${format}`,
|
||||
{ signal: controller.signal },
|
||||
);
|
||||
|
||||
// Expected no-op: pubDay is today/future ("must lie in the past").
|
||||
if (res.status === 400) return null;
|
||||
if (!res.ok) {
|
||||
throw new Error(`DÖE fetch failed (${format}): ${res.status}`);
|
||||
}
|
||||
|
||||
return Buffer.from(await res.arrayBuffer());
|
||||
} finally {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2,15 +2,17 @@ import { Logger, Module, OnModuleInit } from '@nestjs/common';
|
||||
import { ModuleRegistryModule } from '../module-registry/module-registry.module';
|
||||
import { ModuleRegistryService } from '../module-registry/module-registry.service';
|
||||
import { PrismaService } from '../prisma/prisma.service';
|
||||
import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter';
|
||||
import { seedTendersModule } from './tenders.seed';
|
||||
|
||||
/**
|
||||
* NestJS module for the Ausschreibungs-Radar feature.
|
||||
*
|
||||
* Wave 2 (this plan) only registers the module in the marketplace and
|
||||
* seeds the singleton DÖE poll config so the shared poll is
|
||||
* admin-drivable. Ingestion/adapter/scheduler/controller wiring is
|
||||
* added in Plans 03-05.
|
||||
* Wave 2 registered the module in the marketplace and seeded the
|
||||
* singleton DÖE poll config so the shared poll is admin-drivable. This
|
||||
* plan (03) adds the DÖE source adapter + normalizer (parse+map core of
|
||||
* INGEST-01/SCHEMA-01). Ingestion orchestration/scheduler/controller
|
||||
* wiring is added in Plans 04-05.
|
||||
*
|
||||
* PrismaModule is global (no explicit import needed).
|
||||
*
|
||||
@@ -20,7 +22,7 @@ import { seedTendersModule } from './tenders.seed';
|
||||
@Module({
|
||||
imports: [ModuleRegistryModule],
|
||||
controllers: [],
|
||||
providers: [],
|
||||
providers: [DoeOpenDataAdapter],
|
||||
})
|
||||
export class TendersModule implements OnModuleInit {
|
||||
private readonly logger = new Logger(TendersModule.name);
|
||||
|
||||
Reference in New Issue
Block a user