diff --git a/app/web/app.js b/app/web/app.js index 9590c64..41b37b7 100644 --- a/app/web/app.js +++ b/app/web/app.js @@ -16,6 +16,9 @@ function overrides() { } let busy = false; +let activeWs = null; // aktive WS-Verbindung (fuer Barge-in von aussen erreichbar) +let activeSources = []; // laufende AudioBufferSourceNodes (stoppbar per Barge-in) + // Session-ID pro Nutzer (sonst "gehoert einem anderen Nutzer"-Konflikt). Wird aus // /api/me abgeleitet; bis dahin null -> ensureSession() wartet auf loadMe(). let sessionId = null; @@ -99,6 +102,32 @@ function addMessage(role, text) { return div; } +// ---------- Barge-in-Hilfsfunktionen ---------- +function stopAudio() { + activeSources.forEach((s) => { try { s.stop(); } catch (e) {} }); + activeSources = []; +} + +function setMicIdle() { + micBtn.textContent = "🎤"; + micBtn.title = "Mikrofon"; + micBtn.classList.remove("bg-red-600", "animate-pulse", "bg-amber-500"); + micBtn.classList.add("bg-emerald-600"); +} + +function setMicRecording() { + micBtn.classList.remove("bg-emerald-600", "bg-amber-500"); + micBtn.classList.add("bg-red-600", "animate-pulse"); + micBtn.title = "Aufnahme stoppen"; +} + +function setMicBusy() { + micBtn.textContent = "⏹"; + micBtn.title = "KI unterbrechen (Barge-in)"; + micBtn.classList.remove("bg-emerald-600", "bg-red-600", "animate-pulse"); + micBtn.classList.add("bg-amber-500"); +} + // ---------- Audio-Wiedergabe (PCM s16le) ---------- let audioCtx = null; function playPcm(chunks, sampleRate) { @@ -116,6 +145,8 @@ function playPcm(chunks, sampleRate) { const src = audioCtx.createBufferSource(); src.buffer = buf; src.connect(audioCtx.destination); + activeSources.push(src); + src.onended = () => { activeSources = activeSources.filter((s) => s !== src); }; src.start(); } @@ -129,6 +160,7 @@ function wsUrl(path) { function runTurn(path, onopen) { return new Promise((resolve) => { const ws = new WebSocket(wsUrl(path)); + activeWs = ws; ws.binaryType = "arraybuffer"; const pcm = []; let answerEl = null; @@ -136,7 +168,7 @@ function runTurn(path, onopen) { ws.onopen = () => onopen(ws); ws.onerror = () => { statusEl.textContent = "Verbindungsfehler"; resolve(); }; - ws.onclose = () => resolve(); + ws.onclose = () => { activeWs = null; resolve(); }; ws.onmessage = (event) => { if (typeof event.data !== "string") { pcm.push(event.data); return; } @@ -179,6 +211,7 @@ function runTurn(path, onopen) { async function sendText(text) { if (busy || !text.trim()) return; busy = true; + setMicBusy(); await ensureSession(); addMessage("user", text); statusEl.textContent = "denkt …"; @@ -186,6 +219,7 @@ async function sendText(text) { ws.send(JSON.stringify({ text, stream: true, audio_stream: true, ...overrides() })); }); busy = false; + setMicIdle(); } formEl.addEventListener("submit", (e) => { @@ -236,20 +270,19 @@ async function startRecording() { sendVoice(bytes, mime.includes("mp4") || mime.includes("mpeg") ? "mp4" : "webm"); }; mediaRecorder.start(); - micBtn.classList.remove("bg-emerald-600"); - micBtn.classList.add("bg-red-600", "animate-pulse"); + setMicRecording(); statusEl.textContent = "Aufnahme … (zum Stoppen erneut tippen)"; } function stopRecording() { if (mediaRecorder && mediaRecorder.state !== "inactive") mediaRecorder.stop(); - micBtn.classList.remove("bg-red-600", "animate-pulse"); - micBtn.classList.add("bg-emerald-600"); + setMicIdle(); } async function sendVoice(bytes, fmt = "webm") { if (busy) return; busy = true; + setMicBusy(); await ensureSession(); statusEl.textContent = "verarbeite Sprache …"; await runTurn("/ws/voice", (ws) => { @@ -258,11 +291,21 @@ async function sendVoice(bytes, fmt = "webm") { ws.send(JSON.stringify({ type: "end" })); }); busy = false; + setMicIdle(); } micBtn.addEventListener("click", () => { - if (mediaRecorder && mediaRecorder.state === "recording") stopRecording(); - else startRecording(); + if (busy) { + // Barge-in: laufende KI-Antwort sofort unterbrechen + if (activeWs) activeWs.send(JSON.stringify({ type: "interrupt" })); + stopAudio(); + statusEl.textContent = "Unterbrochen"; + // busy + Button-Reset erfolgen sobald die WS schliesst (-> runTurn resolve -> sendVoice/sendText) + } else if (mediaRecorder && mediaRecorder.state === "recording") { + stopRecording(); + } else { + startRecording(); + } }); $("#reload-users").addEventListener("click", loadUsers); diff --git a/scripts/voice_loop.py b/scripts/voice_loop.py index 23cd9b8..3e30359 100644 --- a/scripts/voice_loop.py +++ b/scripts/voice_loop.py @@ -32,6 +32,7 @@ import signal import subprocess import sys import tempfile +import select import time import wave @@ -256,6 +257,22 @@ async def _player_feeder(player, queue: "asyncio.Queue") -> None: await asyncio.to_thread(_feed_player, player, chunk) +async def _stdin_watcher() -> None: + """Wartet per Polling auf Enter-Tastendruck fuer Barge-in. + + Kein blockierender Thread: select.select mit Timeout 0 ist nicht-blockierend + und funktioniert zuverlaessig auf Linux/macOS mit echtem tty-stdin. + Canceln ist sicher (kein Zeichen wird konsumiert, ausser Enter wird tatsaechlich + gedrueckt und diese Funktion war diejenige, die ihn gelesen hat). + """ + while True: + ready, _, _ = select.select([sys.stdin], [], [], 0) + if ready: + sys.stdin.readline() # Enter-Zeilenumbruch konsumieren + return + await asyncio.sleep(0.05) + + def _start_frame(args) -> dict: frame = {"type": "start", "format": "wav"} # audio_stream ist Tri-State: True/False explizit, None -> Server-Default entscheidet. @@ -325,6 +342,10 @@ async def one_turn(ws, wav: bytes, start_frame: dict, queue: "asyncio.Queue | No async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None: """Eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via session_id). + Barge-in: Waehrend der Server antwortet, wartet _stdin_watcher() auf Enter. + Wird Enter gedrueckt, sendet der Client {"type":"interrupt"}, der Server bricht + den Turn ab, der Player wird sofort beendet. + Es wird immer ein durchgehender Player benutzt: der spielt sowohl viele Haeppchen (Stream-Audio) als auch ein einzelnes Komplett-Audio luekenlos ab. Wird kein Player gefunden, wird das Audio gesammelt und als WAV abgespielt. Die Wiedergabe blockiert @@ -333,17 +354,52 @@ async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None: queue: asyncio.Queue | None = asyncio.Queue() if player is not None else None feeder = asyncio.create_task(_player_feeder(player, queue)) if queue is not None else None audio = b"" + interrupted = False try: # ping_timeout/max_queue defensiv: lange Antworten + Echtzeit-Wiedergabe ueberleben. async with websockets.connect(url, max_size=None, max_queue=None, ping_timeout=60) as ws: - audio = await one_turn(ws, wav, start_frame, queue) + turn_task = asyncio.create_task(one_turn(ws, wav, start_frame, queue)) + watcher_task = asyncio.create_task(_stdin_watcher()) + done, _ = await asyncio.wait( + {turn_task, watcher_task}, + return_when=asyncio.FIRST_COMPLETED, + ) + if turn_task in done: + # Normales Ende — Watcher canceln (kein Zeichen konsumiert) + watcher_task.cancel() + audio = turn_task.result() + else: + # Barge-in: Enter wurde gedrueckt + interrupted = True + try: + await ws.send(json.dumps({"type": "interrupt"})) + except Exception: + pass + turn_task.cancel() + try: + await turn_task + except asyncio.CancelledError: + pass finally: if feeder is not None: - queue.put_nowait(None) # Feeder beenden, restliche Haeppchen noch abspielen - await feeder + queue.put_nowait(None) # Feeder beenden + if interrupted and player is not None: + player.terminate() # sofort stoppen statt Restpuffer abwarten + try: + await feeder + except Exception: + pass if player is not None: - await asyncio.to_thread(_close_player, player) - if audio: # nur wenn kein Streaming-Player verfuegbar war + if interrupted: + try: + player.wait(timeout=2) + except subprocess.TimeoutExpired: + player.kill() + else: + await asyncio.to_thread(_close_player, player) + if interrupted: + print(" ↩ Unterbrochen — drücke [Enter] für neue Aufnahme …") + elif audio: # nur wenn kein Streaming-Player verfuegbar war play_pcm(audio)