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 { ModuleRegistryModule } from '../module-registry/module-registry.module';
|
||||||
import { ModuleRegistryService } from '../module-registry/module-registry.service';
|
import { ModuleRegistryService } from '../module-registry/module-registry.service';
|
||||||
import { PrismaService } from '../prisma/prisma.service';
|
import { PrismaService } from '../prisma/prisma.service';
|
||||||
|
import { DoeOpenDataAdapter } from './adapters/doe-opendata.adapter';
|
||||||
import { seedTendersModule } from './tenders.seed';
|
import { seedTendersModule } from './tenders.seed';
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* NestJS module for the Ausschreibungs-Radar feature.
|
* NestJS module for the Ausschreibungs-Radar feature.
|
||||||
*
|
*
|
||||||
* Wave 2 (this plan) only registers the module in the marketplace and
|
* Wave 2 registered the module in the marketplace and seeded the
|
||||||
* seeds the singleton DÖE poll config so the shared poll is
|
* singleton DÖE poll config so the shared poll is admin-drivable. This
|
||||||
* admin-drivable. Ingestion/adapter/scheduler/controller wiring is
|
* plan (03) adds the DÖE source adapter + normalizer (parse+map core of
|
||||||
* added in Plans 03-05.
|
* INGEST-01/SCHEMA-01). Ingestion orchestration/scheduler/controller
|
||||||
|
* wiring is added in Plans 04-05.
|
||||||
*
|
*
|
||||||
* PrismaModule is global (no explicit import needed).
|
* PrismaModule is global (no explicit import needed).
|
||||||
*
|
*
|
||||||
@@ -20,7 +22,7 @@ import { seedTendersModule } from './tenders.seed';
|
|||||||
@Module({
|
@Module({
|
||||||
imports: [ModuleRegistryModule],
|
imports: [ModuleRegistryModule],
|
||||||
controllers: [],
|
controllers: [],
|
||||||
providers: [],
|
providers: [DoeOpenDataAdapter],
|
||||||
})
|
})
|
||||||
export class TendersModule implements OnModuleInit {
|
export class TendersModule implements OnModuleInit {
|
||||||
private readonly logger = new Logger(TendersModule.name);
|
private readonly logger = new Logger(TendersModule.name);
|
||||||
|
|||||||
Reference in New Issue
Block a user