feat(compute): Worker-Selbstanmeldung ueber RVS + Flotten-Anzeige (Stage 2)
Jeder GPU-Dienst meldet sich beim Connect mit worker_hello {instanceId,
service, node, gpus, model} und haelt die Registry per periodischem
worker_ping {instanceId, busy} (~10s) frisch. So weiss ARIA, was wo laeuft.
- xtts/{voxtral,whisper,f5tts}/bridge.py + llm-adapter/adapter.py:
INSTANCE_ID=service@NODE_NAME, _worker_register()-Coroutine (hello + ping),
busy-Quelle je Worker (aktive STT-Sessions / TTS-Render / in-flight LLM);
Task sauber gecancelt bei Reconnect.
- bridge/aria_bridge.py: self._workers-Registry + Handler worker_hello/
worker_ping (spiegelt sat_hello), _worker_list() (35s-Offline-TTL),
_pick_worker() (Round-Robin freie Instanz, fuer Stage-3-Routing),
/internal/worker-list-Endpoint.
- diagnostic/server.js: workers-Map, worker_hello/worker_ping-Tracking,
worker_update-Broadcast + worker_list-Action + on-connect-Snapshot.
- diagnostic/index.html: "Compute-Flotte"-Panel im Satelliten-Tab — pro Node
gruppiert, mit Dienst/Modell/GPU und frei/beschaeftigt/offline-Status.
Stage 2 von 3. Reine Sichtbarkeit, kein Routing-Verhalten geaendert.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -766,6 +766,11 @@ class ARIABridge:
|
|||||||
# requestId → Future (sat_devices / sat_result), analog _pending_flux.
|
# requestId → Future (sat_devices / sat_result), analog _pending_flux.
|
||||||
self._satellites: dict[str, dict] = {}
|
self._satellites: dict[str, dict] = {}
|
||||||
self._pending_sat: dict[str, asyncio.Future] = {}
|
self._pending_sat: dict[str, asyncio.Future] = {}
|
||||||
|
# Compute-Fleet: GPU-Worker (voxtral/whisper/f5tts/llm) melden sich per
|
||||||
|
# worker_hello, halten sich per worker_ping frisch. instanceId →
|
||||||
|
# {service, node, gpus, model, busy, last_seen}. Genutzt fuer die
|
||||||
|
# Diagnostic-Flotten-Anzeige und (Stage 3) targetInstance-Routing.
|
||||||
|
self._workers: dict[str, dict] = {}
|
||||||
# FLUX-Render-Requests die aktuell auf Antwort der flux-bridge (Gamebox) warten.
|
# FLUX-Render-Requests die aktuell auf Antwort der flux-bridge (Gamebox) warten.
|
||||||
# requestId → Future mit dem flux_response-Payload (oder None bei Fehler).
|
# requestId → Future mit dem flux_response-Payload (oder None bei Fehler).
|
||||||
self._pending_flux: dict[str, asyncio.Future] = {}
|
self._pending_flux: dict[str, asyncio.Future] = {}
|
||||||
@@ -3528,6 +3533,40 @@ class ARIABridge:
|
|||||||
self._satellites[sid]["caps"], self._satellites[sid]["control"])
|
self._satellites[sid]["caps"], self._satellites[sid]["control"])
|
||||||
return
|
return
|
||||||
|
|
||||||
|
elif msg_type == "worker_hello":
|
||||||
|
iid = (payload.get("instanceId") or "").strip()
|
||||||
|
if iid:
|
||||||
|
prev = self._workers.get(iid, {})
|
||||||
|
self._workers[iid] = {
|
||||||
|
"instanceId": iid,
|
||||||
|
"service": payload.get("service") or "",
|
||||||
|
"node": payload.get("node") or "",
|
||||||
|
"gpus": payload.get("gpus") or "",
|
||||||
|
"model": payload.get("model") or "",
|
||||||
|
"busy": bool(prev.get("busy", False)),
|
||||||
|
"last_seen": time.time(),
|
||||||
|
}
|
||||||
|
logger.info("[worker] online: %s (service=%s node=%s gpus=%s model=%s)",
|
||||||
|
iid, self._workers[iid]["service"], self._workers[iid]["node"],
|
||||||
|
self._workers[iid]["gpus"] or "?", self._workers[iid]["model"] or "?")
|
||||||
|
return
|
||||||
|
|
||||||
|
elif msg_type == "worker_ping":
|
||||||
|
iid = (payload.get("instanceId") or "").strip()
|
||||||
|
if iid:
|
||||||
|
w = self._workers.get(iid)
|
||||||
|
if w is None:
|
||||||
|
# Ping ohne vorheriges hello (Bridge-Neustart) → Minimal-Eintrag,
|
||||||
|
# service aus der instanceId ableiten (Form: "service@node").
|
||||||
|
svc = iid.split("@", 1)[0]
|
||||||
|
w = self._workers[iid] = {
|
||||||
|
"instanceId": iid, "service": svc, "node": "", "gpus": "",
|
||||||
|
"model": "", "busy": False, "last_seen": 0.0,
|
||||||
|
}
|
||||||
|
w["busy"] = bool(payload.get("busy", False))
|
||||||
|
w["last_seen"] = time.time()
|
||||||
|
return
|
||||||
|
|
||||||
elif msg_type in ("sat_devices", "sat_result"):
|
elif msg_type in ("sat_devices", "sat_result"):
|
||||||
req_id = payload.get("requestId", "")
|
req_id = payload.get("requestId", "")
|
||||||
future = self._pending_sat.get(req_id)
|
future = self._pending_sat.get(req_id)
|
||||||
@@ -4435,6 +4474,9 @@ class ARIABridge:
|
|||||||
elif method == "POST" and path == "/internal/satellite-list":
|
elif method == "POST" and path == "/internal/satellite-list":
|
||||||
# Brain fragt: welche Satelliten/Netze sind online + Capabilities.
|
# Brain fragt: welche Satelliten/Netze sind online + Capabilities.
|
||||||
await _send_response(writer, 200, {"ok": True, "satellites": self._satellite_list()})
|
await _send_response(writer, 200, {"ok": True, "satellites": self._satellite_list()})
|
||||||
|
elif method in ("GET", "POST") and path == "/internal/worker-list":
|
||||||
|
# Diagnostic/Brain fragt: welche Compute-Worker sind online (Flotte).
|
||||||
|
await _send_response(writer, 200, {"ok": True, "workers": self._worker_list()})
|
||||||
elif method == "POST" and path == "/internal/satellite":
|
elif method == "POST" and path == "/internal/satellite":
|
||||||
# Brain-Tool: Discovery oder Command an einen Satelliten.
|
# Brain-Tool: Discovery oder Command an einen Satelliten.
|
||||||
# body: {op:'discover'|'command', satellite, device?, action?, params?}
|
# body: {op:'discover'|'command', satellite, device?, action?, params?}
|
||||||
@@ -4684,6 +4726,44 @@ class ARIABridge:
|
|||||||
})
|
})
|
||||||
return out
|
return out
|
||||||
|
|
||||||
|
# worker_ping kommt alle ~10s; nach 35s ohne Ping gilt ein Worker als offline.
|
||||||
|
WORKER_OFFLINE_S = 35
|
||||||
|
|
||||||
|
def _worker_list(self) -> list[dict]:
|
||||||
|
"""Bekannte Compute-Worker (Flotte). online = kuerzlich per Ping gesehen."""
|
||||||
|
now = time.time()
|
||||||
|
out = []
|
||||||
|
for w in self._workers.values():
|
||||||
|
out.append({
|
||||||
|
"instanceId": w["instanceId"], "service": w.get("service") or "",
|
||||||
|
"node": w.get("node") or "", "gpus": w.get("gpus") or "",
|
||||||
|
"model": w.get("model") or "", "busy": bool(w.get("busy")),
|
||||||
|
"online": (now - w.get("last_seen", 0)) < self.WORKER_OFFLINE_S,
|
||||||
|
})
|
||||||
|
return out
|
||||||
|
|
||||||
|
def _pick_worker(self, service: str) -> Optional[str]:
|
||||||
|
"""Waehlt eine online, moeglichst freie Instanz des Diensts (Round-Robin
|
||||||
|
ueber die freien). Gibt die instanceId oder None (keine online). Fuer
|
||||||
|
Stage-3-Routing (targetInstance)."""
|
||||||
|
now = time.time()
|
||||||
|
online = [w for w in self._workers.values()
|
||||||
|
if w.get("service") == service
|
||||||
|
and (now - w.get("last_seen", 0)) < self.WORKER_OFFLINE_S]
|
||||||
|
if not online:
|
||||||
|
return None
|
||||||
|
free = [w for w in online if not w.get("busy")]
|
||||||
|
pool = free or online # alle busy → trotzdem eine nehmen (least-bad)
|
||||||
|
# Round-Robin: rotierender Zeiger pro Dienst.
|
||||||
|
rr = getattr(self, "_worker_rr", None)
|
||||||
|
if rr is None:
|
||||||
|
rr = self._worker_rr = {}
|
||||||
|
idx = rr.get(service, 0) % len(pool)
|
||||||
|
rr[service] = idx + 1
|
||||||
|
chosen = pool[idx]
|
||||||
|
chosen["busy"] = True # optimistisch, bis der naechste Ping korrigiert
|
||||||
|
return chosen["instanceId"]
|
||||||
|
|
||||||
async def _satellite_request(self, op: str, satellite: str = "",
|
async def _satellite_request(self, op: str, satellite: str = "",
|
||||||
device: str = "", action: str = "",
|
device: str = "", action: str = "",
|
||||||
params: Optional[dict] = None,
|
params: Optional[dict] = None,
|
||||||
|
|||||||
@@ -1307,6 +1307,22 @@
|
|||||||
<div class="card">
|
<div class="card">
|
||||||
<div id="sat-list" style="font-size:12px;color:#8888AA;">(Lade...)</div>
|
<div id="sat-list" style="font-size:12px;color:#8888AA;">(Lade...)</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
<div style="display:flex;justify-content:space-between;align-items:center;margin:16px 0 8px;">
|
||||||
|
<h2 style="margin:0;">Compute-Flotte 🖥️</h2>
|
||||||
|
<button class="btn secondary" onclick="requestWorkers()" style="padding:4px 10px;font-size:11px;">Aktualisieren</button>
|
||||||
|
</div>
|
||||||
|
<div class="card" style="margin-bottom:8px;">
|
||||||
|
<p style="color:#8888AA;font-size:12px;margin:0;">
|
||||||
|
GPU-Dienste (STT / TTS / lokales LLM), verteilt auf einen oder mehrere
|
||||||
|
Nodes. Jeder Node startet per <code>COMPOSE_PROFILES</code> nur die
|
||||||
|
Dienste, die er anbietet. Hier siehst Du, welche Instanz auf welchem
|
||||||
|
Node laeuft, mit welchem Modell/GPU und ob sie gerade beschaeftigt ist.
|
||||||
|
</p>
|
||||||
|
</div>
|
||||||
|
<div class="card">
|
||||||
|
<div id="worker-list" style="font-size:12px;color:#8888AA;">(Lade...)</div>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
</div><!-- /tab-satellites -->
|
</div><!-- /tab-satellites -->
|
||||||
|
|
||||||
@@ -1958,6 +1974,7 @@
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (msg.type === 'sat_update') { satellites = msg.satellites || []; renderSatellites(); return; }
|
if (msg.type === 'sat_update') { satellites = msg.satellites || []; renderSatellites(); return; }
|
||||||
|
if (msg.type === 'worker_update') { workers = msg.workers || []; renderWorkers(); return; }
|
||||||
if (msg.type === 'sat_devices') {
|
if (msg.type === 'sat_devices') {
|
||||||
if (msg.satellite) { satDevices[msg.satellite] = { devices: msg.devices || [], location: msg.location, ts: Date.now() }; }
|
if (msg.satellite) { satDevices[msg.satellite] = { devices: msg.devices || [], location: msg.location, ts: Date.now() }; }
|
||||||
satScanning = null;
|
satScanning = null;
|
||||||
@@ -4179,9 +4196,50 @@
|
|||||||
loadTriggers();
|
loadTriggers();
|
||||||
} else if (tab === 'satellites') {
|
} else if (tab === 'satellites') {
|
||||||
requestSatellites();
|
requestSatellites();
|
||||||
|
requestWorkers();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ── Compute-Flotte ─────────────────────────────────────
|
||||||
|
let workers = []; // [{instanceId, service, node, gpus, model, busy, online}]
|
||||||
|
const WORKER_SVC_META = {
|
||||||
|
voxtral: { icon: '🎙️', label: 'STT (Voxtral)' },
|
||||||
|
whisper: { icon: '🎙️', label: 'STT (Whisper)' },
|
||||||
|
f5tts: { icon: '🔊', label: 'TTS (F5-TTS)' },
|
||||||
|
llm: { icon: '🧠', label: 'LLM (lokal)' },
|
||||||
|
};
|
||||||
|
function requestWorkers() { send({ action: 'worker_list' }); }
|
||||||
|
function renderWorkers() {
|
||||||
|
const box = document.getElementById('worker-list');
|
||||||
|
if (!box) return;
|
||||||
|
if (!workers.length) {
|
||||||
|
box.innerHTML = '<span style="color:#8888AA;">Kein Compute-Node verbunden. ' +
|
||||||
|
'Starte auf einem GPU-Rechner <code>cd xtts && docker compose up -d</code> ' +
|
||||||
|
'(mit gesetztem <code>COMPOSE_PROFILES</code>).</span>';
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Nach Node gruppieren.
|
||||||
|
const byNode = {};
|
||||||
|
for (const w of workers) { (byNode[w.node || '?'] = byNode[w.node || '?'] || []).push(w); }
|
||||||
|
box.innerHTML = Object.keys(byNode).sort().map(node => {
|
||||||
|
const rows = byNode[node].map(w => {
|
||||||
|
const meta = WORKER_SVC_META[w.service] || { icon: '⚙️', label: w.service };
|
||||||
|
const dot = !w.online ? '#666' : (w.busy ? '#FFB020' : '#3FFF3F');
|
||||||
|
const stat = !w.online ? 'offline' : (w.busy ? 'beschaeftigt' : 'frei');
|
||||||
|
return '<div style="display:flex;align-items:center;gap:8px;padding:4px 0;">' +
|
||||||
|
'<span style="width:8px;height:8px;border-radius:50%;background:' + dot + ';display:inline-block;"></span>' +
|
||||||
|
'<span>' + meta.icon + ' <b>' + escapeHtml(meta.label) + '</b></span>' +
|
||||||
|
'<span style="color:#8888AA;">' + escapeHtml(w.model || '') + '</span>' +
|
||||||
|
(w.gpus ? '<span style="color:#8888AA;">GPU ' + escapeHtml(w.gpus) + '</span>' : '') +
|
||||||
|
'<span style="margin-left:auto;color:' + dot + ';">' + stat + '</span>' +
|
||||||
|
'</div>';
|
||||||
|
}).join('');
|
||||||
|
return '<div style="margin-bottom:10px;">' +
|
||||||
|
'<div style="font-weight:600;color:#AAB;margin-bottom:2px;">🖥️ ' + escapeHtml(node) + '</div>' +
|
||||||
|
rows + '</div>';
|
||||||
|
}).join('');
|
||||||
|
}
|
||||||
|
|
||||||
// ── Satelliten-Ansicht ─────────────────────────────────
|
// ── Satelliten-Ansicht ─────────────────────────────────
|
||||||
let satellites = []; // [{id, location, caps, control, online}]
|
let satellites = []; // [{id, location, caps, control, online}]
|
||||||
let satDevices = {}; // id → {devices, location, ts}
|
let satDevices = {}; // id → {devices, location, ts}
|
||||||
|
|||||||
@@ -490,6 +490,24 @@ function broadcastSatellites() {
|
|||||||
broadcast({ type: "sat_update", satellites: satelliteList() });
|
broadcast({ type: "sat_update", satellites: satelliteList() });
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ── Compute-Fleet: GPU-Worker (voxtral/whisper/f5tts/llm) ──────────
|
||||||
|
// Worker melden sich per worker_hello + halten sich per worker_ping (busy) frisch.
|
||||||
|
const workers = new Map(); // instanceId → {instanceId, service, node, gpus, model, busy, last_seen}
|
||||||
|
const WORKER_OFFLINE_MS = 35000; // ping ~10s; nach 35s ohne Ping = offline
|
||||||
|
|
||||||
|
function workerList() {
|
||||||
|
const now = Date.now();
|
||||||
|
return Array.from(workers.values()).map(w => ({
|
||||||
|
instanceId: w.instanceId, service: w.service, node: w.node,
|
||||||
|
gpus: w.gpus, model: w.model, busy: !!w.busy,
|
||||||
|
online: (now - (w.last_seen || 0)) < WORKER_OFFLINE_MS,
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
|
||||||
|
function broadcastWorkers() {
|
||||||
|
broadcast({ type: "worker_update", workers: workerList() });
|
||||||
|
}
|
||||||
|
|
||||||
// ── OpenClaw Gateway Verbindung ─────────────────────────
|
// ── OpenClaw Gateway Verbindung ─────────────────────────
|
||||||
|
|
||||||
async function connectGateway() {
|
async function connectGateway() {
|
||||||
@@ -935,6 +953,32 @@ function connectRVS(forcePlain) {
|
|||||||
if (p.satellite && satellites.has(p.satellite)) satellites.get(p.satellite).last_seen = Date.now();
|
if (p.satellite && satellites.has(p.satellite)) satellites.get(p.satellite).last_seen = Date.now();
|
||||||
broadcast({ type: "sat_devices", satellite: p.satellite || "",
|
broadcast({ type: "sat_devices", satellite: p.satellite || "",
|
||||||
location: p.location || "", devices: p.devices || [] });
|
location: p.location || "", devices: p.devices || [] });
|
||||||
|
} else if (msg.type === "worker_hello") {
|
||||||
|
// Ein Compute-Worker (GPU-Dienst) meldet sich mit seiner Identitaet.
|
||||||
|
const p = msg.payload || {};
|
||||||
|
if (p.instanceId) {
|
||||||
|
const prev = workers.get(p.instanceId) || {};
|
||||||
|
workers.set(p.instanceId, {
|
||||||
|
instanceId: p.instanceId, service: p.service || "",
|
||||||
|
node: p.node || "", gpus: p.gpus || "", model: p.model || "",
|
||||||
|
busy: !!prev.busy, last_seen: Date.now(),
|
||||||
|
});
|
||||||
|
broadcastWorkers();
|
||||||
|
}
|
||||||
|
} else if (msg.type === "worker_ping") {
|
||||||
|
// Heartbeat eines Workers (traegt busy-Status).
|
||||||
|
const p = msg.payload || {};
|
||||||
|
if (p.instanceId) {
|
||||||
|
let w = workers.get(p.instanceId);
|
||||||
|
if (!w) {
|
||||||
|
const svc = String(p.instanceId).split("@")[0];
|
||||||
|
w = { instanceId: p.instanceId, service: svc, node: "", gpus: "", model: "", busy: false, last_seen: 0 };
|
||||||
|
workers.set(p.instanceId, w);
|
||||||
|
}
|
||||||
|
w.busy = !!p.busy;
|
||||||
|
w.last_seen = Date.now();
|
||||||
|
broadcastWorkers();
|
||||||
|
}
|
||||||
} else if (msg.type === "agent_activity") {
|
} else if (msg.type === "agent_activity") {
|
||||||
// Bridge meldet "ARIA denkt/schreibt/tool" oder "idle" — an Browser
|
// Bridge meldet "ARIA denkt/schreibt/tool" oder "idle" — an Browser
|
||||||
// weiterreichen, damit der Thinking-Indikator im Chat erscheint.
|
// weiterreichen, damit der Thinking-Indikator im Chat erscheint.
|
||||||
@@ -2505,6 +2549,8 @@ wss.on("connection", (ws) => {
|
|||||||
if (currentDiskStatus) ws.send(JSON.stringify(currentDiskStatus));
|
if (currentDiskStatus) ws.send(JSON.stringify(currentDiskStatus));
|
||||||
// Aktuell bekannte Satelliten mitgeben (RVS replayt sat_hello nicht).
|
// Aktuell bekannte Satelliten mitgeben (RVS replayt sat_hello nicht).
|
||||||
ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() }));
|
ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() }));
|
||||||
|
// Aktuell bekannte Compute-Worker mitgeben (RVS replayt worker_hello nicht).
|
||||||
|
ws.send(JSON.stringify({ type: "worker_update", workers: workerList() }));
|
||||||
|
|
||||||
ws.on("message", (raw) => {
|
ws.on("message", (raw) => {
|
||||||
try {
|
try {
|
||||||
@@ -2536,6 +2582,9 @@ wss.on("connection", (ws) => {
|
|||||||
} else if (msg.action === "sat_list") {
|
} else if (msg.action === "sat_list") {
|
||||||
// Browser will die aktuelle Satelliten-Liste.
|
// Browser will die aktuelle Satelliten-Liste.
|
||||||
ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() }));
|
ws.send(JSON.stringify({ type: "sat_update", satellites: satelliteList() }));
|
||||||
|
} else if (msg.action === "worker_list") {
|
||||||
|
// Browser will die aktuelle Compute-Flotte.
|
||||||
|
ws.send(JSON.stringify({ type: "worker_update", workers: workerList() }));
|
||||||
} else if (msg.action === "sat_discover") {
|
} else if (msg.action === "sat_discover") {
|
||||||
// Browser triggert einen Geraete-Scan auf einem Satelliten.
|
// Browser triggert einen Geraete-Scan auf einem Satelliten.
|
||||||
sendToRVS_raw({ type: "sat_discover",
|
sendToRVS_raw({ type: "sat_discover",
|
||||||
|
|||||||
@@ -60,6 +60,15 @@ RVS_TOKEN = os.getenv("RVS_TOKEN", "").strip()
|
|||||||
# f5ttsCkptFile, f5ttsVocabFile, f5ttsCfgStrength, f5ttsNfeStep).
|
# f5ttsCkptFile, f5ttsVocabFile, f5ttsCfgStrength, f5ttsNfeStep).
|
||||||
F5TTS_DEVICE = os.getenv("F5TTS_DEVICE", "cuda") # nur Bootstrap
|
F5TTS_DEVICE = os.getenv("F5TTS_DEVICE", "cuda") # nur Bootstrap
|
||||||
|
|
||||||
|
# ── Compute-Fleet: Worker-Identitaet & Registrierung ──────────────
|
||||||
|
# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping.
|
||||||
|
NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node"
|
||||||
|
GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip()
|
||||||
|
WORKER_SERVICE = "f5tts"
|
||||||
|
INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}"
|
||||||
|
WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10"))
|
||||||
|
_tts_busy = False # True waehrend eine Synthese laeuft (busy-Report im ping)
|
||||||
|
|
||||||
DEFAULT_F5TTS_MODEL = "F5TTS_v1_Base"
|
DEFAULT_F5TTS_MODEL = "F5TTS_v1_Base"
|
||||||
DEFAULT_F5TTS_CKPT_FILE = "" # leer = Default-Checkpoint von HF
|
DEFAULT_F5TTS_CKPT_FILE = "" # leer = Default-Checkpoint von HF
|
||||||
DEFAULT_F5TTS_VOCAB_FILE = "" # leer = Default-Vocab vom Modell
|
DEFAULT_F5TTS_VOCAB_FILE = "" # leer = Default-Vocab vom Modell
|
||||||
@@ -460,13 +469,16 @@ _tts_queue: asyncio.Queue[tuple] = asyncio.Queue()
|
|||||||
|
|
||||||
async def _tts_worker(ws, runner: F5Runner) -> None:
|
async def _tts_worker(ws, runner: F5Runner) -> None:
|
||||||
"""Serialisiert Synthesen — GPU kann sonst OOM gehen."""
|
"""Serialisiert Synthesen — GPU kann sonst OOM gehen."""
|
||||||
|
global _tts_busy
|
||||||
while True:
|
while True:
|
||||||
text, voice, request_id, message_id, language, speed = await _tts_queue.get()
|
text, voice, request_id, message_id, language, speed = await _tts_queue.get()
|
||||||
|
_tts_busy = True
|
||||||
try:
|
try:
|
||||||
await _do_tts(ws, runner, text, voice, request_id, message_id, language, speed)
|
await _do_tts(ws, runner, text, voice, request_id, message_id, language, speed)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("TTS-Worker Fehler")
|
logger.exception("TTS-Worker Fehler")
|
||||||
finally:
|
finally:
|
||||||
|
_tts_busy = False
|
||||||
_tts_queue.task_done()
|
_tts_queue.task_done()
|
||||||
|
|
||||||
|
|
||||||
@@ -808,6 +820,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None:
|
|||||||
await _send(ws, "service_status", payload)
|
await _send(ws, "service_status", payload)
|
||||||
|
|
||||||
|
|
||||||
|
async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None:
|
||||||
|
"""Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die
|
||||||
|
Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch."""
|
||||||
|
try:
|
||||||
|
await _send(ws, "worker_hello", {
|
||||||
|
"instanceId": INSTANCE_ID, "service": WORKER_SERVICE,
|
||||||
|
"node": NODE_NAME, "gpus": GPU_IDS, "model": model,
|
||||||
|
})
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(WORKER_PING_INTERVAL_S)
|
||||||
|
busy = bool(busy_fn()) if busy_fn else False
|
||||||
|
await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy})
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
return # Socket tot → still beenden; run_loop reconnectet + startet neu
|
||||||
|
|
||||||
|
|
||||||
async def run_loop(runner: F5Runner) -> None:
|
async def run_loop(runner: F5Runner) -> None:
|
||||||
use_tls = RVS_TLS
|
use_tls = RVS_TLS
|
||||||
retry_s = 2
|
retry_s = 2
|
||||||
@@ -855,6 +885,9 @@ async def run_loop(runner: F5Runner) -> None:
|
|||||||
|
|
||||||
# TTS-Worker fuer diese Verbindung starten
|
# TTS-Worker fuer diese Verbindung starten
|
||||||
worker = asyncio.create_task(_tts_worker(ws, runner))
|
worker = asyncio.create_task(_tts_worker(ws, runner))
|
||||||
|
ping_task = asyncio.create_task(_worker_register(
|
||||||
|
ws, model=runner.model_id,
|
||||||
|
busy_fn=lambda: _tts_busy or not _tts_queue.empty()))
|
||||||
|
|
||||||
try:
|
try:
|
||||||
async for raw in ws:
|
async for raw in ws:
|
||||||
@@ -958,6 +991,7 @@ async def run_loop(runner: F5Runner) -> None:
|
|||||||
_last_diag_voice = ""
|
_last_diag_voice = ""
|
||||||
finally:
|
finally:
|
||||||
worker.cancel()
|
worker.cancel()
|
||||||
|
ping_task.cancel()
|
||||||
try:
|
try:
|
||||||
await worker
|
await worker
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
|||||||
@@ -47,6 +47,15 @@ RVS_TOKEN = os.getenv("RVS_TOKEN", "").strip()
|
|||||||
LLAMA_URL = os.getenv("LLAMA_URL", "http://llama:8081").rstrip("/")
|
LLAMA_URL = os.getenv("LLAMA_URL", "http://llama:8081").rstrip("/")
|
||||||
LLM_MODEL = os.getenv("LLM_MODEL", "qwen3-8b")
|
LLM_MODEL = os.getenv("LLM_MODEL", "qwen3-8b")
|
||||||
LLM_TIMEOUT_SEC = float(os.getenv("LLM_TIMEOUT_SEC", "60"))
|
LLM_TIMEOUT_SEC = float(os.getenv("LLM_TIMEOUT_SEC", "60"))
|
||||||
|
|
||||||
|
# ── Compute-Fleet: Worker-Identitaet & Registrierung ──────────────
|
||||||
|
# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping.
|
||||||
|
NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node"
|
||||||
|
GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip()
|
||||||
|
WORKER_SERVICE = "llm"
|
||||||
|
INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}"
|
||||||
|
WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10"))
|
||||||
|
_inflight = 0 # laufende llm_requests (busy-Report im ping)
|
||||||
# Qwen3 hat Thinking-Mode default AN — dann verbraet es Tokens in einem
|
# Qwen3 hat Thinking-Mode default AN — dann verbraet es Tokens in einem
|
||||||
# <think>-Block und liefert (bei kleinem max_tokens) leeren/abgeschnittenen
|
# <think>-Block und liefert (bei kleinem max_tokens) leeren/abgeschnittenen
|
||||||
# content, ausserdem 3x langsamer. ARIAs schnelles Tier will KEIN Grübeln
|
# content, ausserdem 3x langsamer. ARIAs schnelles Tier will KEIN Grübeln
|
||||||
@@ -124,7 +133,34 @@ async def _emit_llm_status(ws, state: str, model: str, **extra) -> None:
|
|||||||
{"service": "llm", "state": state, "model": model, **extra})
|
{"service": "llm", "state": state, "model": model, **extra})
|
||||||
|
|
||||||
|
|
||||||
|
async def _worker_register(ws) -> None:
|
||||||
|
"""Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die
|
||||||
|
Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch."""
|
||||||
|
try:
|
||||||
|
await _send(ws, "worker_hello", {
|
||||||
|
"instanceId": INSTANCE_ID, "service": WORKER_SERVICE,
|
||||||
|
"node": NODE_NAME, "gpus": GPU_IDS, "model": LLM_MODEL,
|
||||||
|
})
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(WORKER_PING_INTERVAL_S)
|
||||||
|
await _send(ws, "worker_ping",
|
||||||
|
{"instanceId": INSTANCE_ID, "busy": _inflight > 0})
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
return # Socket tot → still beenden; _run reconnectet + startet neu
|
||||||
|
|
||||||
|
|
||||||
async def _handle_llm_request(ws, payload: dict) -> None:
|
async def _handle_llm_request(ws, payload: dict) -> None:
|
||||||
|
global _last_model, _inflight
|
||||||
|
_inflight += 1
|
||||||
|
try:
|
||||||
|
await _do_llm_request(ws, payload)
|
||||||
|
finally:
|
||||||
|
_inflight -= 1
|
||||||
|
|
||||||
|
|
||||||
|
async def _do_llm_request(ws, payload: dict) -> None:
|
||||||
global _last_model
|
global _last_model
|
||||||
req_id = payload.get("requestId", "")
|
req_id = payload.get("requestId", "")
|
||||||
messages = payload.get("messages") or []
|
messages = payload.get("messages") or []
|
||||||
@@ -201,6 +237,7 @@ async def _run() -> None:
|
|||||||
logger.info("RVS verbunden — llm-adapter online")
|
logger.info("RVS verbunden — llm-adapter online")
|
||||||
retry_s = 2
|
retry_s = 2
|
||||||
tls_fallback_tried = False
|
tls_fallback_tried = False
|
||||||
|
ping_task = asyncio.create_task(_worker_register(ws))
|
||||||
async for raw in ws:
|
async for raw in ws:
|
||||||
try:
|
try:
|
||||||
msg = json.loads(raw)
|
msg = json.loads(raw)
|
||||||
@@ -214,6 +251,10 @@ async def _run() -> None:
|
|||||||
asyncio.create_task(_handle_llm_request(ws, payload))
|
asyncio.create_task(_handle_llm_request(ws, payload))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("RVS-Verbindung verloren/fehlgeschlagen: %s", e)
|
logger.warning("RVS-Verbindung verloren/fehlgeschlagen: %s", e)
|
||||||
|
try:
|
||||||
|
ping_task.cancel()
|
||||||
|
except NameError:
|
||||||
|
pass
|
||||||
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
||||||
tls_fallback_tried = True
|
tls_fallback_tried = True
|
||||||
use_tls = False
|
use_tls = False
|
||||||
|
|||||||
@@ -58,6 +58,16 @@ VOXTRAL_MODEL = os.getenv("VOXTRAL_MODEL", "mistralai/Voxtral-Mini-3B-2507")
|
|||||||
VOXTRAL_LANGUAGE = os.getenv("VOXTRAL_LANGUAGE", "de")
|
VOXTRAL_LANGUAGE = os.getenv("VOXTRAL_LANGUAGE", "de")
|
||||||
VOXTRAL_DEVICE = os.getenv("VOXTRAL_DEVICE", "cuda")
|
VOXTRAL_DEVICE = os.getenv("VOXTRAL_DEVICE", "cuda")
|
||||||
|
|
||||||
|
# ── Compute-Fleet: Worker-Identitaet & Registrierung ──────────────
|
||||||
|
# Jeder Node meldet sich bei der aria-bridge (worker_hello) und haelt die
|
||||||
|
# Registry per periodischem worker_ping frisch. INSTANCE_ID adressiert diesen
|
||||||
|
# Worker bei Redundanz (targetInstance-Routing, Stage 3).
|
||||||
|
NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node"
|
||||||
|
GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip()
|
||||||
|
WORKER_SERVICE = "voxtral"
|
||||||
|
INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}"
|
||||||
|
WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10"))
|
||||||
|
|
||||||
STREAM_TRANSCRIBE_INTERVAL_MS = int(os.getenv("STREAM_TRANSCRIBE_INTERVAL_MS", "1000"))
|
STREAM_TRANSCRIBE_INTERVAL_MS = int(os.getenv("STREAM_TRANSCRIBE_INTERVAL_MS", "1000"))
|
||||||
STREAM_DEFAULT_ENDPOINT_MS = 2400
|
STREAM_DEFAULT_ENDPOINT_MS = 2400
|
||||||
STREAM_DEFAULT_HARD_CAP_MS = 300000
|
STREAM_DEFAULT_HARD_CAP_MS = 300000
|
||||||
@@ -695,6 +705,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None:
|
|||||||
await _send(ws, "service_status", payload)
|
await _send(ws, "service_status", payload)
|
||||||
|
|
||||||
|
|
||||||
|
async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None:
|
||||||
|
"""Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die
|
||||||
|
Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch."""
|
||||||
|
try:
|
||||||
|
await _send(ws, "worker_hello", {
|
||||||
|
"instanceId": INSTANCE_ID, "service": WORKER_SERVICE,
|
||||||
|
"node": NODE_NAME, "gpus": GPU_IDS, "model": model,
|
||||||
|
})
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(WORKER_PING_INTERVAL_S)
|
||||||
|
busy = bool(busy_fn()) if busy_fn else False
|
||||||
|
await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy})
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
return # Socket tot → still beenden; run_loop reconnectet + startet neu
|
||||||
|
|
||||||
|
|
||||||
async def run_loop(sessions: SessionManager) -> None:
|
async def run_loop(sessions: SessionManager) -> None:
|
||||||
use_tls = RVS_TLS
|
use_tls = RVS_TLS
|
||||||
retry_s = 2
|
retry_s = 2
|
||||||
@@ -713,6 +741,9 @@ async def run_loop(sessions: SessionManager) -> None:
|
|||||||
sessions.attach_ws(ws)
|
sessions.attach_ws(ws)
|
||||||
await _broadcast_status(ws, "ready", model=VOXTRAL_MODEL)
|
await _broadcast_status(ws, "ready", model=VOXTRAL_MODEL)
|
||||||
await _send(ws, "config_request", {"service": "voxtral"})
|
await _send(ws, "config_request", {"service": "voxtral"})
|
||||||
|
ping_task = asyncio.create_task(_worker_register(
|
||||||
|
ws, model=VOXTRAL_MODEL,
|
||||||
|
busy_fn=lambda: bool(sessions._sessions)))
|
||||||
async for raw in ws:
|
async for raw in ws:
|
||||||
try:
|
try:
|
||||||
msg = json.loads(raw)
|
msg = json.loads(raw)
|
||||||
@@ -797,6 +828,10 @@ async def run_loop(sessions: SessionManager) -> None:
|
|||||||
"AN" if SPEAKER_ID_ENABLED else "AUS")
|
"AN" if SPEAKER_ID_ENABLED else "AUS")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("RVS-Verbindung verloren: %s — retry in %ds", e, retry_s)
|
logger.warning("RVS-Verbindung verloren: %s — retry in %ds", e, retry_s)
|
||||||
|
try:
|
||||||
|
ping_task.cancel()
|
||||||
|
except NameError:
|
||||||
|
pass
|
||||||
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
||||||
use_tls = False
|
use_tls = False
|
||||||
tls_fallback_tried = True
|
tls_fallback_tried = True
|
||||||
|
|||||||
@@ -59,6 +59,14 @@ WHISPER_DEVICE = os.getenv("WHISPER_DEVICE", "cuda")
|
|||||||
WHISPER_COMPUTE_TYPE = os.getenv("WHISPER_COMPUTE_TYPE", "float16")
|
WHISPER_COMPUTE_TYPE = os.getenv("WHISPER_COMPUTE_TYPE", "float16")
|
||||||
WHISPER_LANGUAGE = os.getenv("WHISPER_LANGUAGE", "de")
|
WHISPER_LANGUAGE = os.getenv("WHISPER_LANGUAGE", "de")
|
||||||
|
|
||||||
|
# ── Compute-Fleet: Worker-Identitaet & Registrierung ──────────────
|
||||||
|
# Meldet sich bei der aria-bridge (worker_hello) + periodischer worker_ping.
|
||||||
|
NODE_NAME = os.getenv("NODE_NAME", "node").strip() or "node"
|
||||||
|
GPU_IDS = os.getenv("NVIDIA_VISIBLE_DEVICES", "").strip()
|
||||||
|
WORKER_SERVICE = "whisper"
|
||||||
|
INSTANCE_ID = f"{WORKER_SERVICE}@{NODE_NAME}"
|
||||||
|
WORKER_PING_INTERVAL_S = int(os.getenv("WORKER_PING_INTERVAL_S", "10"))
|
||||||
|
|
||||||
ALLOWED_MODELS = {"tiny", "base", "small", "medium", "large-v3"}
|
ALLOWED_MODELS = {"tiny", "base", "small", "medium", "large-v3"}
|
||||||
|
|
||||||
# Streaming-Parameter (Defaults — koennen pro Session vom App-Payload ueberschrieben werden)
|
# Streaming-Parameter (Defaults — koennen pro Session vom App-Payload ueberschrieben werden)
|
||||||
@@ -822,6 +830,24 @@ async def _broadcast_status(ws, state: str, **extra) -> None:
|
|||||||
await _send(ws, "service_status", payload)
|
await _send(ws, "service_status", payload)
|
||||||
|
|
||||||
|
|
||||||
|
async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None:
|
||||||
|
"""Meldet diesen Worker bei der aria-bridge an (worker_hello) und haelt die
|
||||||
|
Flotten-Registry per periodischem worker_ping (mit busy-Status) frisch."""
|
||||||
|
try:
|
||||||
|
await _send(ws, "worker_hello", {
|
||||||
|
"instanceId": INSTANCE_ID, "service": WORKER_SERVICE,
|
||||||
|
"node": NODE_NAME, "gpus": GPU_IDS, "model": model,
|
||||||
|
})
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(WORKER_PING_INTERVAL_S)
|
||||||
|
busy = bool(busy_fn()) if busy_fn else False
|
||||||
|
await _send(ws, "worker_ping", {"instanceId": INSTANCE_ID, "busy": busy})
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
return # Socket tot → still beenden; run_loop reconnectet + startet neu
|
||||||
|
|
||||||
|
|
||||||
# ──────────────────────────────────────────────────────────────
|
# ──────────────────────────────────────────────────────────────
|
||||||
# WS-LOOP
|
# WS-LOOP
|
||||||
# ──────────────────────────────────────────────────────────────
|
# ──────────────────────────────────────────────────────────────
|
||||||
@@ -862,6 +888,9 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None:
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("Initial-Handshake crashed: %s", e)
|
logger.exception("Initial-Handshake crashed: %s", e)
|
||||||
asyncio.create_task(_initial_handshake())
|
asyncio.create_task(_initial_handshake())
|
||||||
|
ping_task = asyncio.create_task(_worker_register(
|
||||||
|
ws, model=(runner.model_size or WHISPER_MODEL),
|
||||||
|
busy_fn=lambda: bool(sessions._sessions)))
|
||||||
|
|
||||||
async for raw in ws:
|
async for raw in ws:
|
||||||
try:
|
try:
|
||||||
@@ -1038,6 +1067,10 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None:
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning("Verbindung verloren: %s", e)
|
logger.warning("Verbindung verloren: %s", e)
|
||||||
sessions.detach_ws()
|
sessions.detach_ws()
|
||||||
|
try:
|
||||||
|
ping_task.cancel()
|
||||||
|
except NameError:
|
||||||
|
pass
|
||||||
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
if use_tls and RVS_TLS_FALLBACK and not tls_fallback_tried:
|
||||||
logger.info("TLS-Verbindung fehlgeschlagen — Fallback auf ws://")
|
logger.info("TLS-Verbindung fehlgeschlagen — Fallback auf ws://")
|
||||||
use_tls = False
|
use_tls = False
|
||||||
|
|||||||
Reference in New Issue
Block a user