From ec78eb8efe9489adc97bcfa1bae00ff374d950f7 Mon Sep 17 00:00:00 2001 From: duffyduck Date: Fri, 18 Sep 2026 12:27:51 +0200 Subject: [PATCH] feat(compute): Redundanz-Routing per targetInstance (Stage 3) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Mehrere Instanzen pro Dienst nutzbar: Anfragen werden gezielt an eine freie Instanz adressiert statt an alle gebroadcastet. Mehrere f5tts → naechstes freies; mehrere LLM → parallele Turns (Multitasking); STT-Redundanz ueber mehrere Apps via Lease. - Worker (alle vier): filtern am Loop-Eingang — targetInstance gesetzt und != eigener INSTANCE_ID → Nachricht ignorieren. Feld fehlt → wie bisher. - bridge (TTS/LLM, emittiert die Bridge selbst): _pick_worker() waehlt eine online+freie Instanz (Round-Robin), stempelt targetInstance auf xtts_request / llm_request. Keine Instanz bekannt → Broadcast. - bridge (STT-Lease, emittiert die App): neuer stt_lease_request-Handler → _pick_stt_worker() (voxtral vor whisper) → stt_lease {instanceId}. - app (audio.ts): requestSttLease() vor dem Stream, stempelt targetInstance auf stt_stream_start / stt_audio_chunk / stt_stream_end (+cancel). Kurzer Timeout → '' (Broadcast), Aufnahme haengt nie. Voll rueckwaertskompatibel: Routing aktiviert sich erst, wenn Worker sich per worker_hello (Stage 2) registriert haben — sonst bleibt alles Broadcast. Braucht APK-Rebuild fuer die STT-Lease-Seite. Co-Authored-By: Claude Opus 4.8 --- android/src/services/audio.ts | 51 +++++++++++++++++++++++++++-- bridge/aria_bridge.py | 60 ++++++++++++++++++++++++++++------- xtts/f5tts/bridge.py | 5 +++ xtts/llm-adapter/adapter.py | 5 +++ xtts/voxtral/bridge.py | 6 ++++ xtts/whisper/bridge.py | 5 +++ 6 files changed, 118 insertions(+), 14 deletions(-) diff --git a/android/src/services/audio.ts b/android/src/services/audio.ts index 3f38f6a..0f5a1fb 100644 --- a/android/src/services/audio.ts +++ b/android/src/services/audio.ts @@ -204,6 +204,39 @@ export async function transcribeBlob(pcmBase64: string, timeoutMs = 2500): Promi }); } +/** Fragt die Bridge vor dem Aufnahme-Stream, welche STT-Instanz adressiert + * werden soll (Redundanz ueber mehrere STT-Nodes/Apps). Schickt + * stt_lease_request, wartet kurz auf stt_lease (matching requestId). + * Rueckgabe: instanceId (z.B. "voxtral@box-a") oder '' bei Timeout/keine + * Instanz — dann streamt die App wie bisher an ALLE (Broadcast, Single-Node + * unveraendert). Bewusst kurzer Timeout, damit die Aufnahme nie haengt. */ +export async function requestSttLease(timeoutMs = 250): Promise { + const requestId = `sttlease_${Date.now()}_${Math.floor(Math.random() * 100000)}`; + return new Promise((resolve) => { + let done = false; + let unsub: (() => void) | null = null; + const timer = setTimeout(() => finish(''), timeoutMs); + function finish(val: string) { + if (done) return; + done = true; + try { unsub && unsub(); } catch {} + clearTimeout(timer); + resolve(val); + } + try { + unsub = rvs.onMessage((msg: any) => { + if (msg?.type !== 'stt_lease') return; + const p = (msg as any).payload || {}; + if (String(p.requestId || '') !== requestId) return; + finish(typeof p.instanceId === 'string' ? p.instanceId : ''); + }); + rvs.send('stt_lease_request' as any, { requestId }); + } catch { + finish(''); + } + }); +} + export async function loadSttEndpointMs(): Promise { try { const raw = await AsyncStorage.getItem(STT_ENDPOINT_STORAGE_KEY); @@ -383,6 +416,9 @@ class AudioService { // lich Chunks einer alten Session in eine neue mischen. private streamRequestId: string = ''; private streamAudioRequestId: string = ''; + // Adressierte STT-Instanz fuer diesen Stream (Redundanz-Routing). '' = + // Broadcast an alle STT-Nodes (Single-Node / kein Lease = wie bisher). + private streamTargetInstance: string = ''; // Latch: ist endpointListeners fuer den aktuellen Session-Cycle schon gefeuert // worden? Wird auf false gesetzt beim startStreamingRecording, auf true beim // ersten Endpoint (egal ob via RVS oder Fallback). Verhindert Doppel-Fires. @@ -1126,6 +1162,15 @@ class AudioService { const requestId = `sttstr_${Date.now()}_${Math.floor(Math.random() * 100000)}`; this.streamRequestId = requestId; this.streamAudioRequestId = opts.audioRequestId || ''; + // Redundanz-Routing: freie STT-Instanz leasen BEVOR Chunks fliessen, damit + // start + alle Chunks + end dieselbe Instanz adressieren. Kurzer Timeout → + // '' (Broadcast) falls keine Instanz/keine Antwort. Nie blockierend genug + // um die Aufnahme spuerbar zu verzoegern. + try { + this.streamTargetInstance = await requestSttLease(); + } catch { + this.streamTargetInstance = ''; + } this.streamGotPartial = false; this.streamEndpointFired = false; this.recordingStartTime = Date.now(); @@ -1146,6 +1191,7 @@ class AudioService { requestId: sessionId, pcm: String(e?.pcm || ''), seq: Number(e?.seq || 0), + targetInstance: this.streamTargetInstance, }); }); this.streamPcmErrorSub = emitter.addListener('PcmStreamError', (e: any) => { @@ -1181,6 +1227,7 @@ class AudioService { hardCapMs: typeof opts.hardCapMs === 'number' ? opts.hardCapMs : 60000, sampleRate: 16000, projectId: opts.projectId || '', + targetInstance: this.streamTargetInstance, }); // No-Speech-Watchdog — ersetzt den alten VAD-noSpeechTimer. @@ -1230,7 +1277,7 @@ class AudioService { if (!reqId) return; const audioReqId = this.streamAudioRequestId; try { - rvs.send('stt_stream_end' as any, { requestId: reqId, reason }); + rvs.send('stt_stream_end' as any, { requestId: reqId, reason, targetInstance: this.streamTargetInstance }); } catch (e) { console.warn('[Audio] stt_stream_end senden fehlgeschlagen:', e); } @@ -1265,7 +1312,7 @@ class AudioService { if (!reqId) return; const audioReqId = this.streamAudioRequestId; try { - rvs.send('stt_stream_end' as any, { requestId: reqId, reason: `cancel:${reason}` }); + rvs.send('stt_stream_end' as any, { requestId: reqId, reason: `cancel:${reason}`, targetInstance: this.streamTargetInstance }); } catch {} this._cleanupStreamLocal(`cancel:${reason}`); // Listener feuern damit ChatScreen reagieren kann (endConversation etc.) diff --git a/bridge/aria_bridge.py b/bridge/aria_bridge.py index 7a935f1..ed40a63 100644 --- a/bridge/aria_bridge.py +++ b/bridge/aria_bridge.py @@ -1684,20 +1684,27 @@ class ARIABridge: if len(self._xtts_request_to_message) > 100: oldest = next(iter(self._xtts_request_to_message)) self._xtts_request_to_message.pop(oldest, None) + # Redundanz: freie f5tts-Instanz waehlen und gezielt adressieren. + # None (keine Instanz bekannt / alle offline) → kein targetInstance, + # Broadcast wie bisher (Single-Node laeuft unveraendert). + tts_target = self._pick_worker("f5tts") + tts_payload = { + "text": tts_text, + "voice": xtts_voice, + "speed": xtts_speed, + "language": "de", + "requestId": xtts_request_id, + "messageId": message_id, + } + if tts_target: + tts_payload["targetInstance"] = tts_target await self._send_to_rvs({ "type": "xtts_request", - "payload": { - "text": tts_text, - "voice": xtts_voice, - "speed": xtts_speed, - "language": "de", - "requestId": xtts_request_id, - "messageId": message_id, - }, + "payload": tts_payload, "timestamp": int(asyncio.get_event_loop().time() * 1000), }) - logger.info("[core] XTTS-Request gesendet (voice=%s, speed=%.2fx): '%s'", - xtts_voice or "default", xtts_speed, tts_text[:60]) + logger.info("[core] XTTS-Request gesendet (voice=%s, speed=%.2fx, target=%s): '%s'", + xtts_voice or "default", xtts_speed, tts_target or "(broadcast)", tts_text[:60]) except Exception as e: logger.error("[core] XTTS-Request fehlgeschlagen: %s — kein Audio", e) @@ -3551,6 +3558,23 @@ class ARIABridge: self._workers[iid]["gpus"] or "?", self._workers[iid]["model"] or "?") return + elif msg_type == "stt_lease_request": + # Die App fragt vor dem Aufnahme-Stream, welche STT-Instanz sie + # adressieren soll (Redundanz ueber mehrere Apps/Nodes). Wir waehlen + # eine freie STT-Instanz und antworten per stt_lease. Ist keine + # Instanz bekannt (instanceId leer), streamt die App wie bisher an + # ALLE (Broadcast) — Single-Node bleibt unveraendert. + req_id = (payload.get("requestId") or "").strip() + iid = self._pick_stt_worker() or "" + await self._send_to_rvs({ + "type": "stt_lease", + "payload": {"requestId": req_id, "instanceId": iid}, + "timestamp": int(time.time() * 1000), + }) + logger.info("[stt-lease] req=%s → %s", req_id[:8] if req_id else "?", + iid or "(broadcast)") + return + elif msg_type == "worker_ping": iid = (payload.get("instanceId") or "").strip() if iid: @@ -3982,8 +4006,14 @@ class ARIABridge: req_payload["tools"] = tools if model: req_payload["model"] = model - logger.info("[rvs] llm_request → llm-adapter (id=%s, msgs=%d, max_tokens=%d, tools=%d, model=%s)", - request_id[:8], len(messages), max_tokens, len(tools) if tools else 0, model or "-") + # Redundanz/Multitasking: freie llm-Instanz gezielt adressieren; None + # → Broadcast wie bisher. Mehrere Instanzen → parallele Turns. + llm_target = self._pick_worker("llm") + if llm_target: + req_payload["targetInstance"] = llm_target + logger.info("[rvs] llm_request → llm-adapter (id=%s, msgs=%d, max_tokens=%d, tools=%d, model=%s, target=%s)", + request_id[:8], len(messages), max_tokens, len(tools) if tools else 0, + model or "-", llm_target or "(broadcast)") ok = await self._send_to_rvs({ "type": "llm_request", "payload": req_payload, @@ -4764,6 +4794,12 @@ class ARIABridge: chosen["busy"] = True # optimistisch, bis der naechste Ping korrigiert return chosen["instanceId"] + def _pick_stt_worker(self) -> Optional[str]: + """Waehlt eine STT-Instanz fuer ein App-Lease. Voxtral (Default-STT) hat + Vorrang, Whisper ist der Fallback. None → keine online (App streamt dann + ohne targetInstance = heutiges Broadcast-Verhalten).""" + return self._pick_worker("voxtral") or self._pick_worker("whisper") + async def _satellite_request(self, op: str, satellite: str = "", device: str = "", action: str = "", params: Optional[dict] = None, diff --git a/xtts/f5tts/bridge.py b/xtts/f5tts/bridge.py index b1bc087..dfd0cbc 100644 --- a/xtts/f5tts/bridge.py +++ b/xtts/f5tts/bridge.py @@ -897,6 +897,11 @@ async def run_loop(runner: F5Runner) -> None: continue mtype = msg.get("type", "") payload = msg.get("payload", {}) or {} + # Redundanz-Routing: gezielt an eine andere Instanz + # adressiert → ignorieren. Ohne targetInstance → wie bisher. + tgt = payload.get("targetInstance") + if tgt and tgt != INSTANCE_ID: + continue if mtype == "xtts_request": try: diff --git a/xtts/llm-adapter/adapter.py b/xtts/llm-adapter/adapter.py index bdd157a..3591bed 100644 --- a/xtts/llm-adapter/adapter.py +++ b/xtts/llm-adapter/adapter.py @@ -246,6 +246,11 @@ async def _run() -> None: if msg.get("type") != "llm_request": continue payload = msg.get("payload", {}) or {} + # Redundanz-Routing: gezielt an eine andere Instanz adressiert + # → ignorieren. Ohne targetInstance → wie bisher (jeder nimmt). + tgt = payload.get("targetInstance") + if tgt and tgt != INSTANCE_ID: + continue # Jede Anfrage nebenlaeufig — llama.cpp serialisiert intern, # aber wir blockieren so nicht den Empfang weiterer Messages. asyncio.create_task(_handle_llm_request(ws, payload)) diff --git a/xtts/voxtral/bridge.py b/xtts/voxtral/bridge.py index c3ae726..f1cc02b 100644 --- a/xtts/voxtral/bridge.py +++ b/xtts/voxtral/bridge.py @@ -751,6 +751,12 @@ async def run_loop(sessions: SessionManager) -> None: continue mtype = msg.get("type", "") payload = msg.get("payload", {}) or {} + # Redundanz-Routing: ist die Anfrage gezielt an eine andere + # Instanz adressiert, ignorieren. Ohne targetInstance (Feld + # fehlt) → wie bisher, jeder Worker nimmt sie an. + tgt = payload.get("targetInstance") + if tgt and tgt != INSTANCE_ID: + continue if mtype == "stt_stream_start": sessions.start_session(payload) elif mtype == "stt_audio_chunk": diff --git a/xtts/whisper/bridge.py b/xtts/whisper/bridge.py index 7a5d41e..32bc382 100644 --- a/xtts/whisper/bridge.py +++ b/xtts/whisper/bridge.py @@ -899,6 +899,11 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None: continue mtype = msg.get("type", "") payload = msg.get("payload", {}) or {} + # Redundanz-Routing: gezielt an eine andere Instanz adressiert + # → ignorieren. Ohne targetInstance → wie bisher. + tgt = payload.get("targetInstance") + if tgt and tgt != INSTANCE_ID: + continue if mtype == "stt_request": req_id = payload.get("requestId", "?")