feat(barge-in): Unterbrechen der KI-Antwort per Mikrofon-Button und Enter

Web-Interface (app.js):
- Mic-Button hat jetzt 3 Zustände: idle (grün 🎤) / Aufnahme (rot pulsierend) /
  KI antwortet (amber ⏹, klicken = Barge-in)
- Barge-in sendet {"type":"interrupt"} auf der aktiven WS-Verbindung (activeWs)
  und stoppt alle laufenden AudioBufferSourceNodes sofort (stopAudio)
- busy + Button-Reset erfolgen automatisch wenn die WS schliesst

Terminal (voice_loop.py):
- _stdin_watcher(): polling via select.select (nicht-blockierend, kein Thread-Leak)
  wartet auf Enter-Tastendruck waehrend Server antwortet
- _send_and_play() raced one_turn gegen _stdin_watcher mit asyncio.wait;
  bei Barge-in: interrupt-Frame an Server, Player.terminate(), feeder beenden
- Feedback: "↩ Unterbrochen — drücke [Enter] für neue Aufnahme …"

Server-Seite war bereits vollstaendig implementiert (ws.py, kein Handlungsbedarf).

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
Dieter Schlüter 2026-06-18 18:38:33 +02:00
commit 0d907a1886
2 changed files with 111 additions and 12 deletions

View file

