diff --git a/diagnostic/server.js b/diagnostic/server.js index 3c53c0f..1ca7722 100644 --- a/diagnostic/server.js +++ b/diagnostic/server.js @@ -508,6 +508,64 @@ function broadcastWorkers() { broadcast({ type: "worker_update", workers: workerList() }); } +// ── Voice-Flotte: zentraler Stimmen-Store + Auto-Provisioning ────── +// Der Diagnostic-Server ist der "Stimmen-Bibliothekar": Stimmen liegen +// zentral als /shared/voices/{name}.tar.gz (das f5tts-Export-Artefakt = +// wav+txt). Meldet sich eine f5tts-Box (worker_hello), gleichen wir ab und +// schieben ihr fehlende Stimmen (xtts_import_voice, targetInstance) bzw. ziehen +// bei ihr vorhandene, zentral fehlende Stimmen (xtts_export_voice) in den Store. +const CENTRAL_VOICES_DIR = "/shared/voices"; +const centralExportPending = new Map(); // requestId -> name (unsere eigenen Export-Anfragen) + +function ensureCentralVoicesDir() { + try { fs.mkdirSync(CENTRAL_VOICES_DIR, { recursive: true }); } catch (_) {} +} +function centralVoiceNames() { + ensureCentralVoicesDir(); + try { + return fs.readdirSync(CENTRAL_VOICES_DIR) + .filter(f => f.endsWith(".tar.gz")) + .map(f => f.slice(0, -7)); + } catch (_) { return []; } +} +function readCentralVoiceB64(name) { + try { return fs.readFileSync(`${CENTRAL_VOICES_DIR}/${name}.tar.gz`).toString("base64"); } + catch (_) { return null; } +} +function writeCentralVoice(name, dataB64) { + ensureCentralVoicesDir(); + try { fs.writeFileSync(`${CENTRAL_VOICES_DIR}/${name}.tar.gz`, Buffer.from(dataB64, "base64")); return true; } + catch (e) { log("warn", "voice", `zentral schreiben ${name} fehlgeschlagen: ${e.message}`); return false; } +} +function deleteCentralVoice(name) { + try { fs.unlinkSync(`${CENTRAL_VOICES_DIR}/${name}.tar.gz`); log("info", "voice", `zentral geloescht: ${name}`); } + catch (_) {} +} +function requestCentralExport(name) { + // Broadcast-Export-Anfrage; die Box, die die Stimme hat, antwortet. Wir + // korrelieren die Antwort ueber requestId (nur unsere eigenen verarbeiten). + const requestId = "central_" + Date.now() + "_" + Math.random().toString(36).slice(2, 8); + centralExportPending.set(requestId, name); + setTimeout(() => centralExportPending.delete(requestId), 30000); + sendToRVS_raw({ type: "xtts_export_voice", payload: { name, requestId }, timestamp: Date.now() }); +} +function provisionVoiceToInstance(name, instanceId) { + const data = readCentralVoiceB64(name); + if (!data) return; + sendToRVS_raw({ type: "xtts_import_voice", + payload: { name, data, targetInstance: instanceId }, timestamp: Date.now() }); + log("info", "voice", `provisioniere '${name}' → ${instanceId}`); +} +// Abgleich beim worker_hello einer f5tts-Box: push (zentral→Box) + pull (Box→zentral, seed). +function reconcileVoices(instanceId, boxVoices) { + const central = centralVoiceNames(); + const boxSet = new Set(Array.isArray(boxVoices) ? boxVoices : []); + const centralSet = new Set(central); + for (const name of central) if (!boxSet.has(name)) provisionVoiceToInstance(name, instanceId); + for (const name of boxSet) if (!centralSet.has(name)) requestCentralExport(name); + log("info", "voice", `reconcile ${instanceId}: box=${boxSet.size} central=${central.length}`); +} + // ── OpenClaw Gateway Verbindung ───────────────────────── async function connectGateway() { @@ -964,6 +1022,8 @@ function connectRVS(forcePlain) { busy: !!prev.busy, last_seen: Date.now(), }); broadcastWorkers(); + // Voice-Flotte: f5tts-Box gerade online → Stimmen abgleichen/provisionieren. + if ((p.service || "") === "f5tts") reconcileVoices(p.instanceId, p.voices); } } else if (msg.type === "worker_ping") { // Heartbeat eines Workers (traegt busy-Status). @@ -979,6 +1039,30 @@ function connectRVS(forcePlain) { w.last_seen = Date.now(); broadcastWorkers(); } + } else if (msg.type === "xtts_voice_saved") { + // Neue Stimme (App- ODER Diagnostic-Upload, via RVS-Broadcast) → zentral + // sichern. Online-Boxen haben sie durch den voice_upload-Broadcast schon; + // der zentrale Store macht sie persistent + fuer spaeter joinende Boxen + // verfuegbar (die holt dann reconcileVoices ab). + const p = msg.payload || {}; + if (p.name && !p.error) { + log("info", "voice", `Stimme '${p.name}' gespeichert → zentraler Ingest`); + requestCentralExport(p.name); + } + } else if (msg.type === "xtts_voice_exported") { + // Antwort auf eine UNSERER zentralen Export-Anfragen (requestId-Match) → + // in den zentralen Store schreiben. Browser-initiierte Exports tragen + // keinen centralExportPending-requestId und werden hier ignoriert. + const p = msg.payload || {}; + const rid = p.requestId || ""; + if (rid && centralExportPending.has(rid)) { + centralExportPending.delete(rid); + if (p.ok && p.name && p.data) writeCentralVoice(p.name, p.data); + } + } else if (msg.type === "xtts_delete_voice") { + // App-initiierter Delete (Broadcast) → zentrale Kopie mitloeschen. + const p = msg.payload || {}; + if (p.name) deleteCentralVoice(p.name); } else if (msg.type === "agent_activity") { // Bridge meldet "ARIA denkt/schreibt/tool" oder "idle" — an Browser // weiterreichen, damit der Thinking-Indikator im Chat erscheint. @@ -2660,9 +2744,11 @@ wss.on("connection", (ws) => { // tar.gz (base64) an XTTS-Bridge schicken — die packt aus sendToRVS_withResponse("xtts_import_voice", { name: msg.name, data: msg.data }, "xtts_voice_imported", ws); } else if (msg.action === "xtts_delete_voice") { - // Weiterleiten an XTTS-Bridge, die antwortet mit neuer Liste + // Weiterleiten an alle f5tts-Boxen (Broadcast) + zentrale Kopie loeschen. + // (Der eigene Broadcast kommt nicht zu uns zurueck, daher hier direkt.) sendToRVS_raw({ type: "xtts_delete_voice", payload: { name: msg.name }, timestamp: Date.now() }); - log("info", "server", `Voice-Delete '${msg.name}' an XTTS-Bridge gesendet`); + if (msg.name) deleteCentralVoice(msg.name); + log("info", "server", `Voice-Delete '${msg.name}' an f5tts-Boxen + zentral geloescht`); } else if (msg.action === "delete_chat_message") { // Bubble loeschen — Bridge raeumt chat_backup.jsonl + Brain-conversation // + broadcastet chat_message_deleted via RVS. diff --git a/xtts/f5tts/bridge.py b/xtts/f5tts/bridge.py index 3d0ee9f..768b7f7 100644 --- a/xtts/f5tts/bridge.py +++ b/xtts/f5tts/bridge.py @@ -675,6 +675,18 @@ async def handle_voice_upload(ws, payload: dict) -> None: await _send(ws, "xtts_voice_saved", {"name": name, "error": str(e)[:200]}) +def _local_voice_names() -> list: + """Namen der lokal vorhandenen Nutzer-Stimmen (ohne den Box-lokalen + default_ref-Fallback) — fuer die Provisioning-Reconciliation im worker_hello.""" + names = [] + if VOICES_DIR.exists(): + for wav in sorted(VOICES_DIR.glob("*.wav")): + if wav.stem == "default_ref": + continue + names.append(wav.stem) + return names + + async def handle_list_voices(ws) -> None: try: voices = [] @@ -711,13 +723,16 @@ async def handle_delete_voice(ws, payload: dict) -> None: async def handle_export_voice(ws, payload: dict) -> None: """Packt eine Stimme (.wav + .txt) als tar.gz und sendet sie base64 zurueck.""" name = (payload.get("name") or "").strip() + # requestId aus dem Request zuruueckspiegeln — der Diagnostic-Bibliothekar + # korreliert damit zentrale Exports gegen browser-initiierte. + req_id = payload.get("requestId", "") or "" if not name: - await _send(ws, "xtts_voice_exported", {"ok": False, "error": "name fehlt"}) + await _send(ws, "xtts_voice_exported", {"ok": False, "requestId": req_id, "error": "name fehlt"}) return try: wav, txt = voice_paths(name) if not wav.exists(): - await _send(ws, "xtts_voice_exported", {"ok": False, "name": name, "error": "Stimme nicht gefunden"}) + await _send(ws, "xtts_voice_exported", {"ok": False, "requestId": req_id, "name": name, "error": "Stimme nicht gefunden"}) return import io, tarfile buf = io.BytesIO() @@ -727,10 +742,10 @@ async def handle_export_voice(ws, payload: dict) -> None: tar.add(txt, arcname=txt.name) data = base64.b64encode(buf.getvalue()).decode("ascii") logger.info("Voice exportiert: %s (%d KB tar.gz)", name, len(buf.getvalue()) // 1024) - await _send(ws, "xtts_voice_exported", {"ok": True, "name": name, "data": data}) + await _send(ws, "xtts_voice_exported", {"ok": True, "requestId": req_id, "name": name, "data": data}) except Exception as e: logger.exception("handle_export_voice Fehler") - await _send(ws, "xtts_voice_exported", {"ok": False, "name": name, "error": str(e)[:200]}) + await _send(ws, "xtts_voice_exported", {"ok": False, "requestId": req_id, "name": name, "error": str(e)[:200]}) async def handle_import_voice(ws, payload: dict) -> None: @@ -827,6 +842,7 @@ async def _worker_register(ws, *, model: str = "", busy_fn=None) -> None: await _send(ws, "worker_hello", { "instanceId": INSTANCE_ID, "service": WORKER_SERVICE, "node": NODE_NAME, "gpus": GPU_IDS, "model": model, + "voices": _local_voice_names(), # fuer Voice-Provisioning-Reconciliation }) while True: await asyncio.sleep(WORKER_PING_INTERVAL_S)