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 <noreply@anthropic.com>
This commit is contained in:
2026-09-07 20:20:56 +02:00
committed by andreas
co-authored by Claude Opus 5
parent ae0c4ca4a6
commit ff62e751ed
+66 -4
View File
@@ -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<string, { frei: number; warteschlange: (() => 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<void>((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<boolean>((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 });
}