feat(compute): Redundanz-Routing per targetInstance (Stage 3)
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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<string> {
|
||||||
|
const requestId = `sttlease_${Date.now()}_${Math.floor(Math.random() * 100000)}`;
|
||||||
|
return new Promise<string>((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<number> {
|
export async function loadSttEndpointMs(): Promise<number> {
|
||||||
try {
|
try {
|
||||||
const raw = await AsyncStorage.getItem(STT_ENDPOINT_STORAGE_KEY);
|
const raw = await AsyncStorage.getItem(STT_ENDPOINT_STORAGE_KEY);
|
||||||
@@ -383,6 +416,9 @@ class AudioService {
|
|||||||
// lich Chunks einer alten Session in eine neue mischen.
|
// lich Chunks einer alten Session in eine neue mischen.
|
||||||
private streamRequestId: string = '';
|
private streamRequestId: string = '';
|
||||||
private streamAudioRequestId: 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
|
// Latch: ist endpointListeners fuer den aktuellen Session-Cycle schon gefeuert
|
||||||
// worden? Wird auf false gesetzt beim startStreamingRecording, auf true beim
|
// worden? Wird auf false gesetzt beim startStreamingRecording, auf true beim
|
||||||
// ersten Endpoint (egal ob via RVS oder Fallback). Verhindert Doppel-Fires.
|
// 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)}`;
|
const requestId = `sttstr_${Date.now()}_${Math.floor(Math.random() * 100000)}`;
|
||||||
this.streamRequestId = requestId;
|
this.streamRequestId = requestId;
|
||||||
this.streamAudioRequestId = opts.audioRequestId || '';
|
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.streamGotPartial = false;
|
||||||
this.streamEndpointFired = false;
|
this.streamEndpointFired = false;
|
||||||
this.recordingStartTime = Date.now();
|
this.recordingStartTime = Date.now();
|
||||||
@@ -1146,6 +1191,7 @@ class AudioService {
|
|||||||
requestId: sessionId,
|
requestId: sessionId,
|
||||||
pcm: String(e?.pcm || ''),
|
pcm: String(e?.pcm || ''),
|
||||||
seq: Number(e?.seq || 0),
|
seq: Number(e?.seq || 0),
|
||||||
|
targetInstance: this.streamTargetInstance,
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
this.streamPcmErrorSub = emitter.addListener('PcmStreamError', (e: any) => {
|
this.streamPcmErrorSub = emitter.addListener('PcmStreamError', (e: any) => {
|
||||||
@@ -1181,6 +1227,7 @@ class AudioService {
|
|||||||
hardCapMs: typeof opts.hardCapMs === 'number' ? opts.hardCapMs : 60000,
|
hardCapMs: typeof opts.hardCapMs === 'number' ? opts.hardCapMs : 60000,
|
||||||
sampleRate: 16000,
|
sampleRate: 16000,
|
||||||
projectId: opts.projectId || '',
|
projectId: opts.projectId || '',
|
||||||
|
targetInstance: this.streamTargetInstance,
|
||||||
});
|
});
|
||||||
|
|
||||||
// No-Speech-Watchdog — ersetzt den alten VAD-noSpeechTimer.
|
// No-Speech-Watchdog — ersetzt den alten VAD-noSpeechTimer.
|
||||||
@@ -1230,7 +1277,7 @@ class AudioService {
|
|||||||
if (!reqId) return;
|
if (!reqId) return;
|
||||||
const audioReqId = this.streamAudioRequestId;
|
const audioReqId = this.streamAudioRequestId;
|
||||||
try {
|
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) {
|
} catch (e) {
|
||||||
console.warn('[Audio] stt_stream_end senden fehlgeschlagen:', e);
|
console.warn('[Audio] stt_stream_end senden fehlgeschlagen:', e);
|
||||||
}
|
}
|
||||||
@@ -1265,7 +1312,7 @@ class AudioService {
|
|||||||
if (!reqId) return;
|
if (!reqId) return;
|
||||||
const audioReqId = this.streamAudioRequestId;
|
const audioReqId = this.streamAudioRequestId;
|
||||||
try {
|
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 {}
|
} catch {}
|
||||||
this._cleanupStreamLocal(`cancel:${reason}`);
|
this._cleanupStreamLocal(`cancel:${reason}`);
|
||||||
// Listener feuern damit ChatScreen reagieren kann (endConversation etc.)
|
// Listener feuern damit ChatScreen reagieren kann (endConversation etc.)
|
||||||
|
|||||||
+48
-12
@@ -1684,20 +1684,27 @@ class ARIABridge:
|
|||||||
if len(self._xtts_request_to_message) > 100:
|
if len(self._xtts_request_to_message) > 100:
|
||||||
oldest = next(iter(self._xtts_request_to_message))
|
oldest = next(iter(self._xtts_request_to_message))
|
||||||
self._xtts_request_to_message.pop(oldest, None)
|
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({
|
await self._send_to_rvs({
|
||||||
"type": "xtts_request",
|
"type": "xtts_request",
|
||||||
"payload": {
|
"payload": tts_payload,
|
||||||
"text": tts_text,
|
|
||||||
"voice": xtts_voice,
|
|
||||||
"speed": xtts_speed,
|
|
||||||
"language": "de",
|
|
||||||
"requestId": xtts_request_id,
|
|
||||||
"messageId": message_id,
|
|
||||||
},
|
|
||||||
"timestamp": int(asyncio.get_event_loop().time() * 1000),
|
"timestamp": int(asyncio.get_event_loop().time() * 1000),
|
||||||
})
|
})
|
||||||
logger.info("[core] XTTS-Request gesendet (voice=%s, speed=%.2fx): '%s'",
|
logger.info("[core] XTTS-Request gesendet (voice=%s, speed=%.2fx, target=%s): '%s'",
|
||||||
xtts_voice or "default", xtts_speed, tts_text[:60])
|
xtts_voice or "default", xtts_speed, tts_target or "(broadcast)", tts_text[:60])
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error("[core] XTTS-Request fehlgeschlagen: %s — kein Audio", 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 "?")
|
self._workers[iid]["gpus"] or "?", self._workers[iid]["model"] or "?")
|
||||||
return
|
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":
|
elif msg_type == "worker_ping":
|
||||||
iid = (payload.get("instanceId") or "").strip()
|
iid = (payload.get("instanceId") or "").strip()
|
||||||
if iid:
|
if iid:
|
||||||
@@ -3982,8 +4006,14 @@ class ARIABridge:
|
|||||||
req_payload["tools"] = tools
|
req_payload["tools"] = tools
|
||||||
if model:
|
if model:
|
||||||
req_payload["model"] = model
|
req_payload["model"] = model
|
||||||
logger.info("[rvs] llm_request → llm-adapter (id=%s, msgs=%d, max_tokens=%d, tools=%d, model=%s)",
|
# Redundanz/Multitasking: freie llm-Instanz gezielt adressieren; None
|
||||||
request_id[:8], len(messages), max_tokens, len(tools) if tools else 0, model or "-")
|
# → 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({
|
ok = await self._send_to_rvs({
|
||||||
"type": "llm_request",
|
"type": "llm_request",
|
||||||
"payload": req_payload,
|
"payload": req_payload,
|
||||||
@@ -4764,6 +4794,12 @@ class ARIABridge:
|
|||||||
chosen["busy"] = True # optimistisch, bis der naechste Ping korrigiert
|
chosen["busy"] = True # optimistisch, bis der naechste Ping korrigiert
|
||||||
return chosen["instanceId"]
|
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 = "",
|
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,
|
||||||
|
|||||||
@@ -897,6 +897,11 @@ async def run_loop(runner: F5Runner) -> None:
|
|||||||
continue
|
continue
|
||||||
mtype = msg.get("type", "")
|
mtype = msg.get("type", "")
|
||||||
payload = msg.get("payload", {}) or {}
|
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":
|
if mtype == "xtts_request":
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -246,6 +246,11 @@ async def _run() -> None:
|
|||||||
if msg.get("type") != "llm_request":
|
if msg.get("type") != "llm_request":
|
||||||
continue
|
continue
|
||||||
payload = msg.get("payload", {}) or {}
|
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,
|
# Jede Anfrage nebenlaeufig — llama.cpp serialisiert intern,
|
||||||
# aber wir blockieren so nicht den Empfang weiterer Messages.
|
# aber wir blockieren so nicht den Empfang weiterer Messages.
|
||||||
asyncio.create_task(_handle_llm_request(ws, payload))
|
asyncio.create_task(_handle_llm_request(ws, payload))
|
||||||
|
|||||||
@@ -751,6 +751,12 @@ async def run_loop(sessions: SessionManager) -> None:
|
|||||||
continue
|
continue
|
||||||
mtype = msg.get("type", "")
|
mtype = msg.get("type", "")
|
||||||
payload = msg.get("payload", {}) or {}
|
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":
|
if mtype == "stt_stream_start":
|
||||||
sessions.start_session(payload)
|
sessions.start_session(payload)
|
||||||
elif mtype == "stt_audio_chunk":
|
elif mtype == "stt_audio_chunk":
|
||||||
|
|||||||
@@ -899,6 +899,11 @@ async def run_loop(runner: WhisperRunner, sessions: SessionManager) -> None:
|
|||||||
continue
|
continue
|
||||||
mtype = msg.get("type", "")
|
mtype = msg.get("type", "")
|
||||||
payload = msg.get("payload", {}) or {}
|
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":
|
if mtype == "stt_request":
|
||||||
req_id = payload.get("requestId", "?")
|
req_id = payload.get("requestId", "?")
|
||||||
|
|||||||
Reference in New Issue
Block a user