Skip to content

Commit 034281b

Browse files
committed
feat(voice): cut over wearable clients to protocol v1
1 parent b0ef0f1 commit 034281b

38 files changed

Lines changed: 2872 additions & 707 deletions

File tree

‎backends/advanced/src/advanced_omi_backend/controllers/websocket_controller.py‎

Lines changed: 3 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -66,11 +66,6 @@
6666
SessionStatus,
6767
SessionStore,
6868
)
69-
from advanced_omi_backend.services.device_audio import (
70-
is_opus_streaming_client,
71-
stop_play_audio,
72-
stream_play_audio_as_opus,
73-
)
7469
from advanced_omi_backend.services.observability import record_event_sync
7570
from advanced_omi_backend.services.plugin_service import get_plugin_router
7671
from advanced_omi_backend.services.response_coordinator import (
@@ -519,10 +514,8 @@ async def subscribe_to_device_downlink(
519514
) -> None:
520515
"""Forward backend→device control messages from Redis Pub/Sub to the device WebSocket.
521516
522-
Any backend component (wake-word service, plugins) can push a message to a
523-
specific device by publishing to ``device:downlink:{client_id}``. Each message
524-
is a Wyoming-style control frame (e.g. ``{"type": "play-audio", "data": {...}}``)
525-
which the HAVPE relay's ``handle_backend_messages`` dispatches to the device.
517+
Any backend component can push a bound voice event or non-audio device control
518+
message to ``device:downlink:{client_id}``.
526519
527520
Runs as a background task for the lifetime of the WebSocket connection.
528521
"""
@@ -532,7 +525,6 @@ async def subscribe_to_device_downlink(
532525
channel = str(device_downlink_channel(client_id))
533526
redis_client = None
534527
pubsub = None
535-
opus_stream = is_opus_streaming_client(client_id_value)
536528
ack_deadline_tasks: set[asyncio.Task] = set()
537529

538530
try:
@@ -541,10 +533,7 @@ async def subscribe_to_device_downlink(
541533
responses = ResponseCoordinator(redis_client, voice_sessions)
542534
pubsub = redis_client.pubsub()
543535
await pubsub.subscribe(channel)
544-
logger.info(
545-
f"🔊 Subscribed to device downlink channel: {channel}"
546-
f"{' (opus streaming)' if opus_stream else ''}"
547-
)
536+
logger.info(f"🔊 Subscribed to device downlink channel: {channel}")
548537

549538
while True:
550539
try:
@@ -607,23 +596,6 @@ async def expire_after(
607596
)
608597
ack_deadline_tasks.add(task)
609598
task.add_done_callback(ack_deadline_tasks.discard)
610-
elif msg_type == "stop-audio":
611-
# Barge-in: stop whatever TTS is playing on the device. For
612-
# Opus clients this cancels the in-flight stream + flushes the
613-
# device; for others there's nothing streaming to cancel, so
614-
# just forward the control frame best-effort.
615-
if opus_stream:
616-
await stop_play_audio(websocket, client_id_value)
617-
else:
618-
await websocket.send_json(payload)
619-
elif opus_stream and msg_type == "play-audio":
620-
# RAM-limited devices can't take a big base64 WAV frame;
621-
# transcode + stream it as small Opus packets instead.
622-
streamed = await stream_play_audio_as_opus(
623-
websocket, payload.get("data") or {}, client_id_value
624-
)
625-
if not streamed:
626-
await websocket.send_json(payload) # fallback
627599
else:
628600
await websocket.send_json(payload)
629601
data = payload.get("data")

‎backends/advanced/src/advanced_omi_backend/plugins/services.py‎

Lines changed: 13 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,14 @@
55
(e.g., close a conversation) or with other plugins (e.g., call Home Assistant to toggle lights).
66
"""
77

8-
import json
98
import logging
109
from typing import TYPE_CHECKING, Optional
1110

1211
from advanced_omi_backend.models.conversation import Conversation
1312
from advanced_omi_backend.redis_factory import create_async_redis
14-
from advanced_omi_backend.redis_keys import ClientId, device_downlink_channel
1513
from advanced_omi_backend.services.audio_stream.session_store import SessionStore
14+
from advanced_omi_backend.services.response_coordinator import ResponseCoordinator
15+
from advanced_omi_backend.services.voice_sessions import VoiceSessionCoordinator
1616
from advanced_omi_backend.users import User
1717

1818
from .base import PluginContext, PluginResult
@@ -115,31 +115,18 @@ async def star_conversation(self, session_id: str) -> bool:
115115
# toggle_star returns a dict on success, JSONResponse on error
116116
return isinstance(result, dict) and "starred" in result
117117

118-
async def stop_playback(self, client_id: str) -> bool:
119-
"""Stop any TTS currently playing on a device (barge-in).
120-
121-
Publishes a ``stop-audio`` control frame to the device's downlink channel.
122-
The WebSocket handler that owns the device connection picks it up and, for
123-
Opus-streaming clients, cancels the in-flight stream and tells the device to
124-
flush (see ``device_audio.stop_play_audio``). Decoupled via Redis so this
125-
works from any process (the button handler runs in the backend, but wake
126-
handlers run in the workers).
127-
128-
Args:
129-
client_id: The device/client whose playback should stop.
130-
131-
Returns:
132-
True if the stop request was published.
133-
"""
134-
if not client_id:
135-
logger.warning("stop_playback called with no client_id")
136-
return False
137-
message = json.dumps({"type": "stop-audio", "data": {}})
138-
client_ref = ClientId.from_value(client_id)
139-
await self._async_redis.publish(
140-
str(device_downlink_channel(client_ref)), message
118+
async def stop_playback(self, user_id: str, client_id: str) -> bool:
119+
"""Cancel the client's current protocol-v1 response at its generation fence."""
120+
if not user_id or not client_id:
121+
raise ValueError("stop_playback requires user_id and client_id")
122+
coordinator = ResponseCoordinator(
123+
self._async_redis,
124+
VoiceSessionCoordinator(self._async_redis),
125+
)
126+
await coordinator.begin_turn(user_id, client_id, reason="barge_in")
127+
logger.info(
128+
"Button requested coordinated playback cancellation for %s", client_id
141129
)
142-
logger.info(f"⏹ Requested stop-audio for {client_id}")
143130
return True
144131

145132
async def call_plugin(

‎backends/advanced/src/advanced_omi_backend/services/device_audio.py‎

Lines changed: 0 additions & 177 deletions
This file was deleted.

0 commit comments

Comments
 (0)