Files
tes 8bb0ce51b1 Plan 3 fix: GA turn_detection auto-create_response + reset state on WS reconnect
GA Realtime API requires create_response:true and interrupt_response:true
inside turn_detection — without them the model never auto-generates a
reply after VAD detects speech_stopped, so the relay idles out at 30s
with no assistant audio. Add both.

Client side: when the /device WS drops (Coolify rolling deploy, etc.),
Session.run_forever reconnects but the StateMachine kept whatever state
it was in — usually LISTENING, which permanently disables the wakeword.
Add StateMachine.force_idle() and a Callbacks.on_disconnected() hook;
main.py wires it to force the state back to IDLE on every disconnect.
Also harden Session's fire-and-forget send_* helpers to swallow
ConnectionClosed instead of bubbling up as unhandled task exceptions.
2026-06-12 06:50:00 +00:00

193 lines
6.2 KiB
Python

"""Smart Assistant Pi client entry point.
CRITICAL: set PA_ALSA_PLUGHW BEFORE importing sounddevice (see findings.md §1).
"""
import os
os.environ.setdefault("PA_ALSA_PLUGHW", "1")
import asyncio
import queue
import signal
import sys
import threading
from pathlib import Path
import numpy as np
import sounddevice as sd # safe to import now
from client.audio import (
AudioStreams,
FRAME_SAMPLES,
resample_24k_to_16k,
)
from client.config import ConfigMissingError, load_config
from client.log import get_logger
from client.playback import PlaybackWorker
from client.session import Session
from client.state import StateMachine
from client.wakeword import WakewordDetector
_log = get_logger("main")
ASSISTANT_STATE_DIR = Path(os.environ["HOME"]) / "assistant" / "state"
WAKE_WAV = ASSISTANT_STATE_DIR / "wake.wav"
SLEEP_WAV = ASSISTANT_STATE_DIR / "sleep.wav"
def _build_bridge_thread(
mic_queue: queue.Queue,
sm: StateMachine,
wakeword: WakewordDetector,
session: Session,
playback: PlaybackWorker,
stop_event: threading.Event,
) -> threading.Thread:
def run():
frame_counter = 0
recent_rms_max = 0.0
while not stop_event.is_set():
try:
pcm16_bytes = mic_queue.get(timeout=0.25)
except queue.Empty:
continue
frame_counter += 1
raw = np.frombuffer(pcm16_bytes, dtype=np.int16)
rms = float(np.sqrt(np.mean(raw.astype(np.float32) ** 2)))
if rms > recent_rms_max:
recent_rms_max = rms
# Every ~5 s emit a heartbeat so we can see if the InputStream is
# actually producing frames and what state we're routing to.
if frame_counter % 62 == 0:
_log.info(
"bridge frames=%d state=%s wake_en=%s up_en=%s qsize=%d rms_max_5s=%.0f",
frame_counter, sm.state.name,
sm.wakeword_enabled, sm.uplink_enabled, mic_queue.qsize(),
recent_rms_max,
)
recent_rms_max = 0.0
if sm.wakeword_enabled:
frame = np.frombuffer(pcm16_bytes, dtype=np.int16)
if frame.size != FRAME_SAMPLES:
continue
feed = resample_24k_to_16k(frame)
if wakeword.predict(feed):
if sm.wake_fired():
if WAKE_WAV.exists():
playback.enqueue_wav(WAKE_WAV, "wake")
else:
# Even with no wav, fire the done marker so the
# playback_done(for=wake) envelope is sent.
playback.enqueue_done_marker("wake")
session.call_soon_threadsafe_send_wake()
elif sm.uplink_enabled:
session.call_soon_threadsafe_send_binary(pcm16_bytes)
# else: ASSISTANT_SPEAKING / WAKE_PENDING / IDLE_PENDING -> drop
t = threading.Thread(target=run, name="bridge", daemon=True)
return t
def _make_callbacks(sm: StateMachine, playback: PlaybackWorker):
class _CB:
def on_session_started(self, conversation_id: str) -> None:
_log.info("session_started conversation_id=%s", conversation_id)
sm.session_started()
def on_assistant_audio(self, pcm16: bytes) -> None:
sm.assistant_audio_arrived()
playback.enqueue_audio(pcm16)
def on_assistant_done(self) -> None:
_log.info("assistant_done")
sm.assistant_done()
def on_session_ended(self, reason: str) -> None:
_log.info("session_ended reason=%s", reason)
sm.session_ended()
if SLEEP_WAV.exists():
playback.enqueue_wav(SLEEP_WAV, "sleep")
else:
playback.enqueue_done_marker("sleep")
def on_disconnected(self) -> None:
# The conversation is gone with the WS — reset to IDLE so the
# wakeword detector can re-arm. Any pending assistant audio left
# in the playback queue still drains, but we drop a "sleep" done
# marker only if we were mid-session.
_log.info("on_disconnected: forcing state to IDLE")
sm.force_idle()
return _CB()
def main() -> None:
try:
cfg = load_config()
except ConfigMissingError as exc:
raise SystemExit(str(exc)) from exc
_log.info("device_id=%s backend=%s", cfg.device_id, cfg.backend_ws)
mic_queue: queue.Queue = queue.Queue(maxsize=64)
streams = AudioStreams.open(sd, mic_queue=mic_queue)
streams.start()
sm = StateMachine()
session_ref: list[Session] = []
def on_playback_done(label: str) -> None:
_log.info("playback_done for=%s", label)
if not session_ref:
return
session = session_ref[0]
session.call_soon_threadsafe_send_playback_done(label)
if label == "sleep":
sm.sleep_played()
playback = PlaybackWorker(streams, on_done=on_playback_done)
playback.start()
wakeword = WakewordDetector()
callbacks = _make_callbacks(sm, playback)
session = Session(
backend_ws=cfg.backend_ws,
device_token=cfg.device_token,
device_id=cfg.device_id,
callbacks=callbacks,
)
session_ref.append(session)
def set_volume_handler(args: dict) -> dict:
level = int(args.get("level", -1))
playback.set_volume(level)
return {"ok": True, "level": level}
session.register_tool("set_volume", set_volume_handler)
stop_event = threading.Event()
bridge = _build_bridge_thread(mic_queue, sm, wakeword, session, playback, stop_event)
bridge.start()
def _shutdown(signum, _frame):
_log.info("signal %d received; shutting down", signum)
stop_event.set()
try:
playback.stop()
finally:
streams.close()
sys.exit(0)
signal.signal(signal.SIGTERM, _shutdown)
signal.signal(signal.SIGINT, _shutdown)
try:
asyncio.run(session.run_forever())
finally:
stop_event.set()
playback.stop()
streams.close()
if __name__ == "__main__":
main()