@ -16,6 +16,9 @@ function overrides() {
} }
let busy = false; 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 // Session-ID pro Nutzer (sonst "gehoert einem anderen Nutzer"-Konflikt). Wird aus
// /api/me abgeleitet; bis dahin null -> ensureSession() wartet auf loadMe(). // /api/me abgeleitet; bis dahin null -> ensureSession() wartet auf loadMe().
let sessionId = null; let sessionId = null;
@ -99,6 +102,32 @@ function addMessage(role, text) {
return div; 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) ---------- // ---------- Audio-Wiedergabe (PCM s16le) ----------
let audioCtx = null; let audioCtx = null;
function playPcm(chunks, sampleRate) { function playPcm(chunks, sampleRate) {
@ -116,6 +145,8 @@ function playPcm(chunks, sampleRate) {
const src = audioCtx.createBufferSource(); const src = audioCtx.createBufferSource();
src.buffer = buf; src.buffer = buf;
src.connect(audioCtx.destination); src.connect(audioCtx.destination);
activeSources.push(src);
src.onended = () => { activeSources = activeSources.filter((s) => s !== src); };
src.start(); src.start();
} }
@ -129,6 +160,7 @@ function wsUrl(path) {
function runTurn(path, onopen) { function runTurn(path, onopen) {
return new Promise((resolve) => { return new Promise((resolve) => {
const ws = new WebSocket(wsUrl(path)); const ws = new WebSocket(wsUrl(path));
activeWs = ws;
ws.binaryType = "arraybuffer"; ws.binaryType = "arraybuffer";
const pcm = []; const pcm = [];
let answerEl = null; let answerEl = null;
@ -136,7 +168,7 @@ function runTurn(path, onopen) {
ws.onopen = () => onopen(ws); ws.onopen = () => onopen(ws);
ws.onerror = () => { statusEl.textContent = "Verbindungsfehler"; resolve(); }; ws.onerror = () => { statusEl.textContent = "Verbindungsfehler"; resolve(); };
ws.onclose = () => resolve(); ws.onclose = () => { activeWs = null; resolve(); };
ws.onmessage = (event) => { ws.onmessage = (event) => {
if (typeof event.data !== "string") { pcm.push(event.data); return; } if (typeof event.data !== "string") { pcm.push(event.data); return; }
@ -179,6 +211,7 @@ function runTurn(path, onopen) {
async function sendText(text) { async function sendText(text) {
if (busy || !text.trim()) return; if (busy || !text.trim()) return;
busy = true; busy = true;
setMicBusy();
await ensureSession(); await ensureSession();
addMessage("user", text); addMessage("user", text);
statusEl.textContent = "denkt …"; statusEl.textContent = "denkt …";
@ -186,6 +219,7 @@ async function sendText(text) {
ws.send(JSON.stringify({ text, stream: true, audio_stream: true, ...overrides() })); ws.send(JSON.stringify({ text, stream: true, audio_stream: true, ...overrides() }));
}); });
busy = false; busy = false;
setMicIdle();
} }
formEl.addEventListener("submit", (e) => { formEl.addEventListener("submit", (e) => {
@ -236,20 +270,19 @@ async function startRecording() {
sendVoice(bytes, mime.includes("mp4") || mime.includes("mpeg") ? "mp4" : "webm"); sendVoice(bytes, mime.includes("mp4") || mime.includes("mpeg") ? "mp4" : "webm");
}; };
mediaRecorder.start(); mediaRecorder.start();
micBtn.classList.remove("bg-emerald-600"); setMicRecording();
micBtn.classList.add("bg-red-600", "animate-pulse");
statusEl.textContent = "Aufnahme … (zum Stoppen erneut tippen)"; statusEl.textContent = "Aufnahme … (zum Stoppen erneut tippen)";
} }
function stopRecording() { function stopRecording() {
if (mediaRecorder && mediaRecorder.state !== "inactive") mediaRecorder.stop(); if (mediaRecorder && mediaRecorder.state !== "inactive") mediaRecorder.stop();
micBtn.classList.remove("bg-red-600", "animate-pulse"); setMicIdle();
micBtn.classList.add("bg-emerald-600");
} }
async function sendVoice(bytes, fmt = "webm") { async function sendVoice(bytes, fmt = "webm") {
if (busy) return; if (busy) return;
busy = true; busy = true;
setMicBusy();
await ensureSession(); await ensureSession();
statusEl.textContent = "verarbeite Sprache …"; statusEl.textContent = "verarbeite Sprache …";
await runTurn("/ws/voice", (ws) => { await runTurn("/ws/voice", (ws) => {
@ -258,11 +291,21 @@ async function sendVoice(bytes, fmt = "webm") {
ws.send(JSON.stringify({ type: "end" })); ws.send(JSON.stringify({ type: "end" }));
}); });
busy = false; busy = false;
setMicIdle();
} }
micBtn.addEventListener("click", () => { micBtn.addEventListener("click", () => {
if (mediaRecorder && mediaRecorder.state === "recording") stopRecording(); if (busy) {
else startRecording(); // 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); $("#reload-users").addEventListener("click", loadUsers);

View file

@ -32,6 +32,7 @@ import signal
import subprocess import subprocess
import sys import sys
import tempfile import tempfile
import select
import time import time
import wave import wave
@ -256,6 +257,22 @@ async def _player_feeder(player, queue: "asyncio.Queue") -> None:
await asyncio.to_thread(_feed_player, player, chunk) 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: def _start_frame(args) -> dict:
frame = {"type": "start", "format": "wav"} frame = {"type": "start", "format": "wav"}
# audio_stream ist Tri-State: True/False explizit, None -> Server-Default entscheidet. # 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: async def _send_and_play(url: str, wav: bytes, start_frame: dict) -> None:
"""Eine kurze Verbindung pro Runde (Gedaechtnis bleibt serverseitig via session_id). """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 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 (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 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 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 feeder = asyncio.create_task(_player_feeder(player, queue)) if queue is not None else None
audio = b"" audio = b""
interrupted = False
try: try:
# ping_timeout/max_queue defensiv: lange Antworten + Echtzeit-Wiedergabe ueberleben. # 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: 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: finally:
if feeder is not None: if feeder is not None:
queue.put_nowait(None) # Feeder beenden, restliche Haeppchen noch abspielen queue.put_nowait(None) # Feeder beenden
await feeder 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: if player is not None:
await asyncio.to_thread(_close_player, player) if interrupted:
if audio: # nur wenn kein Streaming-Player verfuegbar war 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) play_pcm(audio)