/** * Livia Research-PoC — der eine serverseitige Endpoint (§20). * * POST /api/livia/research/refresh * * Ruft die drei Quellen ab, normalisiert, beurteilt einmal per LLM und gibt * die relevanten Leads zurück. Kein Scheduler, keine Queue, keine Datenbank — * der Lauf beginnt und endet mit dieser Anfrage. * * Läuft als serverlose Funktion auf derselben Vercel-Instanz, die schon die * Anwendung ausliefert. Das ist der Grund, warum für den PoC gar keine neue * Infrastruktur nötig war: der Hoster kann bereits Serverfunktionen, es hat * bisher nur niemand eine gebraucht. */ import type { IncomingMessage, ServerResponse } from 'node:http' import { SOURCES, fetchZhkResearch, fetchGreaterZurichResearch, fetchZefixResearch, fetchGenericSource, } from '../../_lib/researchSources.js' import type { ResearchItem } from '../../_lib/researchSources.js' import { MODEL, analyzeResearchItems } from '../../_lib/researchAnalysis.js' import type { ResearchRefreshResult, ResearchSourceResult, TemporalClass, } from '../../../src/domain/researchLead.js' interface Quelle { id: string name: string url: string laden: () => Promise } /** * Für diese drei Adressen gibt es einen eigens geschriebenen Leser. * * Der Schlüssel ist die Adresse und nicht die Kennung: Ein Zugang, den der * Nutzer im Personalblatt anlegt, bekommt eine eigene Kennung, zeigt aber * vielleicht auf dieselbe Seite. Dann soll er den guten Leser bekommen und * nicht den allgemeinen. */ const SPEZIALLESER = new Map Promise>([ [SOURCES.ZHK.url, fetchZhkResearch], [SOURCES.GZA.url, fetchGreaterZurichResearch], [SOURCES.ZEFIX.url, fetchZefixResearch], ]) /** Die gepflegten drei — Rückfall, wenn die Anfrage keine Quellen mitbringt. */ const STANDARDQUELLEN: Quelle[] = [ { ...SOURCES.ZHK, laden: fetchZhkResearch }, { ...SOURCES.GZA, laden: fetchGreaterZurichResearch }, { ...SOURCES.ZEFIX, laden: fetchZefixResearch }, ] /** Was die Oberfläche mitschickt (Runde 10, §§2.3/2.6/2.9). */ interface RefreshBody { sources?: { id: string; name: string; url: string }[] documents?: { id: string; name: string; text: string }[] temporalClasses?: TemporalClass[] } /** Höchstzahl frei angebundener Quellen je Lauf — schützt Laufzeit und Kosten. */ const MAX_SOURCES = 8 /** Höchstzahl abgelegter Dokumente je Lauf. */ const MAX_DOCUMENTS = 12 /** Wie viel Text ein abgelegtes Dokument beisteuert. */ const DOC_TEXT_LIMIT = 12_000 async function leseBody(req: IncomingMessage): Promise { const stuecke: Buffer[] = [] for await (const chunk of req) stuecke.push(chunk as Buffer) if (stuecke.length === 0) return {} try { return JSON.parse(Buffer.concat(stuecke).toString('utf8')) as RefreshBody } catch { // Ein unlesbarer Rumpf ist kein Grund, den Lauf abzubrechen — dann gelten // eben die gepflegten Standardquellen. return {} } } /** * Aus den übergebenen Adressen die Leseaufträge bauen. * * Genau hier wird die Quellenregistry dynamisch: Was das Personalblatt als * aktiven Lesezugang mit Adresse führt, landet in dieser Liste — mit dem * spezialisierten Leser, wenn es einen gibt, und sonst mit dem allgemeinen. */ function baueQuellen(body: RefreshBody): Quelle[] { const angefragt = (body.sources ?? []).filter(q => q.url && q.name).slice(0, MAX_SOURCES) if (angefragt.length === 0) return STANDARDQUELLEN return angefragt.map((q): Quelle => { const speziell = SPEZIALLESER.get(q.url) return { id: q.id, name: q.name, url: q.url, laden: speziell ?? (() => fetchGenericSource(q)), } }) } /** * Abgelegte Dokumente als Einträge — dieselbe Beurteilung wie bei den Seiten. * * Das ist der ganze Mechanismus aus §2.6: Ein hochgeladenes PDF unterscheidet * sich für die Analyse in nichts von einem gelesenen Artikel, sobald sein Text * extrahiert ist. Ein zweiter Analysepfad dafür wäre eine zweite Stelle, an der * sich die Beurteilung ändern könnte. */ function dokumenteAlsEintraege(body: RefreshBody): ResearchItem[] { return (body.documents ?? []) .filter(d => typeof d.text === 'string' && d.text.trim().length > 120) .slice(0, MAX_DOCUMENTS) .map(d => ({ source: 'Manuell abgelegtes Dokument', title: d.name, // Es gibt keine Webadresse — die Kennung des Dokuments tritt an ihre Stelle. url: `dokument:${d.id}`, text: d.text.slice(0, DOC_TEXT_LIMIT), })) } export default async function handler(req: IncomingMessage, res: ServerResponse): Promise { res.setHeader('Content-Type', 'application/json; charset=utf-8') // Ein Live-Lauf soll nie aus einem Zwischenspeicher beantwortet werden. res.setHeader('Cache-Control', 'no-store') if (req.method !== 'POST') { res.statusCode = 405 res.end(JSON.stringify({ error: 'Nur POST' })) return } const body = await leseBody(req) const quellen = baueQuellen(body) const dokumente = dokumenteAlsEintraege(body) // Alle parallel: eine langsame Quelle soll die anderen nicht aufhalten. const ergebnisse = await Promise.all( quellen.map(async (q): Promise<{ meta: ResearchSourceResult; items: ResearchItem[] }> => { try { const items = await q.laden() return { meta: { id: q.id, name: q.name, url: q.url, ok: true, itemCount: items.length }, items, } } catch (err) { // Eine kaputte Quelle darf die Demo nicht zerstören (§21). return { meta: { id: q.id, name: q.name, url: q.url, ok: false, itemCount: 0, error: err instanceof Error ? err.message : 'Unbekannter Fehler', }, items: [], } } }), ) const sources = ergebnisse.map(e => e.meta) // Abgelegte Dokumente erscheinen als eigene Zeile in der Quellenübersicht — // wer sie hochgeladen hat, soll sehen, dass sie tatsächlich gelesen wurden. if (dokumente.length > 0) { sources.push({ id: 'dokumente', name: `Manuell abgelegte Dokumente (${dokumente.length})`, url: '', ok: true, itemCount: dokumente.length, }) } const items = [...ergebnisse.flatMap(e => e.items), ...dokumente] const basis: ResearchRefreshResult = { sources, analyzed: items.length, discarded: 0, watchlist: 0, temporalFiltered: 0, leads: [], refreshedAt: new Date().toISOString(), model: MODEL, } if (items.length === 0) { res.statusCode = 200 res.end(JSON.stringify({ ...basis, analysisError: 'Keine der angebundenen Quellen lieferte Einträge.', })) return } try { const { leads, discarded, watchlist, temporalFiltered } = await analyzeResearchItems(items, body.temporalClasses ?? []) res.statusCode = 200 res.end(JSON.stringify({ ...basis, leads, discarded, watchlist, temporalFiltered })) } catch (err) { // Gelesen wurde trotzdem — das sagen wir offen, statt eine leere Liste als // «nichts gefunden» auszugeben. res.statusCode = 200 res.end(JSON.stringify({ ...basis, analysisError: err instanceof Error ? err.message : 'Die Beurteilung schlug fehl.', })) } }