diff --git a/.gitignore b/.gitignore
new file mode 100644
index 0000000..541dda0
--- /dev/null
+++ b/.gitignore
@@ -0,0 +1,14 @@
+# macOS
+.DS_Store
+
+# JetBrains IDEs
+.idea/
+clients/.idea/
+clients/python-common-sdk/.idea/
+
+# Python bytecode cache
+**/__pycache__/
+*.py[cod]
+*$py.class
+
+.env
\ No newline at end of file
diff --git a/clients/python-common-sdk/README-zh.md b/clients/python-common-sdk/README-zh.md
index 2ec7054..09cab72 100644
--- a/clients/python-common-sdk/README-zh.md
+++ b/clients/python-common-sdk/README-zh.md
@@ -1,111 +1,138 @@
-# EVA OS Python 客户端
+# EVA OS Python 客户端 SDK 接入指南
-本项目是连接 EVA OS 实时多模态 AI 服务的官方 Python 参考实现。它展示了如何通过 API 进行身份验证,使用 LiveKit 建立 WebRTC 连接,并处理实时的双向音视频流。
+本项目是连接硬件设备与 **EVA OS 实时多模态 AI 服务** 的官方 Python 参考实现。它展示了如何进行接口鉴权、建立超低延迟的长连接,并处理实时的双向音视频流。
-## 前置要求
+更重要的是,这份说明书将引导你如何将定制硬件(无论是 PC 还是低功耗 IoT 设备)完美契合进 EVA OS V2 的**“云边协同 (Edge-Cloud Synergy)”**架构之中。
-### 系统依赖
-本项目依赖 `PortAudio` 库。在安装 Python `PyAudio` 库之前,**必须**先安装系统级的开发头文件。
+---
-* **Ubuntu / Debian (及 Rockchip/树莓派等嵌入式系统):**
- ```bash
- sudo apt-get update
- sudo apt-get install libportaudio2 portaudio19-dev
- ```
+## 🏗 架构理念:云边协同与边端状态主权
-* **macOS:**
- ```bash
- brew install portaudio
- ```
+在接入 SDK 之前,强烈建议您了解 EVA OS V2 的设计原则,我们称之为**“边端状态主权”**。
-* **Windows:**
- 通常 `pip` 会自动安装预编译的二进制包(Wheels)。
+### 职责划分
-### Python 环境
-* 需要 Python 3.11 或更高版本。
+**边端(您的硬件与本 SDK)**
+硬件算力珍贵,因此边端必须保持专注和轻量。SDK 仅负责:
+1. **流媒体传输:** 持续不断地推送麦克风阵列数据,并拉取扬声器音频。
+2. **本地唤醒检测:** 在后台运行极低内存占用的 VOSK 模型,实现零延迟的唤醒词检测(如“你好方舟”)。
+3. **状态请求:** 边端**绝对不做** VAD(静音断句检测)和意图分类。当捕获到唤醒词或按键事件时,SDK 只是通过控制平面向云端发送任务切换的意向。
+4. **指令执行:** 边端实时接收来自云端的设备控制与情绪感知消息,实现如调整音量、UI 情绪反馈等即时响应,该链路完全绕过云端核心业务流。
-## 安装步骤
+**云端(EVA OS / Pipecat 引擎)**
+云端作为中央大脑,拥有充足的算力,它持续接收音视频流并统筹复杂的业务逻辑:
+1. **全局 VAD 与 ASR:** 判断用户何时说话结束,并将音频转为文本。
+2. **意图理解与业务流转:** 执行复杂的工作流并响应用户请求。
+3. **任务调度协同:** 云端负责意图理解并下发建议,或对边端的切换请求进行审批。但在我们的架构中,**边端是状态的最终主权拥有者**。只有当边端真正确认并执行了状态切换,云端才会更新其内部的镜像状态,并顺滑地切出对应的 AI Agent 来接管对话。
-1. **克隆代码仓库:**
- ```bash
- git clone git@github.com:AutoArk/EVA-OS.git
- cd EVA-OS/clients/python-common-sdk
- ```
+这种握手机制确保了您的物理设备与云端数字大脑的状态(通过 RTVI 协议)永远保持绝对同步。
+
+---
+
+## 🚀 快速入门
+
+### 1. 安装系统依赖
+本项目底层依赖 `PortAudio`。在安装 Python 依赖前,**必须**先安装系统级的开发库。
-2. **创建并激活虚拟环境 (推荐):**
+* **Ubuntu / Debian (及树莓派/开发板等):**
```bash
- python3 -m venv venv
- source venv/bin/activate
+ sudo apt-get update
+ sudo apt-get install libportaudio2 portaudio19-dev
```
-
-3. **安装 Python 依赖:**
+* **macOS:**
```bash
- pip install livekit requests pyaudio numpy python-dotenv opencv-python
+ brew install portaudio
```
-## 配置说明
-
-在项目根目录下创建一个 `.env` 文件。
+### 2. 配置环境
+克隆代码并安装依赖:
+```bash
+python3 -m venv venv
+source venv/bin/activate
+pip install -r requirements.txt
+# 可选:如果需要摄像头画面,需安装 opencv
+pip install opencv-python
+```
+在项目根目录创建 `.env` 文件:
```ini
-# .env 文件
-# 必填:API Key
-EVA_API_KEY=YOUR_ACTUAL_API_KEY_HERE
+# 从 EVA OS 后台创建应用获取的 Solution API Key
+EVA_API_KEY=sk-your-api-key-here
```
-### EvaClient 参数详细说明
+### 3. 获取音频设备索引(极其重要!)
+`PyAudio` 对音频通道极其敏感(例如尝试用音箱作为输入录音会导致 `Invalid number of channels` 报错)。
+请务必运行设备扫描脚本:
+```bash
+python list_audio_devices.py
+```
+记下你真实麦克风和扬声器对应的 `Index` 数字,填入示例代码的 `mic_index` 和 `spk_index` 中。
-| 参数名 | 类型 | 必填 | 默认值 | 说明 |
-| :--- | :--- | :---: | :--- | :--- |
-| **api_key** | `str` | **是** | 无 | **认证密钥**。
用于访问 Eva API 服务的凭证。 |
-| **mic_index** | `int` | 否 | `0` | **麦克风设备索引**。
指定用于录音的输入设备 ID。`0` 通常代表系统默认麦克风。 |
-| **spk_index** | `int` | 否 | `0` | **扬声器设备索引**。
指定用于播放音频的输出设备 ID。`0` 通常代表系统默认扬声器。 |
-| **mic_sample_rate** | `int` | 否 | `48000` | **麦克风采样率 (Hz)**。
录音时的音频采样频率,默认为 48kHz 。 |
-| **spk_sample_rate** | `int` | 否 | `48000` | **扬声器采样率 (Hz)**。
播放时的音频采样频率,默认为 48kHz。 |
-| **mic_channels** | `int` | 否 | `1` | **麦克风通道数**。
`1` 表示单声道 (Mono),`2` 表示立体声 (Stereo)。 |
-| **spk_channels** | `int` | 否 | `1` | **扬声器通道数**。
`1` 表示单声道 (Mono),`2` 表示立体声 (Stereo)。 |
-| **frame_size_ms** | `int` | 否 | `60` | **音频帧时长 (毫秒)**。
每次处理或传输的音频数据块的时间长度。这会影响延迟和网络包的大小。 |
-| **camera_index** | `int` | 否 | `0` | **摄像头设备索引**。
通常 `0` 对应默认系统摄像头。若有多个摄像头,请依序尝试。 |
-| **video_width** | `int` | 否 | `640` | **视频宽度**。
摄像头捕获图像的水平像素数。ARM 设备建议不要设置过高。 |
-| **video_height** | `int` | 否 | `480` | **视频高度**。
摄像头捕获图像的垂直像素数。 |
-| **video_fps** | `int` | 否 | `30` | **视频帧率**。
每秒传输的帧数 (FPS)。在低性能设备上建议设为 15 或 20 以降低负载。 |
-| **base_url** | `str` | 否 | `...` | **Eva API 地址**。
默认为 `https://eva.autoarkai.com`,用于常规 HTTP 请求。 |
-| **wss_url** | `str` | 否 | `...` | **WebSocket 地址**。
默认为 `wss://rtc.autoarkai.com`,用于实时音频流传输。 |
+### 4. 启动客户端
+```bash
+python python_livekit_example.py
+```
---
-### 如何获取音频设备索引
-由于 `PyAudio` 对设备索引非常敏感,我们提供了一个辅助脚本来列出当前可用的设备。
+## 🛠 高级场景与最佳实践
-1. 运行该脚本:
- ```bash
- python list_audio_devices.py
- ```
-2. 找到对应的 `Index` 数字,作为入参传递给EvaClient。
+### 场景 A:选择合适的传输协议
+SDK 提供了两套平行的参考实现,请根据您的硬件算力进行选择:
-## 使用指南
+1. **`eva_client.py` (LiveKit / WebRTC):**
+ * **适用场景:** PC、Mac、树莓派 4 及以上等能跑完整 WebRTC 协议栈的设备。
+ * **优势:** 极致的低延迟、自适应 UDP 拥塞控制、抗弱网。这是**强烈推荐**的首选协议。
+2. **`eva_ws_client.py` (原生 WebSocket):**
+ * **适用场景:** ESP32、MCU 等算力薄弱,跑不动 WebRTC 的低功耗设备。
+ * **优势:** 使用最基础的 WebSocket 裸传 Opus 数据包,极易用 C/C++ 移植到单片机上。
-运行主程序:
+这两种协议底层共享同一套 RTVI 状态机控制平面,业务代码无需修改即可无缝切换。
-```bash
-python python_example.py
+### 场景 B:本地唤醒与意图无缝衔接
+EVA OS V2 统一使用 **VOSK** 作为本地唤醒引擎(原生支持中英文,低 CPU 开销)。
+
+在 `python_livekit_example.py` 中配置:
+```python
+client = EvaLiveKitClient(
+ # ...
+ wake_word="你好方舟", # 要监听的唤醒词
+ wake_word_target_task="intent_task", # 唤醒后向云端请求切换到的任务分支
+ on_task_change=on_task_change, # 状态切换最终完成后的回调
+)
```
+**最佳实践:** 唤醒后**不要**在本地做任何麦克风静音或截断!请让用户自然、连贯地说话(例如:“你好方舟,播放儿歌”)。SDK 在后台捕获到唤醒词后会自动发起 `task.switch.command` 切换请求,而云端服务具备极强的容错理解能力,能够智能处理带有唤醒前缀的连续语音,并在必要时下发 `task.switch.advice` 建议。您只需监听 `on_task_change` 事件来处理最终的状态流转即可。
-### 预期行为
-1. 客户端启动并初始化音频系统和摄像头。
-2. 请求认证 Token 并连接到 WebRTC 房间。
-3. 激活指定的麦克风设备,开始推送音频流。
-4. 激活指定的摄像头设备,开始推送视频流。
-5. 订阅房间内的音频流,并通过指定的扬声器设备播放。
+### 场景 C:回声抑制与全双工打断 (Barge-in)
+如果您的硬件没有内置硬件级的 AEC(声学回声消除)芯片,SDK 提供了基于音量的纯软件双工打断机制。
-## 故障排除
+```python
+ echo_suppression=True,
+ barge_in_multiplier=1.5,
+ barge_in_offset=500,
+```
+在常规情况下,当云端 AI 正在说话时,SDK 会抑制麦克风的上传以防止回声死循环。但如果用户大声说话(麦克风峰值 > AI音量峰值 * 1.5 + 500),SDK 将触发**强行打断 (Barge-in)**,允许用户在 AI 播报中途直接插话。
-* **安装 `pyaudio` 失败 (fatal error: portaudio.h: No such file):**
- 这是因为缺少系统开发库。请确保运行了 `sudo apt-get install portaudio19-dev`。
+---
+
+## 📚 接口参考字典
-* **`ValueError: DeviceIndexOutOfRange`:**
- `mic_index` 或 `spk_index` 不存在。请重新运行 `list_audio_devices.py` 确认索引。
+### EvaLiveKitClient 与 EvaWebSocketClient 参数
-* **摄像头无法打开:**
- * 检查是否缺少 OpenCV 依赖。
- * 如果设备性能不足,尝试在代码中降低 `video_width`, `video_height` 和 `video_fps`。
+| 参数名 | 类型 | 必填 | 默认值 | 说明 |
+| :--- | :--- | :---: | :--- | :--- |
+| **api_key** | `str` | **是** | 无 | EVA OS 颁发的 Solution API Key。 |
+| **mic_index** | `int` | 否 | `0` | 麦克风的设备索引(通过扫描脚本获取)。 |
+| **spk_index** | `int` | 否 | `0` | 扬声器的设备索引。 |
+| **mic_sample_rate**| `int` | 否 | `48000` | 麦克风原生采样率。云端要求16kHz,SDK内部会自动做重采样。 |
+| **spk_sample_rate**| `int` | 否 | `48000` | 扬声器原生采样率。 |
+| **channels** | `int` | 否 | `1` | 通道数(1代表单声道)。 |
+| **frame_duration_ms**| `int`| 否 | `60` | 每次打包发送的 Opus 音频帧长(毫秒)。 |
+| **wake_word** | `str` | 否 | `None` | 要在后台监听的本地唤醒词。 |
+| **wake_word_model_path**| `str` | 否 | `None` | 自定义 VOSK 模型路径。为 None 时会自动下载轻量级中文模型。 |
+| **wake_word_target_task**| `str`| 否 | `None` | 命中唤醒词后,向云端申请切入的任务分支 ID。 |
+| **on_task_change**| `callable`| 否| `None` | 当状态切换最终完成后的回调函数。 |
+
+### WebSocket 资源回收说明
+如果您使用轻量级的 WebSocket 协议 (`eva_ws_client.py`),在程序退出时必须通知云端释放资源。
+SDK 内部已封装该逻辑,在优雅退出时会自动发起 `DELETE /api/solution/chat-room-ws` 请求(携带 `{session_id}`)断开连接。
diff --git a/clients/python-common-sdk/README.md b/clients/python-common-sdk/README.md
index 9c3922f..4966b56 100644
--- a/clients/python-common-sdk/README.md
+++ b/clients/python-common-sdk/README.md
@@ -1,111 +1,138 @@
-# EVA OS Python Client
+# EVA OS Python Client SDK
-This project is the official Python reference implementation for connecting to the EVA OS Real-Time Multimodal AI Service. It demonstrates how to authenticate via API, establish a WebRTC connection using LiveKit, and handle real-time bidirectional audio and video streams.
+This project is the official Python reference implementation for connecting hardware devices to the **EVA OS Real-Time Multimodal AI Service**. It demonstrates how to authenticate, establish low-latency connections, and handle real-time bidirectional audio/video streams.
-## Prerequisites
+More importantly, this SDK serves as a definitive guide on how to integrate your custom hardware (from PCs to low-power IoT devices) into the EVA OS V2 **Edge-Cloud Synergy** architecture.
-### System Dependencies
-This project depends on the `PortAudio` library. You **must** install the system-level development headers before installing the Python `PyAudio` library.
+---
-* **Ubuntu / Debian (and embedded systems like Rockchip/Raspberry Pi):**
- ```bash
- sudo apt-get update
- sudo apt-get install libportaudio2 portaudio19-dev
- ```
+## 🏗 Architectural Philosophy: Edge-Cloud Synergy & Edge State Sovereignty
-* **macOS:**
- ```bash
- brew install portaudio
- ```
+Before using this SDK, it is highly recommended to understand the design principles of EVA OS V2. We adopt an architecture called **"Edge State Sovereignty"**.
-* **Windows:**
- Typically, `pip` will automatically install pre-compiled binaries (Wheels).
+### The Division of Labor
-### Python Environment
-* Python 3.11 or higher is required.
+**The Edge (Your Hardware & this SDK)**
+Hardware resources are limited, so the edge should be highly focused. The SDK is only responsible for:
+1. **Media Streaming:** Continuously pushing microphone audio and pulling speaker audio.
+2. **Local Wake Word Detection:** Running a lightweight VOSK model locally to catch wake words (e.g., "Hello Ark") with zero latency.
+3. **State Requests:** The edge **NEVER** processes VAD (Voice Activity Detection) or intent classification. When a wake word or a button press is detected, the SDK simply sends an intent switch request to the cloud.
+4. **Action Execution:** The edge receives real-time device control and emotion messages from the cloud to execute device actions (e.g., volume adjustments) or display UI feedback instantly, completely bypassing the cloud's main business flow.
-## Installation Steps
+**The Cloud (EVA OS / Pipecat Engine)**
+The cloud acts as the central brain. It continuously receives audio and handles:
+1. **VAD & ASR:** Determining when the user stops speaking and converting speech to text.
+2. **LLM Intent Understanding & Business Logic:** Executing complex workflows based on the user's intent.
+3. **Task Orchestration:** The cloud provides intelligent task advice based on voice intent, and validates edge requests. However, **the Edge is the sovereign state owner**. Only after the edge confirms and executes a state switch will the cloud update its mirrored state and seamlessly activate the corresponding AI Agent to take over the conversation.
-1. **Clone the repository:**
- ```bash
- git clone git@github.com:AutoArk/EVA-OS.git
- cd EVA-OS/clients/python-common-sdk
- ```
+This ensures that the state of your hardware device and the cloud engine are always perfectly synchronized via the RTVI control protocol.
+
+---
+
+## 🚀 Quick Start
+
+### 1. System Dependencies
+This project depends on the `PortAudio` library. You **must** install the system-level development headers before installing the Python dependencies.
-2. **Create and activate a virtual environment (recommended):**
+* **Ubuntu / Debian (and embedded systems like Raspberry Pi):**
```bash
- python3 -m venv venv
- source venv/bin/activate
+ sudo apt-get update
+ sudo apt-get install libportaudio2 portaudio19-dev
```
-
-3. **Install Python dependencies:**
+* **macOS:**
```bash
- pip install livekit requests pyaudio numpy python-dotenv opencv-python
+ brew install portaudio
```
-## Configuration
-
-Create a `.env` file in the project root directory.
+### 2. Environment Setup
+Clone the repository and install dependencies:
+```bash
+python3 -m venv venv
+source venv/bin/activate
+pip install -r requirements.txt
+# Optional: Install OpenCV if you need video capture
+pip install opencv-python
+```
+Create a `.env` file in the project root:
```ini
-# .env file
-# Required: API Key
-EVA_API_KEY=YOUR_ACTUAL_API_KEY_HERE
+# Your API Key generated from the EVA OS Dashboard
+EVA_API_KEY=sk-your-api-key-here
```
-### EvaClient Parameter Details
+### 3. Finding Your Audio Devices (Crucial!)
+`PyAudio` is extremely sensitive to channel misconfigurations (e.g., trying to record from a speaker will result in an `Invalid number of channels` error).
+Run our helper script to find the correct `Index`:
+```bash
+python list_audio_devices.py
+```
+Note down the `Index` for your microphone and speaker, and use them to configure `mic_index` and `spk_index` in the example scripts.
-| Parameter | Type | Required | Default | Description |
-| :--- | :--- | :---: | :--- | :--- |
-| **api_key** | `str` | **Yes** | None | **Authentication Key**.
Credential used to access Eva API services. |
-| **mic_index** | `int` | No | `0` | **Microphone Device Index**.
ID of the input device used for recording. `0` usually represents the system default microphone. |
-| **spk_index** | `int` | No | `0` | **Speaker Device Index**.
ID of the output device used for playback. `0` usually represents the system default speaker. |
-| **mic_sample_rate** | `int` | No | `48000` | **Microphone Sample Rate (Hz)**.
Sampling frequency for recording. Defaults to 48kHz. |
-| **spk_sample_rate** | `int` | No | `48000` | **Speaker Sample Rate (Hz)**.
Sampling frequency for playback. Defaults to 48kHz. |
-| **mic_channels** | `int` | No | `1` | **Microphone Channels**.
`1` for Mono, `2` for Stereo. |
-| **spk_channels** | `int` | No | `1` | **Speaker Channels**.
`1` for Mono, `2` for Stereo. |
-| **frame_size_ms** | `int` | No | `60` | **Audio Frame Duration (ms)**.
Time length of audio data blocks processed or transmitted per cycle. Affects latency and packet size. |
-| **camera_index** | `int` | No | `0` | **Camera Device Index**.
Usually `0` represents the system default cameras. If multiple cameras exist, try sequentially. |
-| **video_width** | `int` | No | `640` | **Video Width**.
Horizontal pixel count for camera capture. It is recommended not to set this too high on ARM devices. |
-| **video_height** | `int` | No | `480` | **Video Height**.
Vertical pixel count for camera capture. |
-| **video_fps** | `int` | No | `30` | **Video Frame Rate**.
Frames per second (FPS). Recommended to set to 15 or 20 on low-performance devices to reduce load. |
-| **base_url** | `str` | No | `...` | **Eva API Address**.
Defaults to `https://eva.autoarkai.com`, used for standard HTTP requests. |
-| **wss_url** | `str` | No | `...` | **WebSocket Address**.
Defaults to `wss://rtc.autoarkai.com`, used for real-time audio stream transmission. |
+### 4. Run the Client
+```bash
+python python_livekit_example.py
+```
---
-### How to Get Audio Device Indices
-Since `PyAudio` is very sensitive to device indices, we provide a helper script to list currently available devices.
+## 🛠 Advanced Scenarios & Best Practices
-1. Run the script:
- ```bash
- python list_audio_devices.py
- ```
-2. Find the corresponding `Index` number and pass it as an argument to `EvaClient`.
+### Scenario A: Choosing the Right Transport Protocol
+This SDK provides two reference implementations depending on your hardware's capabilities:
-## Usage Guide
+1. **`eva_client.py` (LiveKit / WebRTC):**
+ * **Use case:** PCs, Raspberry Pi 4+, or any device that can run a full WebRTC stack.
+ * **Pros:** Ultra-low latency, adaptive UDP bitrate, built-in network recovery. This is the **strongly recommended** protocol.
+2. **`eva_ws_client.py` (Raw WebSocket):**
+ * **Use case:** Low-power MCUs (ESP32) or constrained environments where WebRTC is too heavy.
+ * **Pros:** Sends raw Opus frames over standard WebSocket. Easy to port to C/C++.
-Run the main program:
+Both protocols share the exact same RTVI control plane. Your business logic doesn't need to change if you switch transports.
-```bash
-python python_example.py
+### Scenario B: Wake Word & Seamless Intent Handoff
+EVA OS V2 uses **VOSK** as the unified local wake word engine (supports Chinese/English natively with low CPU footprint).
+
+In `python_livekit_example.py`:
+```python
+client = EvaLiveKitClient(
+ # ...
+ wake_word="你好方舟", # The wake word to listen for
+ wake_word_target_task="intent_task", # Request cloud to switch to this task
+ on_task_change=on_task_change, # Callback when a task switch is finalized
+)
```
+**Best Practice:** Do NOT attempt to mute the microphone after a wake word is detected. Speak naturally ("Hello Ark, play some music"). The local VOSK engine will flag the wake word and send a `task.switch.command` silently in the background. The cloud service has excellent fault-tolerance and is designed to process continuous speech, intelligently understanding your intent and providing `task.switch.advice`. You simply need to listen for the `on_task_change` callback to handle the final state change.
-### Expected Behavior
-1. The client starts and initializes the audio system and camera.
-2. Requests an authentication Token and connects to the WebRTC room.
-3. Activates the specified microphone device and starts publishing the audio stream.
-4. Activates the specified camera device and starts publishing the video stream.
-5. Subscribes to audio streams in the room and plays them through the specified speaker device.
+### Scenario C: Echo Suppression & Voice Barge-in (Full Duplex)
+If your hardware lacks hardware-level AEC (Acoustic Echo Cancellation), the SDK provides a volume-based software suppression mechanism.
-## Troubleshooting
+```python
+ echo_suppression=True,
+ barge_in_multiplier=1.5,
+ barge_in_offset=500,
+```
+When the bot is speaking through the speaker, the microphone is normally suppressed to prevent echo loops. However, if the user speaks loudly enough (Mic Volume > Bot Volume * 1.5 + 500), the SDK triggers a **Barge-in**, allowing the user to forcefully interrupt the AI's playback.
-* **Installation of `pyaudio` failed (fatal error: portaudio.h: No such file):**
- This is due to missing system development libraries. Ensure you have run `sudo apt-get install portaudio19-dev`.
+---
+
+## 📚 API Reference
-* **`ValueError: DeviceIndexOutOfRange`:**
- The `mic_index` or `spk_index` does not exist. Please re-run `list_audio_devices.py` to confirm the index.
+### EvaLiveKitClient & EvaWebSocketClient Parameters
-* **Camera fails to open:**
- * Check if OpenCV dependencies are missing.
- * If device performance is insufficient, try lowering `video_width`, `video_height`, and `video_fps` in the code.
+| Parameter | Type | Required | Default | Description |
+| :--- | :--- | :---: | :--- | :--- |
+| **api_key** | `str` | **Yes** | None | Your EVA OS Solution API Key. |
+| **mic_index** | `int` | No | `0` | Microphone device ID (from `list_audio_devices.py`). |
+| **spk_index** | `int` | No | `0` | Speaker device ID. |
+| **mic_sample_rate**| `int` | No | `48000` | Native mic sample rate. Cloud expects 16kHz, SDK automatically resamples. |
+| **spk_sample_rate**| `int` | No | `48000` | Native speaker sample rate. |
+| **channels** | `int` | No | `1` | Audio channels (1=Mono, 2=Stereo). |
+| **frame_duration_ms**| `int`| No | `60` | Opus audio frame duration. |
+| **wake_word** | `str` | No | `None` | Local wake word to trigger VOSK background detection. |
+| **wake_word_model_path**| `str` | No | `None` | Path to custom VOSK model. If None, auto-downloads a lightweight Chinese model. |
+| **wake_word_target_task**| `str`| No | `None` | The task to request from the cloud when wake word hits. |
+| **on_task_change**| `callable`| No| `None` | Hook triggered when a task switch is finalized. |
+
+### WebSocket Disconnection Flow
+If using the WebSocket transport (`eva_ws_client.py`), you must explicitly terminate the session to release cloud resources.
+The SDK automatically sends a `DELETE /api/solution/chat-room-ws` request with your `{session_id}` upon graceful shutdown.
diff --git a/clients/python-common-sdk/client/__init__.py b/clients/python-common-sdk/client/__init__.py
new file mode 100644
index 0000000..912256a
--- /dev/null
+++ b/clients/python-common-sdk/client/__init__.py
@@ -0,0 +1 @@
+# empty init
diff --git a/clients/python-common-sdk/client/base.py b/clients/python-common-sdk/client/base.py
new file mode 100644
index 0000000..02e0a3c
--- /dev/null
+++ b/clients/python-common-sdk/client/base.py
@@ -0,0 +1,105 @@
+import abc
+import asyncio
+import logging
+from typing import Optional, Callable
+import pyaudio
+
+from shared.fsm import RTVITaskNegotiator
+from shared.wake_word import WakeWordEngine, WakeWordEvent
+from shared.audio import find_audio_device_index
+
+logger = logging.getLogger("BaseEvaClient")
+
+class BaseEvaClient(abc.ABC):
+ """
+ Base class for Eva Clients, encapsulating common properties and logic.
+ """
+
+ def __init__(
+ self,
+ api_key: str,
+ base_url: str = "https://eva.autoarkai.com",
+ mic_index: Optional[int] = 0,
+ spk_index: Optional[int] = 0,
+ mic_sample_rate: int = 48000,
+ spk_sample_rate: int = 48000,
+ channels: int = 1,
+ frame_duration_ms: int = 60,
+ camera_index: Optional[int] = 0,
+ video_width: int = 640,
+ video_height: int = 480,
+ video_fps: int = 30,
+ wake_word_engine: Optional[WakeWordEngine] = None,
+ wake_word_target_task: Optional[str] = None,
+ on_task_change: Optional[Callable[[str, str], None]] = None,
+ ):
+ if not api_key:
+ raise ValueError("api_key is required")
+ if not base_url:
+ raise ValueError("base_url is required")
+
+ self.api_key = api_key
+ self.base_url = base_url.rstrip("/")
+
+ # Audio Config
+ self.mic_sample_rate = mic_sample_rate
+ self.spk_sample_rate = spk_sample_rate
+ self.channels = channels
+ self.frame_duration_ms = frame_duration_ms
+ self.mic_frame_size = int(mic_sample_rate * frame_duration_ms / 1000)
+ self.spk_frame_size = int(spk_sample_rate * frame_duration_ms / 1000)
+
+ # Video Config
+ self.camera_index = camera_index
+ self.video_width = video_width
+ self.video_height = video_height
+ self.video_fps = video_fps
+
+ self.pa = pyaudio.PyAudio()
+ self.mic_index = find_audio_device_index(self.pa, mic_index, is_input=True)
+ self.spk_index = find_audio_device_index(self.pa, spk_index, is_input=False)
+
+ self._shutdown_event = asyncio.Event()
+
+ # Task FSM
+ self.task_negotiator = RTVITaskNegotiator(on_task_change=on_task_change)
+
+ # Wake Word
+ self.wake_word_engine = wake_word_engine
+ self._wake_word_target_task = wake_word_target_task
+ if self.wake_word_engine:
+ self.wake_word_engine.set_callback(self._on_wake_word_detected)
+
+ def _on_wake_word_detected(self, event: WakeWordEvent):
+ if self._wake_word_target_task:
+ logger.info(f"Wake word '{event.keyword}' detected (confidence: {event.confidence}), requesting task switch to {self._wake_word_target_task}")
+ self.switch_task(
+ target_task=self._wake_word_target_task,
+ reason=f"wake_word_detected:{event.keyword}"
+ )
+
+ def switch_task(self, target_task: str, reason: str = "manual_request"):
+ """Manually trigger a hard task switch from the edge."""
+ command = self.task_negotiator.request_switch(
+ target_task=target_task,
+ reason=reason
+ )
+ if command:
+ self._send_command_async(command)
+ logger.info(f"Sent task.switch.command to cloud for {target_task}")
+
+ @abc.abstractmethod
+ def _send_command_async(self, command: dict):
+ """Send command to server asynchronously."""
+ pass
+
+ @abc.abstractmethod
+ async def run(self):
+ """Run the client."""
+ pass
+
+ def stop(self):
+ """Stop the client."""
+ self._shutdown_event.set()
+ if self.wake_word_engine:
+ self.wake_word_engine.stop()
diff --git a/clients/python-common-sdk/eva_client.py b/clients/python-common-sdk/client/livekit_client.py
similarity index 64%
rename from clients/python-common-sdk/eva_client.py
rename to clients/python-common-sdk/client/livekit_client.py
index 1247870..639d171 100644
--- a/clients/python-common-sdk/eva_client.py
+++ b/clients/python-common-sdk/client/livekit_client.py
@@ -1,121 +1,70 @@
+from __future__ import annotations
+
import asyncio
+import json
import logging
-import requests
+import uuid
+from typing import Optional, Callable
+
import numpy as np
import pyaudio
-import queue
-import cv2
-from livekit import rtc
+import requests
+
+try:
+ from livekit import rtc
+except ImportError:
+ rtc = None
+
+from shared.audio import AudioOutputBuffer, find_audio_device_index, AudioConsumer
+from shared.fsm import RTVITaskNegotiator
+from shared.wake_word import WakeWordEngine
+from shared.protocol import wrap_rtvi_envelope
+from .base import BaseEvaClient
# --- Logging Configuration ---
logging.basicConfig(
level=logging.INFO, format="%(asctime)s - [%(name)s] %(levelname)s - %(message)s"
)
-logger = logging.getLogger("EvaClient")
-
-
-class AudioOutputBuffer:
- """
- Streaming buffer designed to handle the impedance mismatch between
- fixed-size network packets and variable-size hardware callback requests.
- """
-
- def __init__(self):
- self.queue = queue.Queue(maxsize=200)
- self.remainder = None
-
- def put(self, data: bytes):
- """
- Converts raw bytes to a numpy array and pushes it into the thread-safe queue.
- """
- np_data = np.frombuffer(data, dtype=np.int16)
- self.queue.put(np_data)
+logger = logging.getLogger("EvaLiveKitClient")
- def get_chunk(self, frames_needed):
- """
- Retrieves the exact number of frames requested by the hardware.
- Returns data as bytes.
- """
- out_list = []
- frames_collected = 0
-
- # 1. Process data remaining from the previous callback, if any
- if self.remainder is not None:
- n = len(self.remainder)
- if n > frames_needed:
- # Remainder is larger than needed; take what is required and save the rest
- out_list.append(self.remainder[:frames_needed])
- self.remainder = self.remainder[frames_needed:]
- return np.concatenate(out_list).tobytes()
- else:
- # Remainder is smaller or equal; consume it entirely
- out_list.append(self.remainder)
- frames_collected += n
- self.remainder = None
- # 2. Fetch new packets from the queue until the required frame count is met
- while frames_collected < frames_needed:
- try:
- new_packet = self.queue.get_nowait()
-
- # Ensure shape consistency
- packet_len = len(new_packet)
- needed = frames_needed - frames_collected
-
- if packet_len > needed:
- # Packet is larger than remaining space; slice and save overflow
- out_list.append(new_packet[:needed])
- self.remainder = new_packet[needed:]
- frames_collected += needed
- else:
- # Packet fits entirely
- out_list.append(new_packet)
- frames_collected += packet_len
-
- except queue.Empty:
- # Queue is empty; fill with silence to prevent hardware underrun artifacts
- needed = frames_needed - frames_collected
- silence = np.zeros(needed, dtype=np.int16)
- out_list.append(silence)
- frames_collected += needed
-
- # Concatenate all segments and convert back to bytes for PyAudio
- return np.concatenate(out_list).tobytes()
-
-
-class EvaClient:
+class EvaLiveKitClient(BaseEvaClient):
"""
- Initializes the Eva client.
+ Initializes the Eva LiveKit client.
Args:
- api_key (str):
+ api_key (str):
Required. The Eva API key.
- mic_index (int, optional):
+ mic_index (int, optional):
Index of the microphone device. Defaults to 0.
- spk_index (int, optional):
+ spk_index (int, optional):
Index of the speaker device. Defaults to 0.
- mic_sample_rate (int, optional):
+ mic_sample_rate (int, optional):
Microphone input sample rate in Hz. Defaults to 48000.
- spk_sample_rate (int, optional):
+ spk_sample_rate (int, optional):
Speaker output sample rate in Hz. Defaults to 48000.
- mic_channels (int, optional):
- Number of microphone input channels. Defaults to 1.
- spk_channels (int, optional):
- Number of speaker output channels. Defaults to 1.
- frame_size_ms (int, optional):
+ channels (int, optional):
+ Number of audio channels. Defaults to 1.
+ frame_duration_ms (int, optional):
Duration of a single audio frame in milliseconds. Defaults to 60ms.
- camera_index (int, optional):
+ camera_index (int, optional):
Index of the camera device. Defaults to 0.
- video_width (int, optional):
+ video_width (int, optional):
Video capture width. Defaults to 640.
- video_height (int, optional):
+ video_height (int, optional):
Video capture height. Defaults to 480.
- video_fps (int, optional):
+ video_fps (int, optional):
Video frame rate (FPS). Defaults to 30.
- base_url (str, optional):
+ base_url (str, optional):
Base URL for the HTTP API service.
- wss_url (str, optional):
+ wss_url (str, optional):
WebSocket URL for RTC.
+ wake_word_engine (WakeWordEngine, optional):
+ Wake word engine for voice activation.
+ wake_word_target_task (str, optional):
+ Target task to switch to when wake word is detected.
+ on_task_change (callable, optional):
+ Callback function(old_task, new_task) when task changes.
"""
def __init__(
self,
@@ -124,45 +73,50 @@ def __init__(
spk_index=0,
mic_sample_rate=48000,
spk_sample_rate=48000,
- mic_channels=1,
- spk_channels=1,
- frame_size_ms=60,
+ channels=1,
+ frame_duration_ms=60,
camera_index=0,
video_width=640,
video_height=480,
video_fps=30,
base_url="https://eva.autoarkai.com",
wss_url="wss://rtc.autoarkai.com",
+ wake_word_engine: Optional[WakeWordEngine] = None,
+ wake_word_target_task: Optional[str] = None,
+ on_task_change: Optional[Callable[[str, str], None]] = None,
):
- self.api_key = api_key
- self.base_url = base_url
+ super().__init__(
+ api_key=api_key,
+ base_url=base_url,
+ mic_index=mic_index,
+ spk_index=spk_index,
+ mic_sample_rate=mic_sample_rate,
+ spk_sample_rate=spk_sample_rate,
+ channels=channels,
+ frame_duration_ms=frame_duration_ms,
+ camera_index=camera_index,
+ video_width=video_width,
+ video_height=video_height,
+ video_fps=video_fps,
+ wake_word_engine=wake_word_engine,
+ wake_word_target_task=wake_word_target_task,
+ on_task_change=on_task_change,
+ )
+
+ if rtc is None:
+ raise ImportError(
+ "livekit is not installed. Please install it using: "
+ "pip install livekit livekit-api"
+ )
+
self.wss_url = wss_url
- # Audio Configuration
- self.pa = pyaudio.PyAudio()
- self.mic_sample_rate = mic_sample_rate
- self.spk_sample_rate = spk_sample_rate
- self.mic_channels = mic_channels
- self.spk_channels = spk_channels
- self.frame_size_ms = frame_size_ms
-
- # Video Configuration
- self.camera_index = camera_index
- self.video_width = video_width
- self.video_height = video_height
- self.video_fps = video_fps
-
- # Resolve audio device indices
- self.mic_index = self._find_device_index(mic_index, input=True)
- self.spk_index = self._find_device_index(spk_index, input=False)
-
- if not all([self.base_url, self.wss_url, self.api_key]):
+ if not self.wss_url:
raise ValueError("Missing required environment variables.")
self.audio_buffer = AudioOutputBuffer()
self.room = None
- self._shutdown_event = asyncio.Event()
-
+
# Sources & Tracks
self.mic_source = None
self.video_source = None
@@ -171,40 +125,33 @@ def __init__(
# PyAudio streams
self.input_stream = None
self.output_stream = None
+ self._loop = None
- def _find_device_index(self, value, input=True):
- """
- Resolves the audio device index.
- """
- if value is None:
- return None
-
- if type(value) is str and value.strip() == "":
- return None
-
- try:
- return int(value)
- except ValueError:
- pass
-
- target = value.strip()
- count = self.pa.get_device_count()
- logger.info(f"Searching for audio device matching: '{target}'...")
-
- for i in range(count):
- info = self.pa.get_device_info_by_index(i)
- name = info.get("name", "")
- if input and info.get("maxInputChannels") == 0:
- continue
- if not input and info.get("maxOutputChannels") == 0:
- continue
+ # RTVI data channel for task state coordination
+ self._rtvi_topic = "task-ir-control"
- if target in name:
- logger.info(f"Found device '{name}' at index {i}")
- return i
+ def _send_command_async(self, command: dict):
+ if getattr(self, '_loop', None):
+ asyncio.run_coroutine_threadsafe(
+ self._send_rtvi_message(command), self._loop
+ )
+ else:
+ logger.error("No event loop found to send RTVI message")
- logger.warning(f"Device '{target}' not found. Using system default.")
- return None
+ async def _send_rtvi_message(self, message: dict):
+ """Send RTVI message via LiveKit data channel."""
+ if not self.room:
+ return
+
+ envelope = wrap_rtvi_envelope(message, self._rtvi_topic)
+ payload = json.dumps(envelope).encode("utf-8")
+
+ # Send via data channel
+ await self.room.local_participant.publish_data(
+ payload,
+ reliable=True,
+ topic=self._rtvi_topic,
+ )
def get_room_token(self) -> str:
url = f"{self.base_url}/api/solution/chat-room"
@@ -224,9 +171,11 @@ def get_room_token(self) -> str:
raise
async def run(self):
+ self._loop = asyncio.get_running_loop()
try:
token = self.get_room_token()
- except Exception:
+ except Exception as e:
+ logger.error(f"Failed to get room token: {e}")
return
self.room = rtc.Room()
@@ -248,6 +197,26 @@ def on_disconnected():
logger.info("Disconnected.")
self._shutdown_event.set()
+ @self.room.on("data_received")
+ def on_data_received(data_packet: rtc.DataPacket):
+ data = data_packet.data
+
+ try:
+ msg_dict = json.loads(data.decode("utf-8"))
+
+ # FSM handles unwrapping the RTVI envelope
+ response = self.task_negotiator.handle_message(msg_dict)
+ if response:
+ if getattr(self, '_loop', None):
+ asyncio.run_coroutine_threadsafe(
+ self._send_rtvi_message(response), self._loop
+ )
+ else:
+ logger.error("No event loop found to send RTVI message")
+ except Exception as e:
+ # 忽略解析非 JSON 或者不相关的数据包
+ logger.error(f"Failed to parse or handle message (ignoring): {e}")
+
logger.info(f"Connecting to {self.wss_url}")
try:
await self.room.connect(self.wss_url, token)
@@ -259,8 +228,16 @@ def on_disconnected():
mic_task = asyncio.create_task(self.publish_microphone())
video_task = asyncio.create_task(self.publish_camera())
+ # Start wake word detection if configured
+ if self.wake_word_engine:
+ self.wake_word_engine.start()
+
await self._shutdown_event.wait()
+ # Stop wake word detection
+ if self.wake_word_engine:
+ self.wake_word_engine.stop()
+
# Cleanup Tasks
if mic_task:
mic_task.cancel()
@@ -286,7 +263,7 @@ def on_disconnected():
await self.room.disconnect()
async def publish_microphone(self):
- self.mic_source = rtc.AudioSource(self.mic_sample_rate, self.mic_channels)
+ self.mic_source = rtc.AudioSource(self.mic_sample_rate, self.channels)
track = rtc.LocalAudioTrack.create_audio_track("mic_track", self.mic_source)
options = rtc.TrackPublishOptions()
options.source = rtc.TrackSource.SOURCE_MICROPHONE
@@ -300,16 +277,21 @@ async def publish_microphone(self):
logger.error(f"Failed to publish microphone: {e}")
return
- frames_per_buffer = int(self.mic_sample_rate * self.frame_size_ms / 1000)
+ frames_per_buffer = int(self.mic_sample_rate * self.frame_duration_ms / 1000)
loop = asyncio.get_running_loop()
# PyAudio Callback
def mic_callback(in_data, frame_count, time_info, status):
+ # Feed audio to wake word runner if enabled
+ if self.wake_word_engine and isinstance(self.wake_word_engine, AudioConsumer):
+ # in_data comes as bytes
+ self.wake_word_engine.add_audio_frame(in_data, sample_rate=self.mic_sample_rate, channels=self.channels)
+
# in_data comes as bytes
audio_frame = rtc.AudioFrame(
data=in_data,
sample_rate=self.mic_sample_rate,
- num_channels=self.mic_channels,
+ num_channels=self.channels,
samples_per_channel=frame_count,
)
asyncio.run_coroutine_threadsafe(
@@ -320,7 +302,7 @@ def mic_callback(in_data, frame_count, time_info, status):
try:
self.input_stream = self.pa.open(
format=pyaudio.paInt16,
- channels=self.mic_channels,
+ channels=self.channels,
rate=self.mic_sample_rate,
input=True,
input_device_index=self.mic_index,
@@ -340,6 +322,12 @@ async def publish_camera(self):
"""
Captures video from the camera using OpenCV and publishes it to the room.
"""
+ try:
+ import cv2
+ except ImportError:
+ logger.error("opencv-python is required for video. Install with: pip install opencv-python")
+ return
+
logger.info(f"Opening camera index {self.camera_index}...")
# Initialize OpenCV VideoCapture
@@ -399,24 +387,24 @@ async def publish_camera(self):
async def handle_audio_output(self, track: rtc.RemoteAudioTrack):
audio_stream = rtc.AudioStream(
- track, sample_rate=self.spk_sample_rate, num_channels=self.spk_channels
+ track, sample_rate=self.spk_sample_rate, num_channels=self.channels
)
logger.info(
f"Speaker stream started. Rate: {self.spk_sample_rate}, Index: {self.spk_index}"
)
- frames_per_buffer = int(self.spk_sample_rate * self.frame_size_ms / 1000)
+ frames_per_buffer = int(self.spk_sample_rate * self.frame_duration_ms / 1000)
# PyAudio Output Callback
def spk_callback(in_data, frame_count, time_info, status):
# Retrieve the exact number of bytes needed from the buffer
- data = self.audio_buffer.get_chunk(frame_count)
+ data = self.audio_buffer.get_chunk(frame_count * self.channels)
return (data, pyaudio.paContinue)
try:
self.output_stream = self.pa.open(
format=pyaudio.paInt16,
- channels=self.spk_channels,
+ channels=self.channels,
rate=self.spk_sample_rate,
output=True,
output_device_index=self.spk_index,
diff --git a/clients/python-common-sdk/client/ws_client.py b/clients/python-common-sdk/client/ws_client.py
new file mode 100644
index 0000000..f5e15b1
--- /dev/null
+++ b/clients/python-common-sdk/client/ws_client.py
@@ -0,0 +1,810 @@
+"""
+Eva WebSocket Client
+====================
+
+Lightweight WebSocket transport client for the Eva platform. Intended for
+embedded devices (ESP32, Raspberry Pi, IoT) and environments where a full
+WebRTC stack is not available.
+
+Protocol:
+ Eva Binary Protocol V1
+ - 9-byte big-endian header: version(u8) | media_type(u8) | stream_id(u8) | sequence(u16) | delta_ms(u32)
+ - Opus audio frames (media_type=0x01) and JPEG video frames (media_type=0x02)
+ - JSON text messages for control and events
+
+Flow:
+ 1. POST /api/solution/chat-room {transport_type: "websocket"}
+ 2. Connect WebSocket: attach_url?token=attach_token
+ 3. Send hello (JSON) with declared streams
+ 4. Receive hello response with selected_streams
+ 5. Publish microphone audio (Opus) + optional camera (JPEG)
+ 6. Receive audio (Opus) from the server, decode, play
+"""
+
+from __future__ import annotations
+
+import asyncio
+import json
+import logging
+import struct
+import time
+import uuid
+from typing import Optional, Callable
+
+import numpy as np
+import pyaudio
+import requests
+import websockets
+
+from shared.audio import AudioOutputBuffer, find_audio_device_index, AudioConsumer
+from shared.fsm import RTVITaskNegotiator
+from shared.wake_word import WakeWordEngine
+from shared.protocol import wrap_rtvi_envelope
+from .base import BaseEvaClient
+
+logging.basicConfig(
+ level=logging.INFO, format="%(asctime)s - [%(name)s] %(levelname)s - %(message)s"
+)
+logger = logging.getLogger("EvaWSClient")
+
+# --- Eva Binary Protocol V1 ---
+HEADER_FORMAT = "!BBBHI"
+HEADER_SIZE = 9
+PROTOCOL_VERSION = 1
+MEDIA_OPUS_AUDIO = 0x01
+MEDIA_JPEG_VIDEO = 0x02
+
+# Default stream IDs (overridden by hello negotiation)
+UPLINK_AUDIO_STREAM_ID = 0
+UPLINK_VIDEO_STREAM_ID = 1
+DOWNLINK_AUDIO_STREAM_ID = 2
+
+
+class EvaWebSocketClient(BaseEvaClient):
+ """
+ Eva platform WebSocket client.
+
+ Uses Eva Binary Protocol V1 over a plain WebSocket connection.
+ Suitable for embedded/IoT devices that cannot run a WebRTC stack.
+
+ Args:
+ api_key: Required. The Eva API key (Solution API Key).
+ base_url: HTTP base URL of the Eva platform.
+ mic_index: Microphone device index or name substring. Defaults to 0.
+ spk_index: Speaker device index or name substring. Defaults to 0.
+ sample_rate: Audio sample rate in Hz. Both mic and speaker use the same rate.
+ Defaults to 16000 (server default for WebSocket transport).
+ channels: Number of audio channels. Defaults to 1 (mono).
+ frame_duration_ms: Duration of a single Opus frame in ms. Defaults to 60.
+ camera_index: Camera device index, or None to disable video.
+ video_width: Video capture width. Defaults to 640.
+ video_height: Video capture height. Defaults to 480.
+ video_fps: Video frame rate. Defaults to 2.
+ """
+
+ def __init__(
+ self,
+ api_key: str,
+ base_url: str = "https://eva.autoarkai.com",
+ mic_index=0,
+ spk_index=0,
+ mic_sample_rate: int = 48000,
+ spk_sample_rate: int = 48000,
+ channels: int = 1,
+ frame_duration_ms: int = 60,
+ camera_index=None,
+ video_width: int = 640,
+ video_height: int = 480,
+ video_fps: int = 2,
+ echo_suppression: bool = True,
+ echo_suppression_timeout_ms: int = 500,
+ barge_in_multiplier: float = 1.5,
+ barge_in_offset: int = 500,
+ wake_word_engine: Optional[WakeWordEngine] = None,
+ wake_word_target_task: Optional[str] = None,
+ on_task_change: Optional[Callable[[str, str], None]] = None,
+ ):
+ super().__init__(
+ api_key=api_key,
+ base_url=base_url,
+ mic_index=mic_index,
+ spk_index=spk_index,
+ mic_sample_rate=mic_sample_rate,
+ spk_sample_rate=spk_sample_rate,
+ channels=channels,
+ frame_duration_ms=frame_duration_ms,
+ camera_index=camera_index,
+ video_width=video_width,
+ video_height=video_height,
+ video_fps=video_fps,
+ wake_word_engine=wake_word_engine,
+ wake_word_target_task=wake_word_target_task,
+ on_task_change=on_task_change,
+ )
+
+ # Runtime state
+ self.audio_buffer = AudioOutputBuffer()
+ self.ws = None
+ self.session_id = None
+ self._connect_time_ms = 0.0
+ self._sequence = 0
+ self._bot_speaking = False # True while bot audio is being received/played
+ self._mic_muted = False # True when mic is auto-muted during bot speech
+ self._last_bot_audio_time = 0.0 # Timestamp of last bot audio frame
+ self._loop = None # Event loop reference (set in run())
+
+ # Echo suppression config
+ self._echo_suppression = echo_suppression
+ self._echo_suppression_timeout_s = echo_suppression_timeout_ms / 1000.0
+
+ # Volume-based barge-in (new)
+ self._bot_volume_peak = 0 # Recent bot audio peak amplitude
+ self._barge_in_multiplier = barge_in_multiplier # User volume must be > bot * this
+ self._barge_in_offset = barge_in_offset # Plus this absolute threshold
+
+ # Opus codec
+ self._encoder = None
+ self._decoder = None
+
+ # PyAudio streams
+ self.input_stream = None
+ self.output_stream = None
+ self.video_cap = None
+
+ # Stream IDs (may be overridden by hello negotiation)
+ self.uplink_audio_stream_id = UPLINK_AUDIO_STREAM_ID
+ self.uplink_video_stream_id = UPLINK_VIDEO_STREAM_ID
+ self.downlink_audio_stream_id = DOWNLINK_AUDIO_STREAM_ID
+
+ def _send_command_async(self, command: dict):
+ if getattr(self, '_loop', None) and self.ws:
+ envelope = wrap_rtvi_envelope(command)
+ asyncio.run_coroutine_threadsafe(
+ self.ws.send(json.dumps(envelope)), self._loop
+ )
+ else:
+ logger.warning("[WakeWord] Cannot send command: ws or loop not available")
+
+ # ------------------------------------------------------------------
+ # Session creation (HTTP)
+ # ------------------------------------------------------------------
+ def get_ws_session(self) -> dict:
+ """
+ Request a WebSocket session from the Eva platform.
+
+ Returns:
+ dict with keys: session_id, attach_url, attach_token
+ """
+ url = f"{self.base_url}/api/solution/chat-room"
+ headers = {"Authorization": f"Bearer {self.api_key}"}
+ body = {"transport_type": "websocket"}
+
+ response = requests.post(url, headers=headers, json=body, timeout=30)
+ response.raise_for_status()
+ data = response.json()
+
+ # Handle both wrapped and unwrapped response shapes
+ if isinstance(data, dict) and "data" in data and isinstance(data["data"], dict):
+ data = data["data"]
+
+ session_id = data.get("sessionId")
+ attach_url = data.get("attachUrl")
+ attach_token = data.get("attachToken")
+
+ if not all([session_id, attach_url, attach_token]):
+ raise ValueError(f"Invalid session response: {data}")
+
+ self.session_id = session_id
+ logger.info(f"Session created: {session_id}")
+ return {
+ "session_id": session_id,
+ "attach_url": attach_url,
+ "attach_token": attach_token,
+ }
+
+ # ------------------------------------------------------------------
+ # Binary protocol helpers
+ # ------------------------------------------------------------------
+ @staticmethod
+ def _pack_header(
+ media_type: int, stream_id: int, sequence: int, delta_ms: int
+ ) -> bytes:
+ return struct.pack(
+ HEADER_FORMAT,
+ PROTOCOL_VERSION,
+ media_type,
+ stream_id & 0xFF,
+ sequence & 0xFFFF,
+ delta_ms & 0xFFFFFFFF,
+ )
+
+ @staticmethod
+ def _unpack_header(data: bytes):
+ if len(data) < HEADER_SIZE:
+ raise ValueError(f"Data too short for header: {len(data)} < {HEADER_SIZE}")
+ version, media_type, stream_id, sequence, delta_ms = struct.unpack(
+ HEADER_FORMAT, data[:HEADER_SIZE]
+ )
+ return version, media_type, stream_id, sequence, delta_ms
+
+ def _next_sequence(self) -> int:
+ seq = self._sequence
+ self._sequence = (self._sequence + 1) & 0xFFFF
+ return seq
+
+ def _delta_ms(self) -> int:
+ return int(time.monotonic() * 1000 - self._connect_time_ms) & 0xFFFFFFFF
+
+ # ------------------------------------------------------------------
+ # Hello handshake
+ # ------------------------------------------------------------------
+ def _build_hello(self) -> dict:
+ device_id = f"eva-ws-py-{uuid.uuid4().hex[:8]}"
+ # Server expects 16kHz audio (we resample from device rate)
+ server_sample_rate = 16000
+
+ streams = [
+ {
+ "stream_id": self.uplink_audio_stream_id,
+ "kind": "audio",
+ "direction": "uplink",
+ "source": "mic_main",
+ "codec": "opus",
+ "sample_rate": server_sample_rate,
+ "channels": self.channels,
+ "frame_duration_ms": self.frame_duration_ms,
+ },
+ {
+ "stream_id": self.downlink_audio_stream_id,
+ "kind": "audio",
+ "direction": "downlink",
+ "source": "tts_main",
+ "codec": "opus",
+ "sample_rate": server_sample_rate,
+ "channels": self.channels,
+ "frame_duration_ms": self.frame_duration_ms,
+ },
+ ]
+
+ if self.camera_index is not None:
+ streams.append(
+ {
+ "stream_id": self.uplink_video_stream_id,
+ "kind": "video",
+ "direction": "uplink",
+ "source": "camera_front",
+ "codec": "jpeg",
+ "width": self.video_width,
+ "height": self.video_height,
+ "fps": self.video_fps,
+ }
+ )
+
+ return {
+ "type": "hello",
+ "version": 1,
+ "device_id": device_id,
+ "transport": "websocket",
+ "role": "primary",
+ "streams": streams,
+ }
+
+ # ------------------------------------------------------------------
+ # Opus codec setup
+ # ------------------------------------------------------------------
+ def _init_opus(self):
+ import opuslib # type: ignore
+
+ # Server expects 16kHz audio, so encoder/decoder use 16kHz
+ # Client resamples between device rate (e.g., 48kHz) and server rate (16kHz)
+ server_sample_rate = 16000
+
+ # Encoder: converts resampled PCM (16kHz) to Opus
+ self._encoder = opuslib.Encoder(
+ server_sample_rate, self.channels, opuslib.APPLICATION_VOIP
+ )
+ # Decoder: converts Opus to PCM at 16kHz, then resampled to device rate
+ self._decoder = opuslib.Decoder(server_sample_rate, self.channels)
+ logger.info(
+ f"Opus codec initialized: server={server_sample_rate}Hz, "
+ f"device_mic={self.mic_sample_rate}Hz, device_spk={self.spk_sample_rate}Hz, "
+ f"{self.channels}ch, mic_frame={self.mic_frame_size}, spk_frame={self.spk_frame_size}"
+ )
+
+ # ------------------------------------------------------------------
+ # Microphone input (PCM -> Opus -> WS)
+ # ------------------------------------------------------------------
+ async def _run_microphone(self, loop: asyncio.AbstractEventLoop):
+ frames_per_buffer = self.mic_frame_size
+ frames_sent = [0]
+ last_log_time = [time.monotonic()]
+
+ # Resample if device rate != server rate (server expects 16kHz)
+ server_sample_rate = 16000
+ need_resample = (self.mic_sample_rate != server_sample_rate)
+
+ if need_resample:
+ from scipy import signal
+ logger.info(f"[Mic] Resampling from {self.mic_sample_rate}Hz to {server_sample_rate}Hz")
+
+ def mic_callback(in_data, frame_count, time_info, status):
+ if self._encoder is None or self.ws is None:
+ return (None, pyaudio.paContinue)
+
+ now = time.monotonic()
+
+ # Detect bot stopped speaking (timeout since last audio frame)
+ if self._echo_suppression and self._bot_speaking:
+ if now - self._last_bot_audio_time > self._echo_suppression_timeout_s:
+ self._bot_speaking = False
+ self._bot_volume_peak = 0 # Reset bot volume when silent
+
+ # Diagnostic: check if input is silence (all zeros or near-zeros)
+ pcm = np.frombuffer(in_data, dtype=np.int16)
+ mic_peak = int(np.max(np.abs(pcm))) if len(pcm) > 0 else 0
+
+ # Feed audio to wake word detector (always, regardless of echo suppression)
+ if self.wake_word_engine and isinstance(self.wake_word_engine, AudioConsumer):
+ self.wake_word_engine.add_audio_frame(in_data, sample_rate=self.mic_sample_rate, channels=self.channels)
+
+ # Echo suppression with volume-based barge-in
+ should_send = True
+ if self._echo_suppression and self._bot_speaking:
+ # Calculate barge-in threshold
+ barge_in_threshold = self._bot_volume_peak * self._barge_in_multiplier + self._barge_in_offset
+
+ if mic_peak > barge_in_threshold:
+ # User is speaking louder than bot echo - allow barge-in
+ if self._mic_muted:
+ self._mic_muted = False
+ logger.info(f"Mic unmuted (barge-in: mic={mic_peak} > threshold={int(barge_in_threshold)})")
+ should_send = True
+ else:
+ # Bot is louder - suppress mic
+ if not self._mic_muted:
+ self._mic_muted = True
+ logger.info("Mic auto-muted (bot speaking)")
+ should_send = False
+ else:
+ # Bot not speaking - always send
+ if self._mic_muted:
+ self._mic_muted = False
+ logger.info("Mic unmuted (bot stopped)")
+
+ if should_send:
+ try:
+ # Resample if needed (device rate -> server rate)
+ if need_resample:
+ # Convert to float for resampling
+ pcm_float = pcm.astype(np.float32) / 32768.0
+ # Calculate target length
+ target_length = int(len(pcm_float) * server_sample_rate / self.mic_sample_rate)
+ # Resample
+ pcm_resampled = signal.resample(pcm_float, target_length)
+ # Convert back to int16
+ pcm_to_encode = (pcm_resampled * 32768.0).astype(np.int16)
+ frame_size_for_encode = target_length
+ else:
+ pcm_to_encode = pcm
+ frame_size_for_encode = self.mic_frame_size
+
+ # opuslib.encode expects raw int16 bytes, NOT a Python list
+ opus_bytes = self._encoder.encode(pcm_to_encode.tobytes(), frame_size_for_encode)
+
+ header = self._pack_header(
+ MEDIA_OPUS_AUDIO,
+ self.uplink_audio_stream_id,
+ self._next_sequence(),
+ self._delta_ms(),
+ )
+ asyncio.run_coroutine_threadsafe(
+ self.ws.send(header + opus_bytes), loop
+ )
+ frames_sent[0] += 1
+ except Exception as e:
+ logger.error(f"Mic encode error: {e}")
+
+ # Log every 5 seconds
+ if now - last_log_time[0] >= 5.0:
+ if self._mic_muted:
+ mute_status = " [MUTED]"
+ else:
+ mute_status = ""
+ logger.info(
+ f"Mic: sent {frames_sent[0]} frames in last 5s, "
+ f"peak amplitude={mic_peak}, bot_peak={self._bot_volume_peak}{mute_status}"
+ )
+ frames_sent[0] = 0
+ last_log_time[0] = now
+
+ return (None, pyaudio.paContinue)
+
+ try:
+ self.input_stream = self.pa.open(
+ format=pyaudio.paInt16,
+ channels=self.channels,
+ rate=self.mic_sample_rate,
+ input=True,
+ input_device_index=self.mic_index,
+ frames_per_buffer=frames_per_buffer,
+ stream_callback=mic_callback,
+ )
+ self.input_stream.start_stream()
+ logger.info(
+ f"Mic started: rate={self.mic_sample_rate}, index={self.mic_index}, "
+ f"frame={self.frame_duration_ms}ms"
+ )
+ except Exception as e:
+ logger.error(f"Failed to start mic: {e}")
+ return
+
+ while not self._shutdown_event.is_set():
+ await asyncio.sleep(1)
+
+ # ------------------------------------------------------------------
+ # Camera input (BGR -> JPEG -> WS)
+ # ------------------------------------------------------------------
+ async def _run_camera(self):
+ try:
+ import cv2 # type: ignore
+ except ImportError:
+ logger.error("opencv-python is required for video. Install with: pip install opencv-python")
+ return
+
+ self.video_cap = cv2.VideoCapture(self.camera_index)
+ if not self.video_cap.isOpened():
+ logger.error(f"Could not open camera {self.camera_index}")
+ return
+
+ self.video_cap.set(cv2.CAP_PROP_FRAME_WIDTH, self.video_width)
+ self.video_cap.set(cv2.CAP_PROP_FRAME_HEIGHT, self.video_height)
+
+ interval = 1.0 / max(self.video_fps, 1)
+ logger.info(
+ f"Camera started: {self.video_width}x{self.video_height}@{self.video_fps}fps, "
+ f"index={self.camera_index}"
+ )
+
+ try:
+ while not self._shutdown_event.is_set():
+ ret, frame = self.video_cap.read()
+ if not ret:
+ await asyncio.sleep(0.1)
+ continue
+
+ # Encode BGR -> JPEG bytes
+ ok, buf = cv2.imencode(".jpg", frame, [cv2.IMWRITE_JPEG_QUALITY, 70])
+ if not ok:
+ continue
+
+ header = self._pack_header(
+ MEDIA_JPEG_VIDEO,
+ self.uplink_video_stream_id,
+ self._next_sequence(),
+ self._delta_ms(),
+ )
+
+ if self.ws is not None:
+ await self.ws.send(header + buf.tobytes())
+
+ await asyncio.sleep(interval)
+ except asyncio.CancelledError:
+ pass
+ finally:
+ if self.video_cap is not None and self.video_cap.isOpened():
+ self.video_cap.release()
+
+ # ------------------------------------------------------------------
+ # Speaker output (WS -> Opus decode -> PCM)
+ # ------------------------------------------------------------------
+ async def _run_speaker(self, loop: asyncio.AbstractEventLoop):
+ frames_per_buffer = self.spk_frame_size
+
+ def spk_callback(in_data, frame_count, time_info, status):
+ data = self.audio_buffer.get_chunk(frame_count * self.channels)
+ return (data, pyaudio.paContinue)
+
+ # Try the configured device first; if it fails (e.g. "Invalid number
+ # of channels"), fall back to the system default (index=None).
+ for device_index in (self.spk_index, None):
+ try:
+ self.output_stream = self.pa.open(
+ format=pyaudio.paInt16,
+ channels=self.channels,
+ rate=self.spk_sample_rate,
+ output=True,
+ output_device_index=device_index,
+ frames_per_buffer=frames_per_buffer,
+ stream_callback=spk_callback,
+ )
+ self.output_stream.start_stream()
+ logger.info(
+ f"Speaker started: rate={self.spk_sample_rate}, "
+ f"index={device_index}"
+ )
+ break
+ except Exception as e:
+ if device_index is None:
+ logger.error(f"Failed to start speaker (default): {e}")
+ return
+ logger.warning(
+ f"Speaker index={device_index} failed ({e}), "
+ f"retrying with system default..."
+ )
+
+ while not self._shutdown_event.is_set():
+ await asyncio.sleep(1)
+
+ # ------------------------------------------------------------------
+ # Message receive loop
+ # ------------------------------------------------------------------
+ async def _run_receiver(self):
+ # Diagnostics: count received frames over time
+ frames_received = [0]
+ last_log_time = [time.monotonic()]
+
+ try:
+ async for msg in self.ws:
+ if isinstance(msg, bytes):
+ self._handle_binary(msg, frames_received)
+ elif isinstance(msg, str):
+ # _handle_text may return a message to send back (e.g., commit)
+ reply = self._handle_text(msg)
+ if reply and self.ws:
+ # Wrap TaskIR messages in RTVI envelope
+ if reply.get("type") in ("task.switch.commit", "task.switch.command"):
+ reply = wrap_rtvi_envelope(reply)
+ await self.ws.send(json.dumps(reply))
+
+ # Periodic log
+ now = time.monotonic()
+ if now - last_log_time[0] >= 5.0:
+ logger.info(
+ f"Recv: {frames_received[0]} audio frames in last 5s"
+ )
+ frames_received[0] = 0
+ last_log_time[0] = now
+ except websockets.ConnectionClosed as e:
+ logger.info(f"WebSocket closed: code={e.code}, reason={e.reason}")
+ except asyncio.CancelledError:
+ return
+ except Exception as e:
+ logger.error(f"Receiver error: {e}")
+ finally:
+ self._shutdown_event.set()
+
+ def _handle_binary(self, data: bytes, frames_received=None):
+ if len(data) < HEADER_SIZE:
+ return
+ try:
+ _version, media_type, stream_id, _sequence, _delta_ms = self._unpack_header(data)
+ except ValueError as e:
+ logger.error(f"Failed to unpack binary header: {e}")
+ return
+
+ payload = data[HEADER_SIZE:]
+ if media_type == MEDIA_OPUS_AUDIO and self._decoder is not None:
+ if stream_id != self.downlink_audio_stream_id:
+ return
+ try:
+ # Server sends 16kHz audio, decode it
+ server_sample_rate = 16000
+ server_frame_size = int(server_sample_rate * self.frame_duration_ms / 1000)
+ pcm = self._decoder.decode(payload, server_frame_size)
+
+ # Resample if device rate != server rate
+ if self.spk_sample_rate != server_sample_rate:
+ from scipy import signal
+ pcm_array = np.frombuffer(pcm, dtype=np.int16)
+ target_length = int(len(pcm_array) * self.spk_sample_rate / server_sample_rate)
+ pcm_resampled = signal.resample(pcm_array.astype(np.float32), target_length)
+ pcm = pcm_resampled.astype(np.int16).tobytes()
+
+ self.audio_buffer.put(pcm)
+ if frames_received is not None:
+ frames_received[0] += 1
+ # Mark bot as speaking for echo suppression
+ self._bot_speaking = True
+ self._last_bot_audio_time = time.monotonic()
+ # Track bot volume peak for barge-in detection
+ pcm_array = np.frombuffer(pcm, dtype=np.int16)
+ if len(pcm_array) > 0:
+ self._bot_volume_peak = int(np.max(np.abs(pcm_array)))
+ except Exception as e:
+ logger.error(f"Opus decode error: {e}")
+ elif media_type == MEDIA_JPEG_VIDEO:
+ logger.debug(f"Received JPEG frame: {len(payload)} bytes")
+
+ def _handle_text(self, text: str) -> Optional[dict]:
+ """
+ Handle incoming text message.
+
+ Returns:
+ Message dict to send back (e.g., task.switch.commit), or None
+ """
+ try:
+ msg = json.loads(text)
+ msg_type = msg.get("type")
+ if msg_type == "hello":
+ logger.info(f"Server hello: version={msg.get('version')}")
+ self._apply_selected_streams(msg.get("selected_streams"))
+ return None
+
+ # Pass to TaskFSM for task state coordination
+ if msg_type in ("task.switch.advice", "task.switch.result", "system_config") or \
+ "approved" in msg or "suggested_task" in msg:
+ commit_msg = self.task_negotiator.handle_message(msg)
+ if commit_msg:
+ return commit_msg
+ return None
+
+ logger.debug(f"Text message: {msg_type}")
+ return None
+ except Exception as e:
+ logger.error(f"Text parse error: {e}")
+ return None
+
+ def _apply_selected_streams(self, selected_streams):
+ if not isinstance(selected_streams, list):
+ return
+ for s in selected_streams:
+ kind = s.get("kind")
+ direction = s.get("direction")
+ stream_id = s.get("stream_id")
+ if stream_id is None:
+ continue
+ if kind == "audio" and direction == "uplink":
+ self.uplink_audio_stream_id = stream_id
+ elif kind == "video" and direction == "uplink":
+ self.uplink_video_stream_id = stream_id
+ elif kind == "audio" and direction == "downlink":
+ self.downlink_audio_stream_id = stream_id
+
+ # ------------------------------------------------------------------
+ # Disconnect (HTTP cleanup)
+ # ------------------------------------------------------------------
+ def disconnect_ws_session(self):
+ """Notify the server to clean up the WebSocket session."""
+ if not self.session_id:
+ return
+ url = f"{self.base_url}/api/solution/chat-room-ws"
+ headers = {"Authorization": f"Bearer {self.api_key}"}
+ try:
+ r = requests.delete(
+ url, headers=headers, json={"session_id": self.session_id}, timeout=10
+ )
+ if r.ok:
+ logger.info(f"Session {self.session_id} disconnected on server")
+ else:
+ logger.warning(f"Server disconnect failed: {r.status_code}")
+ except Exception as e:
+ logger.error(f"Failed to notify server disconnect: {e}")
+
+ # ------------------------------------------------------------------
+ # Main entry point
+ # ------------------------------------------------------------------
+ async def run(self):
+ """Connect and run the client until stopped or disconnected."""
+ # 1. Create session
+ try:
+ session = self.get_ws_session()
+ except Exception as e:
+ logger.error(f"Failed to create session: {e}")
+ return
+
+ # 2. Connect WebSocket
+ ws_url = f"{session['attach_url']}?token={session['attach_token']}"
+ logger.info(f"Connecting to {session['attach_url']}")
+
+ try:
+ async with websockets.connect(ws_url) as ws:
+ self.ws = ws
+ self._connect_time_ms = time.monotonic() * 1000
+
+ # 3. Hello handshake
+ import json
+
+ hello = self._build_hello()
+ await ws.send(json.dumps(hello))
+ logger.info("Hello sent, waiting for response...")
+
+ try:
+ raw = await asyncio.wait_for(ws.recv(), timeout=10.0)
+ if isinstance(raw, str):
+ resp = json.loads(raw)
+ if resp.get("type") == "hello":
+ logger.info(f"Handshake OK: version={resp.get('version')}")
+ self._apply_selected_streams(resp.get("selected_streams"))
+ else:
+ logger.warning(f"Expected hello response, got: {resp.get('type')}")
+ else:
+ logger.warning("Expected text hello response, got binary")
+ except asyncio.TimeoutError:
+ logger.error("Hello handshake timeout")
+ return
+ except Exception as e:
+ logger.error(f"Hello handshake failed: {e}")
+ return
+
+ # 4. Init Opus codec
+ self._init_opus()
+
+ # 5. Launch media tasks
+ self._loop = asyncio.get_running_loop()
+ tasks = [
+ asyncio.create_task(self._run_microphone(self._loop)),
+ asyncio.create_task(self._run_speaker(self._loop)),
+ asyncio.create_task(self._run_receiver()),
+ ]
+ if self.camera_index is not None:
+ tasks.append(asyncio.create_task(self._run_camera()))
+
+ # Start wake word detection if configured
+ if self.wake_word_engine:
+ self.wake_word_engine.start()
+ logger.info("[WakeWord] Detection started")
+
+ # Wait for shutdown signal
+ await self._shutdown_event.wait()
+
+ # Cancel remaining tasks
+ for t in tasks:
+ t.cancel()
+ for t in tasks:
+ try:
+ await t
+ except asyncio.CancelledError:
+ pass
+ except websockets.ConnectionClosed as e:
+ logger.info(f"Connection closed: code={e.code}, reason={e.reason}")
+ except Exception as e:
+ logger.error(f"Connection error: {e}")
+ finally:
+ self._cleanup()
+
+ def _cleanup(self):
+ # Stop wake word detection
+ if self.wake_word_engine:
+ self.wake_word_engine.stop()
+ logger.info("[WakeWord] Detection stopped")
+
+ if self.input_stream is not None:
+ try:
+ self.input_stream.stop_stream()
+ self.input_stream.close()
+ except Exception as e:
+ logger.error(f"Error during input stream cleanup: {e}")
+ self.input_stream = None
+
+ if self.output_stream is not None:
+ try:
+ self.output_stream.stop_stream()
+ self.output_stream.close()
+ except Exception as e:
+ logger.error(f"Error during output stream cleanup: {e}")
+ self.output_stream = None
+
+ if self.video_cap is not None:
+ try:
+ if self.video_cap.isOpened():
+ self.video_cap.release()
+ except Exception as e:
+ logger.error(f"Error during video capture cleanup: {e}")
+ self.video_cap = None
+
+ try:
+ self.pa.terminate()
+ except Exception as e:
+ logger.error(f"Error terminating PyAudio: {e}")
+
+ # Best-effort server-side cleanup
+ try:
+ self.disconnect_ws_session()
+ except Exception as e:
+ logger.error(f"Error disconnecting websocket session: {e}")
+
+ logger.info("Client cleaned up")
+
+ def stop(self):
+ """Signal the client to shut down gracefully."""
+ self._shutdown_event.set()
diff --git a/clients/python-common-sdk/list_audio_devices.py b/clients/python-common-sdk/examples/list_audio_devices.py
similarity index 100%
rename from clients/python-common-sdk/list_audio_devices.py
rename to clients/python-common-sdk/examples/list_audio_devices.py
diff --git a/clients/python-common-sdk/examples/python_livekit_example.py b/clients/python-common-sdk/examples/python_livekit_example.py
new file mode 100644
index 0000000..9a32ecb
--- /dev/null
+++ b/clients/python-common-sdk/examples/python_livekit_example.py
@@ -0,0 +1,74 @@
+import asyncio
+import os
+import signal
+from dotenv import load_dotenv
+import sys
+import os
+sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
+from client.livekit_client import EvaLiveKitClient
+from shared.wake_word import VoskWakeWordEngine
+import threading
+
+load_dotenv()
+
+def on_task_change(old_task: str, new_task: str):
+ """Callback when task changes (e.g., after wake word detection)"""
+ print(f"[App] Task changed: {old_task} -> {new_task}")
+
+async def main():
+ # Please adjust the microphone/speaker/camera index according to your actual situation.
+ client = EvaLiveKitClient(
+ api_key=os.getenv("EVA_API_KEY"),
+ base_url=os.getenv("EVA_BASE_URL", "https://eva.autoarkai.com"),
+ wss_url=os.getenv("EVA_WSS_URL", "wss://rtc.autoarkai.com"),
+ mic_index=None, # MacBook Air Microphone
+ spk_index=None, # MacBook Air Speakers
+ mic_sample_rate=48000,
+ spk_sample_rate=48000,
+ camera_index=None, # Set to an integer to enable video
+ video_width=640,
+ video_height=480,
+ video_fps=20,
+
+ # Wake word detection (optional)
+ wake_word_engine=VoskWakeWordEngine(wake_word="你好方舟"),
+ wake_word_target_task="intent_task",
+
+ # Task change callback (called when task switches, e.g., after wake word)
+ on_task_change=on_task_change,
+ )
+
+ def signal_handler():
+ print("\nShutdown.")
+ client.stop()
+
+ loop = asyncio.get_running_loop()
+ for sig in (signal.SIGINT, signal.SIGTERM):
+ loop.add_signal_handler(sig, signal_handler)
+
+ def command_listener():
+ print("\n[Manual Test] You can now type a task ID (e.g. 'intent_task') and press Enter to test edge hard-switch!")
+ while True:
+ try:
+ cmd = input().strip()
+ if cmd == "quit" or cmd == "exit":
+ client.stop()
+ break
+ if cmd:
+ client.switch_task(cmd)
+ except EOFError:
+ break
+ except Exception as e:
+ print(f"Error in command listener: {e}")
+
+ threading.Thread(target=command_listener, daemon=True).start()
+
+ await client.run()
+
+
+if __name__ == "__main__":
+ try:
+ asyncio.run(main())
+ except KeyboardInterrupt:
+ pass
+
\ No newline at end of file
diff --git a/clients/python-common-sdk/examples/python_ws_example.py b/clients/python-common-sdk/examples/python_ws_example.py
new file mode 100644
index 0000000..ccae369
--- /dev/null
+++ b/clients/python-common-sdk/examples/python_ws_example.py
@@ -0,0 +1,107 @@
+"""
+Eva WebSocket Client Example
+============================
+
+Demonstrates using the WebSocket transport (instead of LiveKit/WebRTC) to
+connect to the Eva platform. Suitable for embedded devices and environments
+where a WebRTC stack is not available.
+
+Usage:
+ 1. Set EVA_API_KEY in your environment or .env file
+ 2. Adjust mic_index / spk_index / camera_index for your hardware
+ 3. Run: python python_ws_example.py
+
+Use `list_audio_devices.py` to discover available audio device indices.
+"""
+
+import asyncio
+import os
+import signal
+import threading
+
+from dotenv import load_dotenv
+import sys
+import os
+sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
+
+from client.ws_client import EvaWebSocketClient
+from shared.wake_word import VoskWakeWordEngine
+
+load_dotenv()
+
+
+def on_task_change(old_task: str, new_task: str):
+ """Callback when task changes (e.g., after wake word detection)"""
+ print(f"[App] Task changed: {old_task} -> {new_task}")
+
+
+async def main():
+ client = EvaWebSocketClient(
+ api_key=os.getenv("EVA_API_KEY"),
+ base_url=os.getenv("EVA_BASE_URL", "https://eva.autoarkai.com"),
+ # Adjust these indices for your hardware (use list_audio_devices.py)
+ mic_index=None,
+ spk_index=None, # MacBook Air Speakers
+ # Use device's native sample rate (MacBook Air: 48kHz)
+ # VOSK and Opus will handle resampling internally
+ mic_sample_rate=48000,
+ spk_sample_rate=48000,
+ channels=1,
+ frame_duration_ms=60,
+ # Set camera_index=None to disable video
+ camera_index=None,
+ video_width=640,
+ video_height=480,
+ video_fps=2,
+ # Echo suppression: auto-mute mic while bot is speaking
+ echo_suppression=True,
+ echo_suppression_timeout_ms=500,
+ # Volume-based barge-in: allow interruption when user speaks loudly
+ # User mic peak must exceed: bot_peak * multiplier + offset
+ # Increase multiplier/offset to make barge-in harder (less sensitive)
+ # Decrease to make it easier (more sensitive)
+ barge_in_multiplier=1.5,
+ barge_in_offset=500,
+ # Wake word detection (optional)
+ # VOSK engine: supports Chinese, auto-downloads small model
+ # wake_word_engine=VoskWakeWordEngine(wake_word="你好方舟")
+ # Set wake_word_engine=None to disable
+ wake_word_engine=VoskWakeWordEngine(wake_word="你好方舟"),
+ wake_word_target_task="intent_task",
+ # Task change callback (called when task switches, e.g., after wake word)
+ on_task_change=on_task_change,
+ )
+
+ def signal_handler():
+ print("\nShutdown requested.")
+ client.stop()
+
+ loop = asyncio.get_running_loop()
+ for sig in (signal.SIGINT, signal.SIGTERM):
+ loop.add_signal_handler(sig, signal_handler)
+
+ def command_listener():
+ print("\n[Manual Test] You can now type a task ID (e.g. 'intent_task') and press Enter to test edge hard-switch!")
+ while True:
+ try:
+ cmd = input().strip()
+ if cmd == "quit" or cmd == "exit":
+ client.stop()
+ break
+ if cmd:
+ client.switch_task(cmd)
+ except EOFError:
+ break
+ except Exception as e:
+ print(f"Error in command listener: {e}")
+
+ threading.Thread(target=command_listener, daemon=True).start()
+
+ await client.run()
+
+
+if __name__ == "__main__":
+ try:
+ asyncio.run(main())
+ except KeyboardInterrupt:
+ pass
diff --git a/clients/python-common-sdk/python_example.py b/clients/python-common-sdk/python_example.py
deleted file mode 100644
index 75b7ba6..0000000
--- a/clients/python-common-sdk/python_example.py
+++ /dev/null
@@ -1,39 +0,0 @@
-import asyncio
-import os
-import signal
-from dotenv import load_dotenv
-from eva_client import EvaClient
-
-load_dotenv()
-
-async def main():
- # Please adjust the microphone/speaker/camera index according to your actual situation.
- client = EvaClient(
- api_key=os.getenv("EVA_API_KEY"),
- mic_index=1,
- spk_index=6,
- mic_sample_rate=16000,
- spk_sample_rate=48000,
- camera_index=11,
- video_width=640,
- video_height=480,
- video_fps=20,
- )
-
- def signal_handler():
- print("\nShutdown.")
- client.stop()
-
- loop = asyncio.get_running_loop()
- for sig in (signal.SIGINT, signal.SIGTERM):
- loop.add_signal_handler(sig, signal_handler)
-
- await client.run()
-
-
-if __name__ == "__main__":
- try:
- asyncio.run(main())
- except KeyboardInterrupt:
- pass
-
\ No newline at end of file
diff --git a/clients/python-common-sdk/requirements.txt b/clients/python-common-sdk/requirements.txt
new file mode 100644
index 0000000..0aac582
--- /dev/null
+++ b/clients/python-common-sdk/requirements.txt
@@ -0,0 +1,21 @@
+# Eva Python Client - Dependencies
+#
+# Required:
+requests>=2.28.0
+websockets>=12.0
+pyaudio>=0.2.13
+numpy>=1.24.0
+opuslib>=3.0.1
+python-dotenv>=1.0.0
+scipy>=1.11.0
+
+# Optional (for LiveKit/WebRTC transport):
+livekit>=0.11.0
+livekit-api>=0.5.0
+
+# Optional (for video capture):
+# opencv-python>=4.8.0
+
+# Optional (for wake word detection):
+# VOSK engine (supports Chinese, auto-downloads model):
+# vosk>=0.3.45
diff --git a/clients/python-common-sdk/shared/__init__.py b/clients/python-common-sdk/shared/__init__.py
new file mode 100644
index 0000000..912256a
--- /dev/null
+++ b/clients/python-common-sdk/shared/__init__.py
@@ -0,0 +1 @@
+# empty init
diff --git a/clients/python-common-sdk/shared/audio.py b/clients/python-common-sdk/shared/audio.py
new file mode 100644
index 0000000..b3d007d
--- /dev/null
+++ b/clients/python-common-sdk/shared/audio.py
@@ -0,0 +1,94 @@
+import queue
+import numpy as np
+import pyaudio
+import logging
+from typing import Optional, Protocol, runtime_checkable
+
+logger = logging.getLogger("EvaSharedAudio")
+
+@runtime_checkable
+class AudioConsumer(Protocol):
+ def add_audio_frame(self, pcm_data: bytes, sample_rate: int, channels: int) -> None:
+ """
+ Feed audio data to the consumer.
+ """
+ ...
+
+class AudioOutputBuffer:
+ def __init__(self, maxsize: int = 200):
+ self.queue: "queue.Queue[np.ndarray]" = queue.Queue(maxsize=maxsize)
+ self.remainder = None
+
+ def put(self, data: bytes):
+ np_data = np.frombuffer(data, dtype=np.int16)
+ try:
+ self.queue.put_nowait(np_data)
+ except queue.Full:
+ try:
+ self.queue.get_nowait()
+ except queue.Empty:
+ pass
+ try:
+ self.queue.put_nowait(np_data)
+ except queue.Full:
+ pass
+
+ def get_chunk(self, frames_needed: int) -> bytes:
+ out_list = []
+ frames_collected = 0
+
+ if self.remainder is not None:
+ n = len(self.remainder)
+ if n > frames_needed:
+ out_list.append(self.remainder[:frames_needed])
+ self.remainder = self.remainder[frames_needed:]
+ return np.concatenate(out_list).tobytes()
+ else:
+ out_list.append(self.remainder)
+ frames_collected += n
+ self.remainder = None
+
+ while frames_collected < frames_needed:
+ try:
+ new_packet = self.queue.get_nowait()
+ packet_len = len(new_packet)
+ needed = frames_needed - frames_collected
+
+ if packet_len > needed:
+ out_list.append(new_packet[:needed])
+ self.remainder = new_packet[needed:]
+ frames_collected += needed
+ else:
+ out_list.append(new_packet)
+ frames_collected += packet_len
+ except queue.Empty:
+ needed = frames_needed - frames_collected
+ silence = np.zeros(needed, dtype=np.int16)
+ out_list.append(silence)
+ frames_collected += needed
+
+ return np.concatenate(out_list).tobytes()
+
+def find_audio_device_index(pa: pyaudio.PyAudio, value, is_input: bool = True) -> Optional[int]:
+ if value is None:
+ return None
+ if isinstance(value, str) and value.strip() == "":
+ return None
+ try:
+ return int(value)
+ except (TypeError, ValueError):
+ pass
+
+ target = str(value).strip()
+ count = pa.get_device_count()
+ for i in range(count):
+ info = pa.get_device_info_by_index(i)
+ name = info.get("name", "")
+ if is_input and info.get("maxInputChannels", 0) == 0:
+ continue
+ if not is_input and info.get("maxOutputChannels", 0) == 0:
+ continue
+ if target in name:
+ return i
+ logger.warning(f"Audio device '{target}' not found, using system default")
+ return None
diff --git a/clients/python-common-sdk/shared/fsm.py b/clients/python-common-sdk/shared/fsm.py
new file mode 100644
index 0000000..cefd66b
--- /dev/null
+++ b/clients/python-common-sdk/shared/fsm.py
@@ -0,0 +1,196 @@
+import threading
+import uuid
+import logging
+from enum import Enum, auto
+from typing import Optional, Callable, List, Dict, Any
+from dataclasses import dataclass
+
+from .protocol import unwrap_rtvi_envelope
+
+logger = logging.getLogger("EvaSharedFSM")
+
+# --- Pure FSM Events ---
+
+@dataclass
+class ConfigUpdateEvent:
+ valid_tasks: List[str]
+
+@dataclass
+class SwitchApprovedEvent:
+ command_id: str
+
+@dataclass
+class SwitchDeniedEvent:
+ reason: str
+
+@dataclass
+class SwitchAdviceEvent:
+ advice_id: str
+ suggested_task: str
+ confidence: float
+
+# --- Pure FSM ---
+
+class TaskFSMState(Enum):
+ IDLE = auto()
+ ACTIVE = auto()
+ PENDING_SWITCH = auto()
+
+class TaskFSM:
+ """
+ Pure Finite State Machine for task state management.
+ Does not know about JSON or RTVI protocols.
+ """
+ def __init__(self, on_task_change: Optional[Callable[[str, str], None]] = None):
+ self.state = TaskFSMState.IDLE
+ self.current_task: Optional[str] = None
+ self.pending_task: Optional[str] = None
+ self.revision = 0
+ self.allowed_tasks: List[str] = []
+ self._on_task_change = on_task_change
+ self._lock = threading.Lock()
+
+ def initialize(self, initial_task: str):
+ with self._lock:
+ self.current_task = initial_task
+ self.state = TaskFSMState.ACTIVE
+ logger.info(f"[TaskFSM] Initialized with task: {initial_task}")
+
+ def on_config_update(self, event: ConfigUpdateEvent):
+ with self._lock:
+ self.allowed_tasks = event.valid_tasks
+ logger.info(f"[TaskFSM] Config updated. Allowed tasks: {self.allowed_tasks}")
+
+ def request_switch(self, target_task: str) -> bool:
+ """Returns True if the switch request is valid and state transitioned to PENDING_SWITCH."""
+ with self._lock:
+ if self.allowed_tasks and target_task not in self.allowed_tasks:
+ logger.warning(f"[TaskFSM] Switch to {target_task} denied: not in allowed tasks")
+ return False
+
+ self.pending_task = target_task
+ self.state = TaskFSMState.PENDING_SWITCH
+ logger.info(f"[TaskFSM] Requested switch to {target_task}, state: PENDING_SWITCH")
+ return True
+
+ def on_switch_approved(self, event: SwitchApprovedEvent) -> bool:
+ """Returns True if the switch was successfully applied."""
+ with self._lock:
+ if self.state != TaskFSMState.PENDING_SWITCH:
+ logger.warning(f"[TaskFSM] Received switch approval but not in PENDING_SWITCH state.")
+ return False
+
+ old_task = self.current_task
+ self.current_task = self.pending_task
+ self.state = TaskFSMState.ACTIVE
+ self.revision += 1
+ self.pending_task = None
+
+ logger.info(f"[TaskFSM] Switch approved: {old_task} -> {self.current_task} (revision {self.revision})")
+
+ if self._on_task_change and old_task != self.current_task:
+ try:
+ self._on_task_change(old_task or "", self.current_task or "")
+ except Exception as e:
+ logger.error(f"[TaskFSM] Error in task change callback: {e}")
+
+ return True
+
+ def on_switch_denied(self, event: SwitchDeniedEvent):
+ with self._lock:
+ self.state = TaskFSMState.ACTIVE if self.current_task else TaskFSMState.IDLE
+ self.pending_task = None
+ logger.info(f"[TaskFSM] Switch denied: {event.reason}")
+
+ def on_switch_advice(self, event: SwitchAdviceEvent, threshold: float) -> bool:
+ """Returns True if the advice triggered an auto-switch."""
+ if event.confidence <= threshold:
+ return False
+
+ with self._lock:
+ old_task = self.current_task
+ self.current_task = event.suggested_task
+ self.state = TaskFSMState.ACTIVE
+ self.revision += 1
+ self.pending_task = None
+
+ logger.info(f"[TaskFSM] Auto-switch triggered: {old_task} -> {self.current_task} (revision {self.revision})")
+
+ if self._on_task_change and old_task != self.current_task:
+ try:
+ self._on_task_change(old_task or "", self.current_task or "")
+ except Exception as e:
+ logger.error(f"[TaskFSM] Error in task change callback: {e}")
+
+ return True
+
+
+# --- Protocol Adapter ---
+
+class RTVITaskNegotiator:
+ """
+ Adapter that translates RTVI JSON messages into FSM Events,
+ and FSM State changes into RTVI JSON Commands/Commits.
+ """
+ def __init__(self, on_task_change: Optional[Callable[[str, str], None]] = None, auto_switch_confidence_threshold: float = 0.8):
+ self.fsm = TaskFSM(on_task_change=on_task_change)
+ self.auto_switch_confidence_threshold = auto_switch_confidence_threshold
+
+ def initialize(self, initial_task: str):
+ self.fsm.initialize(initial_task)
+
+ def handle_message(self, raw_message: dict) -> Optional[dict]:
+ """Parses RTVI message, updates FSM, and optionally returns a payload to send."""
+ message = unwrap_rtvi_envelope(raw_message)
+ msg_type = message.get("type", "unknown")
+
+ if msg_type == "system_config":
+ event = ConfigUpdateEvent(valid_tasks=message.get("valid_tasks", []))
+ self.fsm.on_config_update(event)
+ return None
+
+ if msg_type == "task.switch.result" or "approved" in message:
+ approved = message.get("approved", False)
+ if approved:
+ event = SwitchApprovedEvent(command_id=message.get("id", ""))
+ if self.fsm.on_switch_approved(event):
+ return self._build_commit_payload("edge_command", event.command_id)
+ else:
+ event = SwitchDeniedEvent(reason=message.get("reason", ""))
+ self.fsm.on_switch_denied(event)
+ return None
+
+ if msg_type == "task.switch.advice" or "suggested_task" in message:
+ event = SwitchAdviceEvent(
+ advice_id=message.get("id", ""),
+ suggested_task=message.get("suggested_task", ""),
+ confidence=message.get("confidence", 0.0)
+ )
+ if self.fsm.on_switch_advice(event, self.auto_switch_confidence_threshold):
+ return self._build_commit_payload("advice", event.advice_id)
+ return None
+
+ return None
+
+ def request_switch(self, target_task: str, reason: str = "User Request") -> Optional[dict]:
+ """Translates a switch request intent into an FSM state transition and an RTVI command payload."""
+ if self.fsm.request_switch(target_task):
+ return {
+ "type": "task.switch.command",
+ "id": f"cmd_{uuid.uuid4().hex[:8]}",
+ "from_task": self.fsm.current_task or "",
+ "to_task": target_task,
+ "reason": reason,
+ "revision": self.fsm.revision
+ }
+ return None
+
+ def _build_commit_payload(self, ref_source: str, ref_id: str) -> dict:
+ return {
+ "type": "task.switch.commit",
+ "id": f"commit_{uuid.uuid4().hex[:8]}",
+ "final_task": self.fsm.current_task,
+ "revision": self.fsm.revision,
+ "ref_source": ref_source,
+ "ref_id": ref_id
+ }
diff --git a/clients/python-common-sdk/shared/protocol.py b/clients/python-common-sdk/shared/protocol.py
new file mode 100644
index 0000000..38d98b0
--- /dev/null
+++ b/clients/python-common-sdk/shared/protocol.py
@@ -0,0 +1,20 @@
+import uuid
+
+def wrap_rtvi_envelope(message: dict, topic: str = "task-ir-control") -> dict:
+ return {
+ "label": "rtvi-ai",
+ "type": "client-message",
+ "id": f"msg-{uuid.uuid4().hex[:8]}",
+ "data": {
+ "t": topic,
+ "d": message
+ }
+ }
+
+def unwrap_rtvi_envelope(message: dict) -> dict:
+ if "label" in message and message.get("label") == "rtvi-ai":
+ inner = message.get("data", message)
+ if isinstance(inner, dict) and "d" in inner:
+ return inner["d"]
+ return inner
+ return message
diff --git a/clients/python-common-sdk/shared/wake_word.py b/clients/python-common-sdk/shared/wake_word.py
new file mode 100644
index 0000000..ab0daa2
--- /dev/null
+++ b/clients/python-common-sdk/shared/wake_word.py
@@ -0,0 +1,205 @@
+import threading
+import json
+from pathlib import Path
+from typing import Optional, Callable, List, Dict, Any, Protocol, runtime_checkable
+import numpy as np
+import logging
+
+from .audio import AudioConsumer
+
+logger = logging.getLogger("EvaSharedWakeWord")
+
+class WakeWordEvent:
+ def __init__(self, keyword: str, confidence: float = 1.0, metadata: Optional[Dict[str, Any]] = None):
+ self.keyword = keyword
+ self.confidence = confidence
+ self.metadata = metadata or {}
+
+class WakeWordEngine(Protocol):
+ def start(self) -> None:
+ """
+ Start the wake word detection engine.
+ Can be overridden by engines that need explicit startup (e.g. threads, hardware init).
+ """
+ ...
+
+ def stop(self) -> None:
+ """
+ Stop the wake word detection engine.
+ Can be overridden by engines that need explicit teardown.
+ """
+ ...
+
+ def set_callback(self, callback: Callable[[WakeWordEvent], None]) -> None:
+ """
+ Set the callback to be invoked when a wake word is detected.
+ """
+ ...
+
+class VoskWakeWordEngine(WakeWordEngine, AudioConsumer):
+ def __init__(
+ self,
+ wake_word: str,
+ model_path: Optional[str] = None,
+ sample_rate: int = 16000,
+ ):
+ self.wake_word = wake_word
+ self.model_path = model_path
+ self.sample_rate = sample_rate
+
+ self.on_detected: Optional[Callable[[WakeWordEvent], None]] = None
+
+ self._thread: Optional[threading.Thread] = None
+ self._running = False
+ self._audio_buffer: List[np.ndarray] = []
+ self._buffer_lock = threading.Lock()
+ self._data_event = threading.Event()
+ self._available = False
+
+ self._init_vosk()
+
+ def set_callback(self, callback: Callable[[WakeWordEvent], None]) -> None:
+ self.on_detected = callback
+
+ def _init_vosk(self):
+ try:
+ from vosk import Model
+ self._vosk_model_class = Model
+ self._available = True
+ logger.info(f"[WakeWord] VOSK engine loaded, wake word: '{self.wake_word}'")
+ except ImportError:
+ logger.warning("[WakeWord] vosk not installed. Install with: pip install vosk")
+ self._vosk_model_class = None
+
+ def _ensure_vosk_model(self) -> Optional[str]:
+ if self.model_path:
+ return self.model_path
+
+ model_name = "vosk-model-small-cn-0.22"
+ cache_dir = Path.home() / ".cache" / "eva_ws_client"
+ model_dir = cache_dir / model_name
+
+ if model_dir.exists():
+ return str(model_dir)
+
+ cache_dir.mkdir(parents=True, exist_ok=True)
+ url = f"https://alphacephei.com/vosk/models/{model_name}.zip"
+ zip_path = cache_dir / f"{model_name}.zip"
+
+ logger.info(f"[WakeWord] Downloading VOSK model: {model_name} (~50MB)...")
+ try:
+ import urllib.request
+ urllib.request.urlretrieve(url, str(zip_path))
+
+ logger.info("[WakeWord] Extracting model...")
+ import zipfile
+ with zipfile.ZipFile(str(zip_path), 'r') as zf:
+ zf.extractall(str(cache_dir))
+ zip_path.unlink()
+
+ logger.info(f"[WakeWord] Model ready at: {model_dir}")
+ return str(model_dir)
+ except Exception as e:
+ logger.error(f"[WakeWord] Failed to download VOSK model: {e}")
+ return None
+
+ def add_audio_frame(self, pcm_data: bytes, sample_rate: int, channels: int) -> None:
+ if not self._running:
+ return
+
+ audio_data = np.frombuffer(pcm_data, dtype=np.int16)
+
+ if channels > 1:
+ # Simple downmix by dropping extra channels
+ audio_data = audio_data[::channels]
+
+ if sample_rate != self.sample_rate:
+ if sample_rate % self.sample_rate == 0:
+ # Simple decimation
+ ratio = sample_rate // self.sample_rate
+ audio_data = audio_data[::ratio]
+ else:
+ if not getattr(self, "_warned_sr", False):
+ logger.warning(f"[WakeWord] Audio sample rate {sample_rate} differs from VOSK expected {self.sample_rate} and cannot be easily decimated. Wake word detection might fail.")
+ self._warned_sr = True
+ with self._buffer_lock:
+ self._audio_buffer.append(audio_data)
+ max_samples = self.sample_rate * 2
+ total_samples = sum(len(frame) for frame in self._audio_buffer)
+ while total_samples > max_samples and self._audio_buffer:
+ removed = self._audio_buffer.pop(0)
+ total_samples -= len(removed)
+ self._data_event.set()
+
+ def start(self):
+ if not self._available:
+ logger.warning("[WakeWord] Cannot start: VOSK not available")
+ return
+
+ self._running = True
+ self._thread = threading.Thread(target=self._vosk_detection_loop, daemon=True)
+ self._thread.start()
+ logger.info("[WakeWord] Detection started (engine=vosk)")
+
+ def stop(self):
+ self._running = False
+ self._data_event.set()
+ if self._thread:
+ self._thread.join(timeout=1.0)
+ logger.info("[WakeWord] Detection stopped")
+
+ def _vosk_detection_loop(self):
+ try:
+ model_path = self._ensure_vosk_model()
+ if not model_path:
+ logger.error("[WakeWord] No VOSK model available, stopping detection")
+ self._running = False
+ return
+
+ from vosk import KaldiRecognizer
+ import vosk
+ vosk.SetLogLevel(-1)
+
+ model = self._vosk_model_class(model_path)
+ recognizer = KaldiRecognizer(model, self.sample_rate)
+
+ logger.info(f"[WakeWord] VOSK recognizer ready, listening for '{self.wake_word}'")
+
+ while self._running:
+ self._data_event.wait(timeout=0.1)
+ self._data_event.clear()
+
+ with self._buffer_lock:
+ if not self._audio_buffer:
+ continue
+ audio = np.concatenate(self._audio_buffer)
+ self._audio_buffer.clear()
+
+ audio_bytes = audio.astype(np.int16).tobytes()
+ chunk_size = self.sample_rate * 2 * 2
+ for i in range(0, len(audio_bytes), chunk_size):
+ chunk = audio_bytes[i:i + chunk_size]
+
+ if recognizer.AcceptWaveform(chunk):
+ result = json.loads(recognizer.Result())
+ text = result.get("text", "").replace(" ", "")
+ if text:
+ logger.debug(f"[WakeWord] Recognized: '{text}'")
+
+ if self.wake_word in text:
+ logger.info(f"[WakeWord] Detected! '{self.wake_word}' found in '{text}'")
+ if self.on_detected:
+ self.on_detected(WakeWordEvent(keyword=self.wake_word, metadata={"text": text}))
+ else:
+ partial = json.loads(recognizer.PartialResult())
+ partial_text = partial.get("partial", "").replace(" ", "")
+ if partial_text and self.wake_word in partial_text:
+ logger.info(f"[WakeWord] Detected (partial)! '{self.wake_word}' in '{partial_text}'")
+ if self.on_detected:
+ self.on_detected(WakeWordEvent(keyword=self.wake_word, metadata={"partial": partial_text}))
+ recognizer.Reset()
+
+ except Exception as e:
+ logger.error(f"[WakeWord] VOSK detection loop error: {e}")
+ finally:
+ logger.info("[WakeWord] VOSK detection loop ended")