From ff62e751ed0ed86e934981e533b99a7977390b68 Mon Sep 17 00:00:00 2001 From: Elton Turing Date: Mon, 7 Sep 2026 20:20:55 +0200 Subject: [PATCH] llmrouter: Slot-Leak beheben (42i/intern#1283) Ein abgebrochener Streaming-Aufruf gegen das lokale Qwen (maxParallel 1) hat den einzigen Slot dauerhaft gefressen: die Freigabe hing allein im flush() des TransformStreams, und der laeuft bei cancel() nie. Danach reihte sich jeder weitere Aufruf UNBEGRENZT in die Warteschlange ein, statt zu scheitern -- auto/mass, role/ops (Buzz) und das small_model aller Container-Personas hingen zwei Stunden, ohne dass irgendwo etwas rot wurde. Diagnose von Elton im Selbstheilungslauf, dreimal in Folge belegt. Drei Haerten: - Die Freigabe ist gegen Doppelaufruf gesichert (Flag). Erst dadurch darf sie an mehreren Stellen stehen, ohne Slots zu erfinden. - cancel() gibt jetzt frei, nicht nur flush(). Das ist der eigentliche Fehler: bricht der Client ab, laeuft flush() nie. - Wer laenger als 30 s auf einen Slot wartet, bekommt 503 mit Grund. Ein Fehler, an dem der Aufrufer etwas aendern kann, ist besser als Stille. Ausgeloest habe ich den Fall selbst: zwei eigene Testaufrufe gegen role/ops liefen gegen 18:05 ins Timeout und wurden abgebrochen -- genau die Handlung, die den Slot frisst. Der Fehler im Code ist aelter, der Ausloeser war meiner. Co-Authored-By: Claude Opus 5 --- src/server.ts | 70 ++++++++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 66 insertions(+), 4 deletions(-) diff --git a/src/server.ts b/src/server.ts index 5690d81..a0abf33 100644 --- a/src/server.ts +++ b/src/server.ts @@ -63,13 +63,58 @@ async function altSicherstellen(ref: string): Promise<{ model: string; effort: E // --- Nebenlaeufigkeit je Upstream (lokales Qwen: ein Slot) ------------------ const semaphoren = new Map void)[] }>(); +const SLOT_FRIST_MS = 30_000; + +/** + * Belegt einen Slot des Upstreams und liefert die Freigabe zurueck. + * + * Zwei Haerten, beide aus 42i/intern#1283 (2026-09-07): Ein abgebrochener + * Streaming-Aufruf gegen das lokale Qwen (maxParallel 1) hat den einzigen Slot + * dauerhaft gefressen -- die Freigabe hing allein im flush() des + * TransformStreams, und der laeuft bei cancel() nie. Danach reihte sich jeder + * weitere Aufruf unbegrenzt in die Warteschlange ein, statt zu scheitern: + * `auto/mass`, `role/ops` und das small_model aller Container-Personas hingen + * zwei Stunden lang, ohne dass irgendwo etwas rot wurde. + * + * 1. Die Freigabe ist gegen Doppelaufruf gesichert. Sie darf jetzt an jeder + * Stelle stehen, auch mehrfach und in einem finally -- ohne das Flag wuerde + * ein zweiter Aufruf einen Slot erzeugen, den es nie gab. + * 2. Wer laenger als SLOT_FRIST_MS wartet, bekommt eine Absage statt ewiger + * Stille. Ein 503 ist eine Antwort, an der ein Aufrufer etwas aendern kann; + * ein haengender Aufruf ist es nicht. + */ async function belegen(name: string, u: Upstream): Promise<() => void> { if (u.maxParallel === undefined) return () => {}; let s = semaphoren.get(name); if (s === undefined) { s = { frei: u.maxParallel, warteschlange: [] }; semaphoren.set(name, s); } - if (s.frei > 0) s.frei -= 1; - else await new Promise((res) => s!.warteschlange.push(res)); - return () => { const n = s!.warteschlange.shift(); if (n) n(); else s!.frei += 1; }; + if (s.frei > 0) { + s.frei -= 1; + } else { + let platz: (() => void) | undefined; + const bekommen = await new Promise((res) => { + platz = () => res(true); + s!.warteschlange.push(platz); + setTimeout(() => { + const i = s!.warteschlange.indexOf(platz!); + if (i >= 0) { s!.warteschlange.splice(i, 1); res(false); } + }, SLOT_FRIST_MS); + }); + if (!bekommen) throw new SlotFrist(name); + } + let schon = false; + return () => { + if (schon) return; + schon = true; + const n = s!.warteschlange.shift(); + if (n) n(); else s!.frei += 1; + }; +} + +/** Kein Slot innerhalb der Frist -- der Aufrufer bekommt 503 statt Stille. */ +class SlotFrist extends Error { + constructor(public readonly upstream: string) { + super(`Upstream ${upstream} ist ausgelastet (kein Slot in ${SLOT_FRIST_MS / 1000} s)`); + } } // --- Aufloesung Name -> Ziel --------------------------------------------------- @@ -297,7 +342,20 @@ Bun.serve({ try { log.schreiben(eintrag); } catch (e) { sagen(`Log-Fehler: ${(e as Error).message}`); } }; - let freigeben = await belegen(ziel.upstreamName, ziel.upstream); + // Kein Slot in der Frist -> 503 statt Stille (42i/intern#1283). Der + // Aufrufer sieht, dass der Upstream ausgelastet ist, und kann es spaeter + // erneut versuchen; vorher hing er unbegrenzt. + let freigeben: () => void; + try { + freigeben = await belegen(ziel.upstreamName, ziel.upstream); + } catch (e) { + if (e instanceof SlotFrist) { + eintrag.status = 503; eintrag.fehler = e.message; abschliessen(); + sagen(`503 ${eintrag.key} ${ziel.alias}: ${e.message}`); + return Response.json({ error: { message: e.message, code: 503 } }, { status: 503 }); + } + throw e; + } let antwort: Response; try { antwort = await fetch(`${ziel.upstream.base}${pfad}`, { method: "POST", headers, body: upstreamBody, signal: AbortSignal.timeout(600_000) }); @@ -365,6 +423,10 @@ Bun.serve({ } }, flush() { freigeben(); abschliessen(usage); }, + // Bricht der Client ab, laeuft flush() NIE -- ohne diese Zeile bleibt + // der Slot fuer immer belegt (42i/intern#1283). Die Freigabe ist gegen + // Doppelaufruf gesichert, beide Wege duerfen also feuern. + cancel() { freigeben(); abschliessen(usage); }, }); return new Response(antwort.body.pipeThrough(tee), { status: antwort.status, headers: weiter }); }