import { basename, join } from "node:path"; import { mkdir } from "node:fs/promises"; import { listCompanies, summarizeCompany, writeSummary, type LlmCallStatus, type LlmCallTrace } from "./summarizer"; import { dataPath } from "../../deployPaths"; import { allowedModelFromEnv, isSafeStem, jsonError, validateRequestedModel, withAuth } from "../../security"; const OUTPUTS_DIR = process.env.BMP_SUMMARIZER_OUTPUTS_DIR ?? dataPath(join(import.meta.dir, "outputs"), "summarizer", "outputs"); const PAGE_PATH = join(import.meta.dir, "index.html"); const MAX_PARALLEL_SUMMARIES = 2; const MAX_ACTIVE_JOBS = Number(process.env.BMP_MAX_ACTIVE_SUMMARY_JOBS ?? 10); const EXPECTED_LLM_CALLS_PER_COMPANY = 13; const OBSERVABILITY_DIR = process.env.BMP_OBSERVABILITY_DIR ?? dataPath(join(import.meta.dir, "outputs", "_observability"), "summarizer", "outputs", "_observability"); type JobStatus = "queued" | "running" | "done" | "error" | "cancelled"; interface SummaryJob { id: string; status: JobStatus; stem: string; model: string; total: number; completed: number; current?: string; files: string[]; downloads: SummaryDownload[]; errors: string[]; llmCalls: SafeLlmCallTrace[]; observability: ObservabilityLog; observabilityWrite: Promise; usage: UsageSummary; statusMessage?: string; createdAt: string; finishedAt?: string; } type SafeLlmCallTrace = Pick< LlmCallTrace, "id" | "operation" | "label" | "model" | "status" | "startedAt" | "finishedAt" | "durationMs" | "error" > & { attempt: number; maxAttempts: number; usage?: unknown; }; interface ObservabilityEvent { at: string; type: string; message?: string; data?: unknown; } interface ObservabilityLog { version: 1; jobId: string; stem: string; model: string; status: JobStatus; createdAt: string; finishedAt?: string; total: number; completed: number; current?: string; companies: string[]; files: string[]; downloads: SummaryDownload[]; errors: SerializedError[]; inputs: Record; outputs: Record; llmCalls: LlmCallTrace[]; events: ObservabilityEvent[]; logPath: string; } interface SerializedError { name?: string; message: string; stack?: string; cause?: unknown; } interface SummaryDownload { company: string; label: string; url: string; } interface UsageSummary { cost: number; upstreamInferenceCost: number; promptTokens: number; completionTokens: number; totalTokens: number; reasoningTokens: number; cachedTokens: number; calls: number; } interface LlmProgress { total: number; running: number; done: number; error: number; latest?: { label: string; operation: LlmCallTrace["operation"]; status: LlmCallStatus; startedAt: string; finishedAt?: string; durationMs?: number; }; } function sanitizeTrace(trace: LlmCallTrace): SafeLlmCallTrace { const response = trace.response as { usage?: unknown } | undefined; return { id: trace.id, operation: trace.operation, label: trace.label, model: trace.model, attempt: trace.attempt, maxAttempts: trace.maxAttempts, status: trace.status, startedAt: trace.startedAt, finishedAt: trace.finishedAt, durationMs: trace.durationMs, usage: response?.usage, error: trace.error, }; } function emptyUsageSummary(): UsageSummary { return { cost: 0, upstreamInferenceCost: 0, promptTokens: 0, completionTokens: 0, totalTokens: 0, reasoningTokens: 0, cachedTokens: 0, calls: 0, }; } function summarizeUsage(calls: SafeLlmCallTrace[]): UsageSummary { const summary = emptyUsageSummary(); for (const call of calls) { const usage = readUsage(call.usage); if (!usage) continue; summary.calls += 1; summary.cost += usage.cost; summary.upstreamInferenceCost += usage.upstreamInferenceCost; summary.promptTokens += usage.promptTokens; summary.completionTokens += usage.completionTokens; summary.totalTokens += usage.totalTokens; summary.reasoningTokens += usage.reasoningTokens; summary.cachedTokens += usage.cachedTokens; } return { ...summary, cost: Number(summary.cost.toFixed(8)), upstreamInferenceCost: Number(summary.upstreamInferenceCost.toFixed(8)), }; } function summarizeLlmProgress(job: SummaryJob): LlmProgress { const calls = job.llmCalls; const progress: LlmProgress = { total: Math.max(1, job.total) * EXPECTED_LLM_CALLS_PER_COMPANY, running: 0, done: 0, error: 0, }; for (const call of calls) { progress[call.status] += 1; } const latest = [...calls].sort((a, b) => { const aTime = Date.parse(a.finishedAt ?? a.startedAt); const bTime = Date.parse(b.finishedAt ?? b.startedAt); return bTime - aTime; })[0]; if (latest) { progress.latest = { label: latest.label, operation: latest.operation, status: latest.status, startedAt: latest.startedAt, finishedAt: latest.finishedAt, durationMs: latest.durationMs, }; } return progress; } function readUsage(value: unknown): UsageSummary | undefined { if (typeof value !== "object" || value === null) return undefined; const usage = value as Record; const promptDetails = usage.prompt_tokens_details as Record | undefined; const completionDetails = usage.completion_tokens_details as Record | undefined; const costDetails = usage.cost_details as Record | undefined; return { cost: readNumber(usage.cost ?? usage.total_cost), upstreamInferenceCost: readNumber(costDetails?.upstream_inference_cost ?? costDetails?.total_cost), promptTokens: readNumber(usage.prompt_tokens), completionTokens: readNumber(usage.completion_tokens), totalTokens: readNumber(usage.total_tokens), reasoningTokens: readNumber(completionDetails?.reasoning_tokens), cachedTokens: readNumber(promptDetails?.cached_tokens), calls: 1, }; } function readNumber(value: unknown): number { if (typeof value === "number" && Number.isFinite(value)) return value; if (typeof value === "string") { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : 0; } return 0; } function publicJob(job: SummaryJob) { return { id: job.id, status: job.status, stem: job.stem, model: job.model, total: job.total, completed: job.completed, current: job.current, downloads: job.downloads, errors: job.errors, usage: job.usage, llmProgress: summarizeLlmProgress(job), statusMessage: job.statusMessage, createdAt: job.createdAt, finishedAt: job.finishedAt, }; } function activeJobCount(): number { return [...jobs.values()].filter((job) => job.status === "queued" || job.status === "running").length; } function isCancelled(job: SummaryJob): boolean { return job.status === "cancelled"; } const jobs = new Map(); function serializeError(error: unknown): SerializedError { if (error instanceof Error) { return { name: error.name, message: error.message, stack: error.stack, cause: error.cause, }; } return { message: String(error) }; } function createObservabilityLog(jobId: string, stem: string, model: string, total: number, createdAt: string): ObservabilityLog { return { version: 1, jobId, stem, model, status: "queued", createdAt, total, completed: 0, companies: [], files: [], downloads: [], errors: [], inputs: {}, outputs: {}, llmCalls: [], events: [ { at: createdAt, type: "job.created", message: "Summary job was queued.", }, ], logPath: join(OBSERVABILITY_DIR, `${createdAt.replaceAll(":", "-")}-${jobId}.json`), }; } function recordEvent(job: SummaryJob, type: string, message?: string, data?: unknown): void { job.observability.events.push({ at: new Date().toISOString(), type, message, data, }); } function syncObservabilityState(job: SummaryJob): void { job.observability.status = job.status; job.observability.finishedAt = job.finishedAt; job.observability.completed = job.completed; job.observability.current = job.current; job.observability.files = [...job.files]; job.observability.downloads = [...job.downloads]; } function persistObservability(job: SummaryJob): Promise { syncObservabilityState(job); job.observabilityWrite = job.observabilityWrite .catch(() => undefined) .then(async () => { await mkdir(OBSERVABILITY_DIR, { recursive: true }); await Bun.write(job.observability.logPath, JSON.stringify(job.observability, null, 2)); }) .catch((error) => { console.error("Failed to write summarizer observability log:", error); }); return job.observabilityWrite; } function createJob(stem: string, model: string, total: number): SummaryJob { const id = crypto.randomUUID(); const createdAt = new Date().toISOString(); const job: SummaryJob = { id, status: "queued", stem, model, total, completed: 0, files: [], downloads: [], errors: [], llmCalls: [], observability: createObservabilityLog(id, stem, model, total, createdAt), observabilityWrite: Promise.resolve(), usage: emptyUsageSummary(), createdAt, }; jobs.set(job.id, job); void persistObservability(job); return job; } function labelForOutputFile(file: string): string { if (file.endsWith(".json")) return "JSON herunterladen"; if (file.endsWith(".xlsx")) return "XLSX herunterladen"; if (file.endsWith(".questions.pdf")) return "Fragen-PDF herunterladen"; if (file.endsWith(".questions.html")) return "Fragen-HTML herunterladen"; if (file.endsWith(".pdf")) return "PDF herunterladen"; if (file.endsWith(".report.html")) return "Report-HTML herunterladen"; return basename(file); } function buildDownload(company: string, filePath: string): SummaryDownload { const file = basename(filePath); return { company, label: labelForOutputFile(file), url: `/downloads/summarizer/${encodeURIComponent(company)}/${encodeURIComponent(file)}`, }; } async function runWithConcurrency( items: T[], limit: number, worker: (item: T) => Promise, ): Promise { let nextIndex = 0; async function runner() { while (nextIndex < items.length) { const index = nextIndex++; await worker(items[index]!); } } const count = Math.min(limit, items.length); await Promise.all(Array.from({ length: count }, () => runner())); } async function executeJob(job: SummaryJob): Promise { if (isCancelled(job)) return; job.status = "running"; job.statusMessage = "Auswertung wird vorbereitet."; recordEvent(job, "job.started", job.statusMessage); await persistObservability(job); const companies = job.stem === "__all__" ? await listCompanies() : [job.stem]; job.observability.companies = [...companies]; recordEvent(job, "companies.selected", undefined, { companies }); await persistObservability(job); try { await runWithConcurrency(companies, MAX_PARALLEL_SUMMARIES, async (company) => { if (isCancelled(job)) return; job.current = company; job.statusMessage = `LLM-Auswertung für ${company} läuft.`; recordEvent(job, "company.started", job.statusMessage, { company }); await persistObservability(job); const summary = await summarizeCompany(company, job.model, { onCompanyInput: (companyStem, input) => { job.observability.inputs[companyStem] = input; recordEvent(job, "company.input.loaded", undefined, { company: companyStem }); void persistObservability(job); }, onCompanyOutput: (companyStem, output) => { job.observability.outputs[companyStem] = output; recordEvent(job, "company.output.created", undefined, { company: companyStem }); void persistObservability(job); }, onLlmCall: (trace) => { const existingIndex = job.llmCalls.findIndex((call) => call.id === trace.id); const safeTrace = sanitizeTrace(trace); if (existingIndex >= 0) { job.llmCalls[existingIndex] = safeTrace; } else { job.llmCalls.push(safeTrace); } const fullTraceIndex = job.observability.llmCalls.findIndex((call) => call.id === trace.id); if (fullTraceIndex >= 0) { job.observability.llmCalls[fullTraceIndex] = trace; } else { job.observability.llmCalls.push(trace); } job.usage = summarizeUsage(job.llmCalls); job.statusMessage = trace.status === "running" ? `${trace.label} wird verarbeitet (Versuch ${trace.attempt}/${trace.maxAttempts}).` : trace.status === "error" ? `${trace.label} fehlgeschlagen (Versuch ${trace.attempt}/${trace.maxAttempts}).` : `${trace.label} abgeschlossen.`; recordEvent(job, `llm.${trace.status}`, job.statusMessage, { id: trace.id, operation: trace.operation, label: trace.label, attempt: trace.attempt, maxAttempts: trace.maxAttempts, }); void persistObservability(job); }, }); job.statusMessage = `Downloads für ${company} werden vorbereitet.`; recordEvent(job, "company.summary.created", job.statusMessage, { company }); if (isCancelled(job)) return; const outputPaths = await writeSummary(company, summary, OUTPUTS_DIR); job.files.push(...outputPaths.map((path) => basename(path))); job.downloads.push(...outputPaths.map((path) => buildDownload(company, path))); job.completed += 1; recordEvent(job, "company.finished", `Auswertung für ${company} fertig.`, { company, outputFiles: outputPaths.map((path) => basename(path)), }); await persistObservability(job); }); if (isCancelled(job)) { job.current = undefined; job.statusMessage = "Auswertung abgebrochen."; job.finishedAt = new Date().toISOString(); recordEvent(job, "job.cancelled", job.statusMessage); await persistObservability(job); return; } job.status = "done"; job.current = undefined; job.statusMessage = "Auswertung fertig."; job.finishedAt = new Date().toISOString(); recordEvent(job, "job.done", job.statusMessage); await persistObservability(job); } catch (error) { if (isCancelled(job)) { job.current = undefined; job.statusMessage = "Auswertung abgebrochen."; job.finishedAt = new Date().toISOString(); recordEvent(job, "job.cancelled", job.statusMessage); await persistObservability(job); return; } const serializedError = serializeError(error); job.status = "error"; job.errors.push(serializedError.message); job.observability.errors.push(serializedError); job.statusMessage = `Auswertung fehlgeschlagen: ${serializedError.message}`; job.finishedAt = new Date().toISOString(); recordEvent(job, "job.error", job.statusMessage, serializedError); await persistObservability(job); } } export const routes = { "/summarizer": { GET: withAuth(async () => new Response(Bun.file(PAGE_PATH), { headers: { "Content-Type": "text/html; charset=utf-8" }, })), }, "/api/summarizer/companies": { GET: withAuth(async () => { const companies = await listCompanies(); return Response.json({ companies }); }), }, "/api/summarizer/summarize": { POST: withAuth(async (req: Request) => { let body: { stem?: string; model?: string }; try { body = await req.json(); } catch { return Response.json({ error: "Expected JSON body" }, { status: 400 }); } const { stem } = body; const model = validateRequestedModel(body.model); if (!stem || (stem !== "__all__" && !isSafeStem(stem))) { return Response.json({ error: 'Missing or invalid "stem" field' }, { status: 400 }); } if (!model) { return Response.json( { error: `Model is not allowed for this demo. Use ${allowedModelFromEnv()}.` }, { status: 400 }, ); } if (!process.env.OPENROUTER_API_KEY) { return Response.json( { error: "OPENROUTER_API_KEY is not set in environment" }, { status: 500 }, ); } if (activeJobCount() >= MAX_ACTIVE_JOBS) { return jsonError(`Es laufen bereits ${MAX_ACTIVE_JOBS} Auswertungen. Bitte warten Sie, bis eine davon fertig ist.`, 429); } const companies = await listCompanies(); if (stem !== "__all__" && !companies.includes(stem)) { return Response.json({ error: "Company not found" }, { status: 404 }); } const total = stem === "__all__" ? (await listCompanies()).length : 1; const job = createJob(stem, model, total); setTimeout(() => { void executeJob(job); }, 0); return Response.json({ ok: true, jobId: job.id }); }, { csrf: true, limit: { key: "summarize", max: Math.max(20, MAX_ACTIVE_JOBS * 2), windowMs: 60_000 }, }), }, "/api/summarizer/jobs": { GET: withAuth(async () => { const summaries = [...jobs.values()] .toSorted((a, b) => Date.parse(b.createdAt) - Date.parse(a.createdAt)) .map(publicJob); return Response.json({ jobs: summaries }); }), }, "/api/summarizer/jobs/:id/llm-calls": { GET: withAuth(async (req: Request) => { const id = (req as Request & { params: Record }).params.id; if (!id) { return Response.json({ error: "Job not found" }, { status: 404 }); } const job = jobs.get(id); if (!job) { return Response.json({ error: "Job not found" }, { status: 404 }); } return Response.json({ jobId: job.id, status: job.status, total: job.total, completed: job.completed, current: job.current, usage: job.usage, llmCalls: job.llmCalls, }); }), }, "/api/summarizer/jobs/:id": { GET: withAuth(async (req: Request) => { const id = (req as Request & { params: Record }).params.id; if (!id) { return Response.json({ error: "Job not found" }, { status: 404 }); } const job = jobs.get(id); if (!job) { return Response.json({ error: "Job not found" }, { status: 404 }); } return Response.json(publicJob(job)); }), DELETE: withAuth(async (req: Request) => { const id = (req as Request & { params: Record }).params.id; if (!id) { return Response.json({ error: "Job not found" }, { status: 404 }); } const job = jobs.get(id); if (!job) { return Response.json({ error: "Job not found" }, { status: 404 }); } if (job.status === "queued" || job.status === "running") { job.status = "cancelled"; job.current = undefined; job.statusMessage = "Auswertung abgebrochen."; job.finishedAt = new Date().toISOString(); return Response.json({ ok: true, job: publicJob(job) }); } jobs.delete(id); return Response.json({ ok: true, deleted: true }); }, { csrf: true, limit: { key: "summarizer-job-delete", max: 30, windowMs: 60_000 }, }), }, } as const;