diff --git a/.gitignore b/.gitignore index b813130..a950d6e 100644 --- a/.gitignore +++ b/.gitignore @@ -14,12 +14,15 @@ node_modules/ .terraform/ plugins/plugins/*/.env plugins/plugins/*/.env.* -plugins/plugins/*/.venv plugins/plugins/*/token.json -external/vibeVoice/.venv # Local server config (API keys); keep *example* files tracked config.toml +# Python virtualenvs and bytecode (external sidecars, plugins) +.venv/ +__pycache__/ +*.pyc + # Go build / test artifacts *.exe *.exe~ diff --git a/README.md b/README.md index 5c0b8a0..b16a84f 100644 --- a/README.md +++ b/README.md @@ -113,7 +113,7 @@ StreamCore starts one layer below prompt-and-tool frameworks: the media path. Yo Details and code: [Bring your own agent](./docs/bring-your-own-agent.md) · [Agent runtime](./docs/agent-runtime.md). -Providers: Deepgram, AssemblyAI, OpenAI, Cartesia, ElevenLabs, MiniMax, Speechify, Telnyx, Ollama, VibeVoice (local), xAI Grok Voice (speech-to-speech), pgvector/Supabase for retrieval. See [Providers](./docs/providers.md). +Providers: Deepgram, AssemblyAI, OpenAI, Cartesia, ElevenLabs, MiniMax, Speechify, Telnyx, Ollama, Moonshine (local), VibeVoice (local), xAI Grok Voice (speech-to-speech), pgvector/Supabase for retrieval. See [Providers](./docs/providers.md). OpenAI STT supports `whisper-1`, `gpt-4o-transcribe`, and `gpt-4o-mini-transcribe` through the independent `openai.stt_model` setting. @@ -128,7 +128,7 @@ Telnyx STT fronts a dozen engines behind one key (`telnyx.transcription_engine`, | [Bring your own agent](./docs/bring-your-own-agent.md) | Five ways to own the intelligence, including the HTTP agent endpoint and the `llm.Client` interface | | [Agent runtime](./docs/agent-runtime.md) | Plugins, skills, RAG, document ingestion | | [Developer agent](./docs/developer-agent.md) | Optional GitHub App and Codex integrations: CI investigation, isolated worktrees, confirmation-gated pull requests | -| [Providers](./docs/providers.md) | Grok speech-to-speech, MiniMax, local VibeVoice, per-provider caveats | +| [Providers](./docs/providers.md) | Grok speech-to-speech, MiniMax, local Moonshine and VibeVoice, per-provider caveats | | [Configuration](./docs/configuration.md) | Full annotated `config.toml` reference | | [Protocol](./docs/protocol.md) | WHIP signaling, DataChannel events, auth | | [Architecture](./docs/architecture.md) | Media flow, why Go, package layout | diff --git a/README.zh-CN.md b/README.zh-CN.md index 208daf9..dad82be 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -113,7 +113,7 @@ StreamCore 位于「提示词 + 工具」类框架的下一层:媒体链路。 详情与代码:[接入你自己的智能体](./docs/bring-your-own-agent.zh-CN.md) · [智能体运行时](./docs/agent-runtime.zh-CN.md)。 -服务商:Deepgram、AssemblyAI、OpenAI、Cartesia、ElevenLabs、MiniMax、Speechify、Telnyx、Ollama、VibeVoice(本地)、xAI Grok Voice(语音到语音),检索支持 pgvector / Supabase。见[服务商](./docs/providers.zh-CN.md)。 +服务商:Deepgram、AssemblyAI、OpenAI、Cartesia、ElevenLabs、MiniMax、Speechify、Telnyx、Ollama、Moonshine(本地)、VibeVoice(本地)、xAI Grok Voice(语音到语音),检索支持 pgvector / Supabase。见[服务商](./docs/providers.zh-CN.md)。 OpenAI STT 可通过独立的 `openai.stt_model` 配置选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 @@ -128,7 +128,7 @@ Telnyx STT 用一个 key 前置十余种引擎(`telnyx.transcription_engine` | [接入你自己的智能体](./docs/bring-your-own-agent.zh-CN.md) | 掌控智能体的五种方式,含 HTTP 智能体端点与 `llm.Client` 接口 | | [智能体运行时](./docs/agent-runtime.zh-CN.md) | 插件、技能、RAG、文档入库 | | [开发者智能体](./docs/developer-agent.zh-CN.md) | 可选的 GitHub App 与 Codex 集成:CI 排查、隔离 worktree、需确认的 Pull Request | -| [服务商](./docs/providers.zh-CN.md) | Grok 语音到语音、MiniMax、本地 VibeVoice 及各服务商注意事项 | +| [服务商](./docs/providers.zh-CN.md) | Grok 语音到语音、MiniMax、本地 Moonshine 与 VibeVoice 及各服务商注意事项 | | [配置](./docs/configuration.zh-CN.md) | 完整带注释的 `config.toml` 参考 | | [协议](./docs/protocol.zh-CN.md) | WHIP 信令、DataChannel 事件、鉴权 | | [架构](./docs/architecture.zh-CN.md) | 媒体流转、为什么用 Go、包结构 | diff --git a/config.toml.example b/config.toml.example index 26765eb..1e5b2ff 100644 --- a/config.toml.example +++ b/config.toml.example @@ -110,13 +110,13 @@ turn_merge_ms = 350 # Debounce window for merging finals into one turn. provider = "" # Supported: grok (empty = classic pipeline) [stt] -provider = "deepgram" # Supported: aliyun, assemblyai, deepgram, openai, telnyx, vibevoice, volcengine +provider = "deepgram" # Supported: aliyun, assemblyai, deepgram, moonshine, openai, telnyx, vibevoice, volcengine [llm] provider = "openai" # Supported: openai, ollama, agent [tts] -provider = "cartesia" # Supported: aliyun, cartesia, deepgram, elevenlabs, mimo, minimax, speechify, telnyx, vibevoice, volcengine +provider = "cartesia" # Supported: aliyun, cartesia, deepgram, elevenlabs, mimo, minimax, moonshine, speechify, telnyx, vibevoice, volcengine # Provider credentials @@ -248,6 +248,18 @@ asr_url = "ws://127.0.0.1:8200" # WebSocket URL for VibeVoice ASR server tts_url = "http://127.0.0.1:8300" # HTTP URL for VibeVoice TTS server voice = "en-Emma_woman" # TTS voice name +# Moonshine — local STT and TTS via external Python services, no API key. +# Start the servers first: +# python external/moonshine/moonshineStt/server.py (default port 8210) +# python external/moonshine/moonshineTts/server.py (default port 8310) +# Language and model selection are sidecar flags, not config keys: both are +# fixed when the model loads rather than chosen per request. +[moonshine] +stt_url = "ws://127.0.0.1:8210" # WebSocket URL for the Moonshine STT server +tts_url = "http://127.0.0.1:8310" # HTTP URL for the Moonshine TTS server +voice = "kokoro_af_heart" # TTS voice id. The prefix picks the vocoder: + # kokoro_, piper_, or zipvoice_ + # Alibaba Cloud Model Studio (DashScope). One endpoint and one key serve both # streaming ASR (stt.provider = "aliyun") and streaming TTS # (tts.provider = "aliyun"). Get an API key at https://bailian.console.aliyun.com/ diff --git a/docs/capabilities.md b/docs/capabilities.md index f0b74d4..e944892 100644 --- a/docs/capabilities.md +++ b/docs/capabilities.md @@ -114,9 +114,9 @@ StreamCore can run a complete speech-to-agent-to-speech pipeline, but that is on | AI integration | Providers | |----------------|-----------| -| Streaming STT | Deepgram, AssemblyAI, OpenAI, VibeVoice (local) | +| Streaming STT | Deepgram, AssemblyAI, OpenAI, Moonshine (local), VibeVoice (local) | | LLM | OpenAI, Ollama (local or self-hosted), or your own HTTP agent endpoint (`agent`) | -| Streaming TTS | Cartesia, Deepgram, ElevenLabs, MiniMax, Speechify, VibeVoice (local) | +| Streaming TTS | Cartesia, Deepgram, ElevenLabs, MiniMax, Moonshine (local), Speechify, VibeVoice (local) | | Speech-to-speech | xAI Grok Voice (replaces STT + LLM + TTS in one model) | | Retrieval | pgvector, Supabase | | Custom tools | Python / TypeScript / JavaScript plugins, native Go tools | diff --git a/docs/configuration.md b/docs/configuration.md index 045c04f..084508d 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -73,13 +73,13 @@ turn_merge_ms = 350 # Debounce window for merging finals into o provider = "" # "grok", or empty for the classic pipeline [stt] -provider = "deepgram" # aliyun | assemblyai | deepgram | openai | telnyx | vibevoice | volcengine +provider = "deepgram" # aliyun | assemblyai | deepgram | moonshine | openai | telnyx | vibevoice | volcengine [llm] provider = "openai" # openai | ollama | agent [tts] -provider = "cartesia" # cartesia | deepgram | elevenlabs | mimo | minimax | speechify | telnyx | vibevoice +provider = "cartesia" # cartesia | deepgram | elevenlabs | mimo | minimax | moonshine | speechify | telnyx | vibevoice # [grok] # Used when realtime.provider = "grok" # api_key = "" @@ -173,6 +173,11 @@ asr_url = "ws://127.0.0.1:8200" tts_url = "http://127.0.0.1:8300" voice = "en-Emma_woman" +[moonshine] +stt_url = "ws://127.0.0.1:8210" +tts_url = "http://127.0.0.1:8310" +voice = "kokoro_af_heart" # Prefix picks the vocoder: kokoro_, piper_, or zipvoice_ + # RAG is optional — omit the [rag] section to disable it entirely. # [rag] # provider = "supabase" # "pgvector" or "supabase" diff --git a/docs/configuration.zh-CN.md b/docs/configuration.zh-CN.md index c855463..102e8a9 100644 --- a/docs/configuration.zh-CN.md +++ b/docs/configuration.zh-CN.md @@ -59,13 +59,13 @@ turn_merge_ms = 350 # Debounce window for merging finals into o provider = "" # "grok", or empty for the classic pipeline [stt] -provider = "deepgram" # aliyun | assemblyai | deepgram | openai | telnyx | vibevoice | volcengine +provider = "deepgram" # aliyun | assemblyai | deepgram | moonshine | openai | telnyx | vibevoice | volcengine [llm] provider = "openai" # openai | ollama | agent [tts] -provider = "cartesia" # cartesia | deepgram | elevenlabs | mimo | minimax | speechify | telnyx | vibevoice +provider = "cartesia" # cartesia | deepgram | elevenlabs | mimo | minimax | moonshine | speechify | telnyx | vibevoice # [grok] # Used when realtime.provider = "grok" # api_key = "" @@ -158,6 +158,11 @@ asr_url = "ws://127.0.0.1:8200" tts_url = "http://127.0.0.1:8300" voice = "en-Emma_woman" +[moonshine] +stt_url = "ws://127.0.0.1:8210" +tts_url = "http://127.0.0.1:8310" +voice = "kokoro_af_heart" # Prefix picks the vocoder: kokoro_, piper_, or zipvoice_ + # RAG is optional — omit the [rag] section to disable it entirely. # [rag] # provider = "supabase" # "pgvector" or "supabase" diff --git a/docs/providers.md b/docs/providers.md index ae995c6..90ff875 100644 --- a/docs/providers.md +++ b/docs/providers.md @@ -4,9 +4,9 @@ | Role | Providers | Required credentials | |------|-----------|----------------------| -| STT | `aliyun`, `assemblyai`, `deepgram`, `openai`, `telnyx`, `vibevoice`, `volcengine` | Matching provider API key, or a local VibeVoice ASR server | +| STT | `aliyun`, `assemblyai`, `deepgram`, `moonshine`, `openai`, `telnyx`, `vibevoice`, `volcengine` | Matching provider API key, or a local Moonshine / VibeVoice ASR server | | LLM | `openai`, `ollama`, `agent` | OpenAI API key, an Ollama instance you control, or your own HTTP agent endpoint | -| TTS | `cartesia`, `deepgram`, `elevenlabs`, `mimo`, `minimax`, `speechify`, `telnyx`, `vibevoice` | Matching provider API key, or a local VibeVoice TTS server | +| TTS | `cartesia`, `deepgram`, `elevenlabs`, `mimo`, `minimax`, `moonshine`, `speechify`, `telnyx`, `vibevoice` | Matching provider API key, or a local Moonshine / VibeVoice TTS server | | Speech-to-speech | `grok` | xAI API key — replaces STT, LLM, and TTS together | | RAG (optional) | `pgvector`, `supabase` | Postgres connection string or Supabase URL + key, plus an OpenAI key for embeddings | @@ -17,6 +17,7 @@ Notes: - `llm.provider = "ollama"` targets any Ollama-compatible endpoint via `base_url` — local or on your own infrastructure. - `llm.provider = "agent"` POSTs each turn to an HTTP endpoint you host; your agent owns memory, prompting, and tools, and replies stream back as SSE, chunked text, or JSON. See [Bring your own agent](./bring-your-own-agent.md). - `stt.provider = "vibevoice"` and `tts.provider = "vibevoice"` use local models; start the Python sidecars first. +- `stt.provider = "moonshine"` and `tts.provider = "moonshine"` are the other fully local pair, also behind Python sidecars. The STT sidecar answers in Deepgram's wire format, so it runs through the same transcript handling as Deepgram itself. See [Local Moonshine setup](#local-moonshine-setup). - `tts.provider = "minimax"` covers 40+ languages and is the strongest option for Mandarin. See [MiniMax TTS](#minimax-tts) for the region and model-plan caveats. - `tts.provider = "telnyx"` is Telnyx hosted synthesis over a per-utterance WebSocket; voice availability varies by account. See [Telnyx TTS](#telnyx-tts) for the connection model and the voice catalog. - `tts.provider = "mimo"` is Xiaomi's MiMo TTS, with Chinese and English voices and optional voice cloning on the paid models. @@ -176,9 +177,9 @@ VibeVoice provides fully local STT and TTS with no API keys, using [VibeVoice-AS ```bash # Apple Silicon (MLX) -pip install mlx-audio numpy websockets fastapi uvicorn +pip install mlx-audio numpy websockets onnxruntime requests fastapi uvicorn # OR PyTorch (Linux / CUDA) -pip install torch transformers librosa numpy websockets fastapi uvicorn +pip install torch "transformers>=5.3.0" accelerate librosa numpy websockets onnxruntime requests fastapi uvicorn python external/vibeVoice/vibeVoiceAsr/server.py # ws://127.0.0.1:8200 python external/vibeVoice/vibeVoiceTTS/server.py # http://127.0.0.1:8300 @@ -198,3 +199,34 @@ voice = "en-Emma_woman" ``` The ASR server accepts live PCM over WebSocket and emits JSON transcript events. The TTS server accepts HTTP POST and returns raw PCM. + +## Local Moonshine setup + +[Moonshine](https://moonshine.ai) is the other fully local pair: streaming STT and TTS from one pip package, no API key and no account. Models are downloaded on first run and cached. English STT models are MIT; other languages load under the non-commercial [Moonshine Community License](https://www.moonshine.ai/license). + +```bash +pip install -r external/moonshine/moonshineStt/requirements.txt +pip install -r external/moonshine/moonshineTts/requirements.txt + +python external/moonshine/moonshineStt/server.py # ws://127.0.0.1:8210 +python external/moonshine/moonshineTts/server.py # http://127.0.0.1:8310 +``` + +```toml +[stt] +provider = "moonshine" + +[tts] +provider = "moonshine" + +[moonshine] +stt_url = "ws://127.0.0.1:8210" +tts_url = "http://127.0.0.1:8310" +voice = "kokoro_af_heart" +``` + +Language and model selection are sidecar flags rather than config keys, because both are fixed when the model loads rather than chosen per request: `--language`, `--model-arch` for STT, `--language` and `--voice` for TTS. English resolves to `medium-streaming` by default; `--model-arch tiny-streaming` trades accuracy for a much smaller footprint. The first request for a TTS voice downloads it, so `--preload` is worth setting for anything but a first try. + +Both sidecars speak Deepgram's wire format. The STT server sends `Results`, `SpeechStarted` and `UtteranceEnd` frames, which the server decodes into Deepgram's own types and routes through the same accumulator — so overlapping finals are merged and an immediate repeat is suppressed exactly as they are for Deepgram. It reports no confidence, and the frames leave the field out rather than inventing one; the pipeline reads the absent value as unknown, not as low. The TTS server mirrors `/v1/speak`, writing raw headerless PCM as it synthesizes, so playback starts on the first clause. + +Measured on an M-series Mac: first TTS chunk at ~130ms with synthesis running about 9x faster than playback. diff --git a/docs/providers.zh-CN.md b/docs/providers.zh-CN.md index 4c4c1f5..13b29f1 100644 --- a/docs/providers.zh-CN.md +++ b/docs/providers.zh-CN.md @@ -4,9 +4,9 @@ | 角色 | 服务商 | 所需凭据 | |------|-----------|----------------------| -| STT | `aliyun`、`assemblyai`、`deepgram`、`openai`、`telnyx`、`vibevoice`、`volcengine` | 对应服务商的 API key,或一个本地 VibeVoice ASR 服务 | +| STT | `aliyun`、`assemblyai`、`deepgram`、`moonshine`、`openai`、`telnyx`、`vibevoice`、`volcengine` | 对应服务商的 API key,或一个本地 Moonshine / VibeVoice ASR 服务 | | LLM | `openai`、`ollama`、`agent` | OpenAI API key、你自己掌控的 Ollama 实例,或你自己的 HTTP 智能体端点 | -| TTS | `cartesia`、`deepgram`、`elevenlabs`、`mimo`、`minimax`、`speechify`、`telnyx`、`vibevoice` | 对应服务商的 API key,或一个本地 VibeVoice TTS 服务 | +| TTS | `cartesia`、`deepgram`、`elevenlabs`、`mimo`、`minimax`、`moonshine`、`speechify`、`telnyx`、`vibevoice` | 对应服务商的 API key,或一个本地 Moonshine / VibeVoice TTS 服务 | | 语音到语音 | `grok` | xAI API key —— 一并取代 STT、LLM 与 TTS | | RAG(可选) | `pgvector`、`supabase` | Postgres 连接串或 Supabase URL + key,另需 OpenAI key 用于 embedding | @@ -17,6 +17,7 @@ - `llm.provider = "ollama"` 通过 `base_url` 指向任何兼容 Ollama 的端点 —— 本地或你自己的基础设施均可。 - `llm.provider = "agent"` 把每一轮对话 POST 到你托管的 HTTP 端点;记忆、提示词与工具都由你的智能体掌控,回复以 SSE、分块文本或 JSON 流式返回。见[接入你自己的智能体](./bring-your-own-agent.zh-CN.md)。 - `stt.provider = "vibevoice"` 与 `tts.provider = "vibevoice"` 使用本地模型;请先启动 Python 边车进程。 +- `stt.provider = "moonshine"` 与 `tts.provider = "moonshine"` 是另一组完全本地的方案,同样基于 Python 边车进程。STT 边车按 Deepgram 的线格式应答,因此走的是与 Deepgram 相同的转写处理路径。见[本地 Moonshine 配置](#本地-moonshine-配置)。 - `tts.provider = "minimax"` 覆盖 40+ 语言,是中文场景下最强的选项。区域与套餐相关的坑见 [MiniMax TTS](#minimax-tts)。 - `tts.provider = "telnyx"` 是 Telnyx 托管合成,每个话语一条 WebSocket 连接;音色可用性因账号而异。连接模型与音色目录见 [Telnyx TTS](#telnyx-tts)。 - `tts.provider = "mimo"` 是小米 MiMo TTS,中英文音色齐备,付费模型还支持声音克隆。 @@ -176,9 +177,9 @@ VibeVoice 提供完全本地、无需 API key 的 STT 与 TTS:识别用 [VibeV ```bash # Apple Silicon (MLX) -pip install mlx-audio numpy websockets fastapi uvicorn +pip install mlx-audio numpy websockets onnxruntime requests fastapi uvicorn # 或 PyTorch(Linux / CUDA) -pip install torch transformers librosa numpy websockets fastapi uvicorn +pip install torch "transformers>=5.3.0" accelerate librosa numpy websockets onnxruntime requests fastapi uvicorn python external/vibeVoice/vibeVoiceAsr/server.py # ws://127.0.0.1:8200 python external/vibeVoice/vibeVoiceTTS/server.py # http://127.0.0.1:8300 @@ -198,3 +199,34 @@ voice = "en-Emma_woman" ``` ASR 服务通过 WebSocket 接收实时 PCM 并输出 JSON 转写事件。TTS 服务接收 HTTP POST 并返回裸 PCM。 + +## 本地 Moonshine 配置 + +[Moonshine](https://moonshine.ai) 是另一组完全本地的方案:流式 STT 与 TTS 来自同一个 pip 包,无需 API key,也无需账号。模型首次运行时下载并缓存。英语 STT 模型为 MIT 许可,其他语言使用非商业的 [Moonshine Community License](https://www.moonshine.ai/license)。 + +```bash +pip install -r external/moonshine/moonshineStt/requirements.txt +pip install -r external/moonshine/moonshineTts/requirements.txt + +python external/moonshine/moonshineStt/server.py # ws://127.0.0.1:8210 +python external/moonshine/moonshineTts/server.py # http://127.0.0.1:8310 +``` + +```toml +[stt] +provider = "moonshine" + +[tts] +provider = "moonshine" + +[moonshine] +stt_url = "ws://127.0.0.1:8210" +tts_url = "http://127.0.0.1:8310" +voice = "kokoro_af_heart" +``` + +语言与模型的选择放在边车进程的命令行参数里,而不是配置项里,因为两者都在模型加载时就固定下来,并非按请求选择:STT 用 `--language`、`--model-arch`,TTS 用 `--language`、`--voice`。英语默认解析到 `medium-streaming`;`--model-arch tiny-streaming` 以准确率换取小得多的占用。某个 TTS 音色的首次请求会触发下载,因此除了初次尝试,建议加上 `--preload`。 + +两个边车都使用 Deepgram 的线格式。STT 服务发送 `Results`、`SpeechStarted` 与 `UtteranceEnd` 帧,服务端将其解析为 Deepgram 自己的类型并交给同一个累加器 —— 重叠的 final 会被合并,紧随其后的重复会被抑制,与 Deepgram 完全一致。它不提供 confidence,帧里也就不带这个字段,而不是编造一个;管线把缺失值理解为「未知」而非「低置信度」。TTS 服务对齐 `/v1/speak`,边合成边写出无头 PCM,因此播放可以从第一个短句开始。 + +在 M 系列 Mac 上实测:首个 TTS 分块约 130 ms,合成速度约为播放速度的 9 倍。 diff --git a/docs/quickstart.md b/docs/quickstart.md index 59d41f3..a5e763d 100644 --- a/docs/quickstart.md +++ b/docs/quickstart.md @@ -126,9 +126,9 @@ ollama pull gpt-oss:20b ```bash # Apple Silicon (MLX) -pip install mlx-audio numpy websockets fastapi uvicorn +pip install mlx-audio numpy websockets onnxruntime requests fastapi uvicorn # OR Linux / CUDA -# pip install torch transformers librosa numpy websockets fastapi uvicorn +# pip install torch "transformers>=5.3.0" accelerate librosa numpy websockets onnxruntime requests fastapi uvicorn python external/vibeVoice/vibeVoiceAsr/server.py # ws://127.0.0.1:8200 python external/vibeVoice/vibeVoiceTTS/server.py # http://127.0.0.1:8300 @@ -162,6 +162,26 @@ go run . Fully local realtime voice, no external API dependencies. Model and sidecar details in [Providers → Local VibeVoice setup](./providers.md#local-vibevoice-setup). +Moonshine is the other local option, and the sidecars swap in without touching the rest of this config: + +```bash +pip install -r external/moonshine/moonshineStt/requirements.txt +pip install -r external/moonshine/moonshineTts/requirements.txt + +python external/moonshine/moonshineStt/server.py # ws://127.0.0.1:8210 +python external/moonshine/moonshineTts/server.py # http://127.0.0.1:8310 +``` + +```toml +[stt] +provider = "moonshine" + +[tts] +provider = "moonshine" +``` + +See [Providers → Local Moonshine setup](./providers.md#local-moonshine-setup). + --- Next: [Configuration reference](./configuration.md) · [Providers](./providers.md) · [Agent runtime](./agent-runtime.md) diff --git a/docs/quickstart.zh-CN.md b/docs/quickstart.zh-CN.md index 7b30281..abfb114 100644 --- a/docs/quickstart.zh-CN.md +++ b/docs/quickstart.zh-CN.md @@ -125,9 +125,9 @@ ollama pull gpt-oss:20b ```bash # Apple Silicon (MLX) -pip install mlx-audio numpy websockets fastapi uvicorn +pip install mlx-audio numpy websockets onnxruntime requests fastapi uvicorn # 或 Linux / CUDA -# pip install torch transformers librosa numpy websockets fastapi uvicorn +# pip install torch "transformers>=5.3.0" accelerate librosa numpy websockets onnxruntime requests fastapi uvicorn python external/vibeVoice/vibeVoiceAsr/server.py # ws://127.0.0.1:8200 python external/vibeVoice/vibeVoiceTTS/server.py # http://127.0.0.1:8300 @@ -161,6 +161,26 @@ go run . 完全本地的实时语音,不依赖任何外部 API。模型与边车进程的细节见[服务商 → 本地 VibeVoice 配置](./providers.zh-CN.md#本地-vibevoice-配置)。 +Moonshine 是另一个本地选项,边车进程可以直接替换,其余配置不动: + +```bash +pip install -r external/moonshine/moonshineStt/requirements.txt +pip install -r external/moonshine/moonshineTts/requirements.txt + +python external/moonshine/moonshineStt/server.py # ws://127.0.0.1:8210 +python external/moonshine/moonshineTts/server.py # http://127.0.0.1:8310 +``` + +```toml +[stt] +provider = "moonshine" + +[tts] +provider = "moonshine" +``` + +见[服务商 → 本地 Moonshine 配置](./providers.zh-CN.md#本地-moonshine-配置)。 + --- 下一步:[配置参考](./configuration.zh-CN.md) · [服务商](./providers.zh-CN.md) · [智能体运行时](./agent-runtime.zh-CN.md) diff --git a/external/moonshine/moonshineStt/README.md b/external/moonshine/moonshineStt/README.md new file mode 100644 index 0000000..4ee25df --- /dev/null +++ b/external/moonshine/moonshineStt/README.md @@ -0,0 +1,65 @@ +# Moonshine STT Server + +**English** | [简体中文](./README.zh-CN.md) + +Live streaming speech-to-text server using [Moonshine](https://moonshine.ai). Accepts raw PCM audio over WebSocket and answers in Deepgram's live-transcription frames, so `stt.provider = "moonshine"` runs through the same accumulate-and-dedupe path the server already uses for Deepgram. + +Everything runs on-device. No API key, no account — the model is downloaded on first run and cached. + +## Models + +The catalog picks a default per language; `--model-arch` overrides it. + +| Arch | Notes | +|------|-------| +| `tiny`, `base` | Non-streaming, lowest footprint | +| `tiny-streaming`, `base-streaming` | Built for live audio | +| `small-streaming`, `medium-streaming` | Higher accuracy, more compute | + +English models are MIT. Other languages load under the non-commercial [Moonshine Community License](https://www.moonshine.ai/license), and the server prints a notice when they do. + +## Install + +```bash +pip install -r requirements.txt +``` + +## Run + +```bash +python server.py +# ws://127.0.0.1:8210 + +python server.py --port 9000 --model-arch base-streaming +python server.py --language es --update-interval 0.8 +``` + +## Protocol + +- **Client → Server**: binary WebSocket frames — raw PCM (16 kHz, 16-bit signed LE, mono) +- **Client → Server**: JSON text frames — `{"type": "KeepAlive" | "Finalize" | "CloseStream"}` +- **Server → Client**: Deepgram JSON frames — `Results`, `SpeechStarted`, `UtteranceEnd`, `Metadata` + +```json +{"type":"SpeechStarted","channel":[0,1],"timestamp":0.42} +{"type":"Results","channel_index":[0,1],"start":0.42,"duration":0.9,"is_final":false,"speech_final":false,"channel":{"alternatives":[{"transcript":"turn on"}]}} +{"type":"Results","channel_index":[0,1],"start":0.42,"duration":1.8,"is_final":true,"speech_final":true,"channel":{"alternatives":[{"transcript":"turn on the lights"}]}} +{"type":"UtteranceEnd","channel":[0,1],"last_word_end":2.22} +``` + +Moonshine completes a transcript line only at the end of an utterance, so every `is_final` it sends is also a `speech_final`. Deepgram splits the two because it freezes words mid-sentence. + +There is no confidence field: Moonshine does not report one, and the server sends no number the model did not produce. The Go client decodes the absent field as 0, which the pipeline reads as "unknown" rather than "low". + +## Options + +| Flag | Default | Description | +|------|---------|-------------| +| `--host` | `127.0.0.1` | Bind host | +| `--port` | `8210` | Bind port | +| `--language` | `en` | Two-letter language code | +| `--model-arch` | catalog default | `tiny`, `base`, `tiny-streaming`, `base-streaming`, `small-streaming`, `medium-streaming` | +| `--update-interval` | `0.5` | Seconds of audio between transcription passes | +| `--log-level` | `INFO` | Logging level | + +`--update-interval` is a floor, not a cadence. A pass has to cover at least as much audio as the last one took to produce, so a machine that cannot keep up lets passes grow instead of falling further behind on every one. diff --git a/external/moonshine/moonshineStt/README.zh-CN.md b/external/moonshine/moonshineStt/README.zh-CN.md new file mode 100644 index 0000000..cbc11c0 --- /dev/null +++ b/external/moonshine/moonshineStt/README.zh-CN.md @@ -0,0 +1,65 @@ +# Moonshine STT Server + +[English](./README.md) | **简体中文** + +使用 [Moonshine](https://moonshine.ai) 的实时流式语音转文字服务。通过 WebSocket 接收原始 PCM 音频,返回 Deepgram 的实时转写帧,因此 `stt.provider = "moonshine"` 走的是服务端已经为 Deepgram 准备好的那条合并与去重路径。 + +全部在本机运行。无需 API key,无需账号 —— 模型首次运行时下载并缓存。 + +## 模型 + +目录会按语言选择默认模型,`--model-arch` 可覆盖。 + +| 架构 | 说明 | +|------|-------| +| `tiny`、`base` | 非流式,占用最小 | +| `tiny-streaming`、`base-streaming` | 为实时音频设计 | +| `small-streaming`、`medium-streaming` | 准确率更高,算力开销更大 | + +英语模型为 MIT 许可。其他语言使用非商业的 [Moonshine Community License](https://www.moonshine.ai/license),加载时服务端会打印提示。 + +## 安装 + +```bash +pip install -r requirements.txt +``` + +## 运行 + +```bash +python server.py +# ws://127.0.0.1:8210 + +python server.py --port 9000 --model-arch base-streaming +python server.py --language es --update-interval 0.8 +``` + +## 协议 + +- **客户端 → 服务端**:二进制 WebSocket 帧 —— 原始 PCM(16 kHz、16 位有符号小端、单声道) +- **客户端 → 服务端**:JSON 文本帧 —— `{"type": "KeepAlive" | "Finalize" | "CloseStream"}` +- **服务端 → 客户端**:Deepgram JSON 帧 —— `Results`、`SpeechStarted`、`UtteranceEnd`、`Metadata` + +```json +{"type":"SpeechStarted","channel":[0,1],"timestamp":0.42} +{"type":"Results","channel_index":[0,1],"start":0.42,"duration":0.9,"is_final":false,"speech_final":false,"channel":{"alternatives":[{"transcript":"turn on"}]}} +{"type":"Results","channel_index":[0,1],"start":0.42,"duration":1.8,"is_final":true,"speech_final":true,"channel":{"alternatives":[{"transcript":"turn on the lights"}]}} +{"type":"UtteranceEnd","channel":[0,1],"last_word_end":2.22} +``` + +Moonshine 只在一句话结束时才完成一行转写,所以它发出的每个 `is_final` 同时也是 `speech_final`。Deepgram 之所以区分两者,是因为它会在句子中途就冻结部分词。 + +没有 confidence 字段:Moonshine 不给出该值,服务端也不会编造模型没产出的数字。Go 客户端把缺失的字段解析为 0,管线将其理解为「未知」而非「低置信度」。 + +## 选项 + +| 参数 | 默认值 | 说明 | +|------|---------|-------------| +| `--host` | `127.0.0.1` | 绑定地址 | +| `--port` | `8210` | 绑定端口 | +| `--language` | `en` | 两字母语言代码 | +| `--model-arch` | 目录默认值 | `tiny`、`base`、`tiny-streaming`、`base-streaming`、`small-streaming`、`medium-streaming` | +| `--update-interval` | `0.5` | 两次转写之间的音频秒数 | +| `--log-level` | `INFO` | 日志级别 | + +`--update-interval` 是下限而不是固定节奏。每一轮至少要覆盖与上一轮耗时相当的音频量,因此跟不上的机器会让每轮覆盖更多音频,而不是每一轮都落后得更远。 diff --git a/external/moonshine/moonshineStt/requirements.txt b/external/moonshine/moonshineStt/requirements.txt new file mode 100644 index 0000000..aca0331 --- /dev/null +++ b/external/moonshine/moonshineStt/requirements.txt @@ -0,0 +1,7 @@ +# Moonshine STT sidecar. +# +# moonshine-voice ships the native runtime and pulls numpy in with it; the +# model itself is downloaded on first run and cached. +moonshine-voice>=0.1.5 +websockets>=12.0 +numpy>=1.24.0 diff --git a/external/moonshine/moonshineStt/server.py b/external/moonshine/moonshineStt/server.py new file mode 100644 index 0000000..26d7789 --- /dev/null +++ b/external/moonshine/moonshineStt/server.py @@ -0,0 +1,328 @@ +#!/usr/bin/env python3 +""" +Moonshine STT - live WebSocket speech-to-text server. + +Accepts raw PCM audio (16kHz, 16-bit, mono) over WebSocket and answers with +Deepgram's live-transcription frames, so the Go server can route Moonshine +through the same accumulate-and-dedupe path it already runs for Deepgram. + +Everything is local: moonshine-voice downloads the model on first run and +needs no API key. + +Protocol: + Client -> Server: binary frames (raw PCM, linear16, 16kHz, mono) + text frames {"type": "KeepAlive" | "Finalize" | "CloseStream"} + Server -> Client: text frames Results / SpeechStarted / UtteranceEnd / Metadata + +Moonshine reports no per-transcript confidence, so the Results frames leave +the field out. Deepgram's Go types decode a missing confidence as 0, which +the pipeline reads as "unknown" rather than "low". +""" + +import argparse +import asyncio +import json +import logging +import time +import uuid +from concurrent.futures import ThreadPoolExecutor + +import numpy as np +import websockets + +from moonshine_voice import ( + Transcriber, + TranscriptEventListener, + get_model_for_language, + model_arch_to_string, + string_to_model_arch, +) + +logger = logging.getLogger("moonshine-stt") + +SAMPLE_RATE = 16000 +SAMPLE_WIDTH = 2 # 16-bit + +# Deepgram reports single-channel audio as channel 0 of 1. +CHANNEL_INDEX = [0, 1] + +# Audio is handed to the model in 100ms batches. Inbound frames are 20ms of +# RTP, and a native call per frame pays the fixed per-pass cost five times +# over for the same audio. +FLUSH_SAMPLES = SAMPLE_RATE // 10 + +# Seconds of audio between transcription passes; --update-interval overrides it. +UPDATE_INTERVAL = 0.5 + +# One shared transcriber holds the model weights, and one worker thread runs +# every native call against it. Serialising on the executor rather than a lock +# is what keeps concurrent sessions from reentering the C API — each session +# still gets its own stream, so their transcripts stay independent. +_transcriber = None +_model_arch = None +_executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="moonshine-stt") + + +def load_model(language: str, model_arch): + """Download (first run) and load the transcription model for a language.""" + global _transcriber, _model_arch + + logger.info(f"Loading Moonshine model for language {language!r}...") + model_path, resolved_arch = get_model_for_language(language, model_arch) + _transcriber = Transcriber(model_path=model_path, model_arch=resolved_arch) + _model_arch = resolved_arch + logger.info(f"Model loaded: {model_arch_to_string(resolved_arch)} at {model_path}") + + +# --------------------------------------------------------------------------- +# Deepgram wire frames +# --------------------------------------------------------------------------- + + +def results_frame(line, text: str, is_final: bool) -> dict: + """Build a Deepgram Results frame for one transcript line.""" + return { + "type": "Results", + "channel_index": CHANNEL_INDEX, + "start": round(float(line.start_time), 4), + "duration": round(float(line.duration), 4), + "is_final": is_final, + # Deepgram separates the two because it freezes words mid-sentence and + # only later decides the utterance ended. Moonshine completes a line + # exactly once, at the end of the utterance, so every final it reports + # is a speech_final too. + "speech_final": is_final, + "channel": {"alternatives": [{"transcript": text}]}, + } + + +def speech_started_frame(timestamp: float) -> dict: + return { + "type": "SpeechStarted", + "channel": CHANNEL_INDEX, + "timestamp": round(float(timestamp), 4), + } + + +def utterance_end_frame(last_word_end: float) -> dict: + return { + "type": "UtteranceEnd", + "channel": CHANNEL_INDEX, + "last_word_end": round(float(last_word_end), 4), + } + + +def error_frame(err: Exception) -> dict: + return { + "type": "Error", + "description": str(err), + "message": err.__class__.__name__, + } + + +class DeepgramFrameListener(TranscriptEventListener): + """Turns Moonshine transcript events into Deepgram wire frames. + + The callbacks fire on the worker thread, inside add_audio. Frames are + queued here and sent once that call returns, which keeps websocket sends + on the event loop and keeps them in the order the model produced them. + """ + + def __init__(self): + self.pending = [] + + def on_line_started(self, event): + self.pending.append(speech_started_frame(event.line.start_time)) + + def on_line_text_changed(self, event): + # A line that just completed fires this immediately before + # on_line_completed carrying the same words. The final below covers + # them, so an interim here would only send the sentence twice. + if event.line.is_complete: + return + text = event.line.text.strip() + if text: + self.pending.append(results_frame(event.line, text, is_final=False)) + + def on_line_completed(self, event): + line = event.line + text = line.text.strip() + if text: + self.pending.append(results_frame(line, text, is_final=True)) + # Deepgram sends UtteranceEnd alongside speech_final, and the Go + # callback uses it to flush anything a missing speech_final left + # buffered. Sending both keeps that fallback path exercised. + self.pending.append(utterance_end_frame(line.start_time + line.duration)) + + def on_error(self, event): + logger.error(f"Transcriber error: {event.error}") + self.pending.append(error_frame(event.error)) + + def drain(self): + frames, self.pending = self.pending, [] + return frames + + +# --------------------------------------------------------------------------- +# Session +# --------------------------------------------------------------------------- + + +class STTSession: + """One WebSocket client: its own transcription stream and audio buffer.""" + + def __init__(self, ws, loop, update_interval: float): + self.ws = ws + self.loop = loop + self.update_interval = update_interval + self.listener = DeepgramFrameListener() + self.stream = None + self.request_id = str(uuid.uuid4()) + self.started_at = time.monotonic() + self._pending = np.empty(0, dtype=np.int16) + + async def _run(self, fn, *args): + return await self.loop.run_in_executor(_executor, fn, *args) + + async def open(self): + self.stream = await self._run(self._open_stream) + + def _open_stream(self): + stream = _transcriber.create_stream(update_interval=self.update_interval) + stream.add_listener(self.listener) + stream.start() + return stream + + async def handle_audio(self, data: bytes): + samples = np.frombuffer(data, dtype=np.int16) + self._pending = np.concatenate([self._pending, samples]) + if len(self._pending) >= FLUSH_SAMPLES: + await self._feed_pending() + + async def _feed_pending(self): + if len(self._pending) == 0: + return + chunk, self._pending = self._pending, np.empty(0, dtype=np.int16) + audio = (chunk.astype(np.float32) / 32768.0).tolist() + await self._run(self.stream.add_audio, audio, SAMPLE_RATE) + await self._send_frames() + + async def finalize(self): + """Transcribe everything buffered so far, the way Deepgram's Finalize does.""" + if len(self._pending): + await self._feed_pending() + return + await self._run(self.stream.update_transcription) + await self._send_frames() + + async def _send_frames(self): + for frame in self.listener.drain(): + if frame["type"] == "Results": + text = frame["channel"]["alternatives"][0]["transcript"] + kind = "FINAL" if frame["is_final"] else "PARTIAL" + logger.info(f"{kind}: {text}") + await self.ws.send(json.dumps(frame)) + + async def close(self): + if self.stream is None: + return + try: + await self._feed_pending() + # stop() transcribes what is left in the stream, so the tail of an + # utterance still reaches the caller as a final rather than being + # dropped with the connection. + await self._run(self.stream.stop) + await self._send_frames() + await self.ws.send(json.dumps(self._metadata_frame())) + except websockets.exceptions.ConnectionClosed: + pass + finally: + await self._run(self.stream.close) + self.stream = None + + def _metadata_frame(self) -> dict: + return { + "type": "Metadata", + "request_id": self.request_id, + "channels": 1, + "duration": round(time.monotonic() - self.started_at, 4), + "models": [model_arch_to_string(_model_arch)], + } + + +async def handle_connection(ws): + """Handle one WebSocket client session.""" + logger.info("New STT session connected") + session = STTSession(ws, asyncio.get_running_loop(), UPDATE_INTERVAL) + await session.open() + + try: + async for message in ws: + if isinstance(message, bytes): + await session.handle_audio(message) + continue + + kind = _control_type(message) + if kind == "Finalize": + await session.finalize() + elif kind == "CloseStream": + break + elif kind != "KeepAlive": + logger.debug(f"Ignoring text message: {message[:120]}") + except websockets.exceptions.ConnectionClosed: + logger.info("STT session disconnected") + finally: + await session.close() + + +def _control_type(message: str) -> str: + try: + return json.loads(message).get("type", "") + except (json.JSONDecodeError, AttributeError): + return "" + + +async def main(host: str, port: int): + logger.info(f"Starting Moonshine STT WebSocket server on ws://{host}:{port}") + async with websockets.serve(handle_connection, host, port, max_size=2**20): + await asyncio.Future() # run forever + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description="Moonshine STT WebSocket Server") + parser.add_argument("--host", default="127.0.0.1", help="Bind host") + parser.add_argument("--port", type=int, default=8210, help="Bind port") + parser.add_argument( + "--language", + default="en", + help="Two-letter language code (default: en). Anything but 'en' loads a " + "model under the non-commercial Moonshine Community License", + ) + parser.add_argument( + "--model-arch", + default=None, + help="tiny, base, tiny-streaming, base-streaming, small-streaming, or " + "medium-streaming. Omit to use the catalog default for the language", + ) + parser.add_argument( + "--update-interval", + type=float, + default=0.5, + help="Seconds of audio between transcription passes. This is a floor: " + "a machine that cannot keep up lets passes grow rather than falling " + "further behind on every one", + ) + parser.add_argument("--log-level", default="INFO") + args = parser.parse_args() + + logging.basicConfig( + level=getattr(logging, args.log_level.upper()), + format="%(asctime)s [%(name)s] %(levelname)s: %(message)s", + ) + + UPDATE_INTERVAL = args.update_interval + + arch = string_to_model_arch(args.model_arch) if args.model_arch else None + load_model(args.language, arch) + + asyncio.run(main(args.host, args.port)) diff --git a/external/moonshine/moonshineTts/README.md b/external/moonshine/moonshineTts/README.md new file mode 100644 index 0000000..090f60e --- /dev/null +++ b/external/moonshine/moonshineTts/README.md @@ -0,0 +1,82 @@ +# Moonshine TTS Server + +**English** | [简体中文](./README.zh-CN.md) + +Text-to-speech server using [Moonshine](https://moonshine.ai). `/v1/speak` mirrors Deepgram's endpoint of the same name — text in the body, audio format in the query, raw headerless PCM written back as it is synthesized — so playback starts on the first clause instead of waiting for the whole utterance. + +Everything runs on-device. No API key, no account — the voice is downloaded on first use and cached. + +## Voices + +Voice ids carry a vocoder prefix: `kokoro_`, `piper_`, or `zipvoice_`. `GET /v1/voices` lists what is on disk and what the catalog can fetch. + +```bash +curl 'http://127.0.0.1:8310/v1/voices?language=en_us' +``` + +## Install + +```bash +pip install -r requirements.txt +``` + +## Run + +```bash +python server.py +# http://127.0.0.1:8310 + +python server.py --port 9000 --voice kokoro_am_adam +python server.py --language es --voice kokoro_ef_dora --preload +``` + +The first request for a voice downloads it, which can take a while. `--preload` moves that cost to startup. + +## API + +### `POST /v1/speak` + +| Query param | Default | Description | +|-------------|---------|-------------| +| `voice` | `--voice` | Catalog voice id | +| `language` | `--language` | Synthesis language, e.g. `en_us` | +| `encoding` | `linear16` | Only `linear16` is supported | +| `sample_rate` | `16000` | Output rate; the model's own rate is resampled to it | +| `container` | `none` | Only `none` is supported | + +The body is the text, either `text/plain` or Deepgram's `{"text": "..."}` JSON. The response is raw PCM (16-bit signed LE, mono) at `sample_rate`, streamed one chunk per piece of the reply. + +```bash +curl -s -X POST 'http://127.0.0.1:8310/v1/speak?voice=kokoro_af_heart' \ + -H 'Content-Type: text/plain' \ + --data 'Hello from Moonshine.' --output speech.pcm +``` + +### `GET /v1/voices` + +```json +{"language": "en_us", "present": ["kokoro_af_heart"], "downloadable": ["kokoro_af_bella", "..."]} +``` + +### `GET /health` + +```json +{"status": "ok"} +``` + +## Options + +| Flag | Default | Description | +|------|---------|-------------| +| `--host` | `127.0.0.1` | Bind host | +| `--port` | `8310` | Bind port | +| `--language` | `en_us` | Default synthesis language | +| `--voice` | `kokoro_af_heart` | Default voice id | +| `--preload` | off | Load the default voice at startup rather than on first request | +| `--log-level` | `INFO` | Logging level | + +## Notes + +A synthesizer speaks one thing at a time — the native layer rejects a second generation while one is in flight — so requests are serialized. Measured on an M-series Mac with `kokoro_af_heart`: first chunk at ~130 ms, synthesis running about 9x faster than playback, which is well clear of what a live conversation needs. + +Chunks are resampled with a running fractional position rather than one chunk at a time, so the joins between them do not click. diff --git a/external/moonshine/moonshineTts/README.zh-CN.md b/external/moonshine/moonshineTts/README.zh-CN.md new file mode 100644 index 0000000..082d088 --- /dev/null +++ b/external/moonshine/moonshineTts/README.zh-CN.md @@ -0,0 +1,82 @@ +# Moonshine TTS Server + +[English](./README.md) | **简体中文** + +使用 [Moonshine](https://moonshine.ai) 的文字转语音服务。`/v1/speak` 对齐 Deepgram 的同名接口 —— 文本放在 body、音频格式放在 query、边合成边把无头 PCM 写回 —— 因此播放可以从第一个短句就开始,而不必等整句合成完。 + +全部在本机运行。无需 API key,无需账号 —— 音色首次使用时下载并缓存。 + +## 音色 + +音色 id 带有声码器前缀:`kokoro_`、`piper_` 或 `zipvoice_`。`GET /v1/voices` 会列出本地已有的和目录中可下载的音色。 + +```bash +curl 'http://127.0.0.1:8310/v1/voices?language=en_us' +``` + +## 安装 + +```bash +pip install -r requirements.txt +``` + +## 运行 + +```bash +python server.py +# http://127.0.0.1:8310 + +python server.py --port 9000 --voice kokoro_am_adam +python server.py --language es --voice kokoro_ef_dora --preload +``` + +某个音色的首次请求会触发下载,可能较慢。`--preload` 把这个开销挪到启动时。 + +## 接口 + +### `POST /v1/speak` + +| Query 参数 | 默认值 | 说明 | +|-------------|---------|-------------| +| `voice` | `--voice` | 目录中的音色 id | +| `language` | `--language` | 合成语言,例如 `en_us` | +| `encoding` | `linear16` | 仅支持 `linear16` | +| `sample_rate` | `16000` | 输出采样率;模型自身采样率会重采样到该值 | +| `container` | `none` | 仅支持 `none` | + +body 即文本,可以是 `text/plain`,也可以是 Deepgram 的 `{"text": "..."}` JSON。响应是 `sample_rate` 下的原始 PCM(16 位有符号小端、单声道),按回复的每个片段分块流式返回。 + +```bash +curl -s -X POST 'http://127.0.0.1:8310/v1/speak?voice=kokoro_af_heart' \ + -H 'Content-Type: text/plain' \ + --data 'Hello from Moonshine.' --output speech.pcm +``` + +### `GET /v1/voices` + +```json +{"language": "en_us", "present": ["kokoro_af_heart"], "downloadable": ["kokoro_af_bella", "..."]} +``` + +### `GET /health` + +```json +{"status": "ok"} +``` + +## 选项 + +| 参数 | 默认值 | 说明 | +|------|---------|-------------| +| `--host` | `127.0.0.1` | 绑定地址 | +| `--port` | `8310` | 绑定端口 | +| `--language` | `en_us` | 默认合成语言 | +| `--voice` | `kokoro_af_heart` | 默认音色 id | +| `--preload` | 关闭 | 在启动时而非首次请求时加载默认音色 | +| `--log-level` | `INFO` | 日志级别 | + +## 说明 + +一个合成器同一时间只能说一件事 —— 原生层会拒绝正在生成时发起的第二次生成 —— 所以请求是串行处理的。在 M 系列 Mac 上用 `kokoro_af_heart` 实测:首个分块约 130 ms,合成速度约为播放速度的 9 倍,对实时对话来说余量充足。 + +分块重采样时会沿用连续的小数读取位置,而不是逐块独立重采样,因此块与块的衔接处不会出现爆音。 diff --git a/external/moonshine/moonshineTts/requirements.txt b/external/moonshine/moonshineTts/requirements.txt new file mode 100644 index 0000000..02df578 --- /dev/null +++ b/external/moonshine/moonshineTts/requirements.txt @@ -0,0 +1,8 @@ +# Moonshine TTS sidecar. +# +# moonshine-voice ships the native runtime and pulls numpy in with it; voices +# are downloaded on first use and cached. +moonshine-voice>=0.1.5 +fastapi>=0.110.0 +uvicorn>=0.27.0 +numpy>=1.24.0 diff --git a/external/moonshine/moonshineTts/server.py b/external/moonshine/moonshineTts/server.py new file mode 100644 index 0000000..9ece9a7 --- /dev/null +++ b/external/moonshine/moonshineTts/server.py @@ -0,0 +1,228 @@ +#!/usr/bin/env python3 +""" +Moonshine TTS - HTTP text-to-speech server. + +POST /v1/speak -> raw PCM audio (16kHz, 16-bit signed LE, mono), streamed +GET /v1/voices -> {"language": "...", "present": [...], "downloadable": [...]} +GET /health -> {"status": "ok"} + +/v1/speak mirrors Deepgram's endpoint of the same name: the text goes in the +body, encoding and sample_rate are query parameters, and the response is +headerless PCM written as it is synthesized. That makes the Go client a near +copy of the Deepgram one, and playback starts on the first clause instead of +waiting for the whole utterance. + +Everything is local: moonshine-voice downloads the voice on first use and +needs no API key. +""" + +import argparse +import json +import logging +import threading + +import numpy as np +from fastapi import FastAPI, HTTPException, Query, Request +from fastapi.responses import StreamingResponse +import uvicorn + +from moonshine_voice import ( + MoonshineError, + TextToSpeech, + list_tts_voices, +) + +logger = logging.getLogger("moonshine-tts") + +TARGET_SAMPLE_RATE = 16000 + +# Defaults, overridden by --language / --voice. +DEFAULT_LANGUAGE = "en_us" +DEFAULT_VOICE = "kokoro_af_heart" + +# One engine per (language, voice). A synthesizer speaks one thing at a time — +# the native layer rejects a second generation while one is in flight — so the +# lock is held for the whole response, not just while starting it. +_engines = {} +_engine_lock = threading.Lock() +_synthesis_lock = threading.Lock() + + +def get_engine(language: str, voice: str) -> TextToSpeech: + """Return the loaded engine for a voice, loading (and downloading) it once.""" + key = (language, voice) + with _engine_lock: + engine = _engines.get(key) + if engine is not None: + return engine + + logger.info(f"Loading voice {voice!r} for language {language!r}...") + engine = TextToSpeech().language(language).voice(voice) + engine.load() + _engines[key] = engine + logger.info(f"Voice {voice!r} ready") + return engine + + +class LinearResampler: + """Resamples a stream of float chunks with linear interpolation. + + Resampling each chunk on its own leaves a discontinuity at every join, + which is audible as a click once a sentence is cut into a dozen of them. + Carrying the last source sample and the fractional read position across + calls makes the joins land where they would in one continuous pass. + """ + + def __init__(self, src_rate: int, dst_rate: int): + self.ratio = src_rate / dst_rate + self.next_t = 0.0 # next output position, in source samples + self.base = 0.0 # global index of tail[0] + self.tail = np.empty(0, dtype=np.float32) + + def process(self, samples: np.ndarray) -> np.ndarray: + if self.ratio == 1.0: + return samples + + buf = np.concatenate([self.tail, samples]) + highest = self.base + len(buf) - 1 # last position we can interpolate + + out = np.empty(0, dtype=np.float32) + if self.next_t <= highest: + count = int(np.floor((highest - self.next_t) / self.ratio)) + 1 + positions = self.next_t + self.ratio * np.arange(count) + out = np.interp(positions - self.base, np.arange(len(buf)), buf) + self.next_t += self.ratio * count + + self.tail = buf[-1:] + self.base = self.base + len(buf) - 1 + return out + + +def to_pcm16(samples: np.ndarray) -> bytes: + """Float samples in -1..1 to little-endian signed 16-bit PCM.""" + if len(samples) == 0: + return b"" + return np.clip(samples * 32767.0, -32768, 32767).astype(np.int16).tobytes() + + +def synthesize_stream(engine: TextToSpeech, text: str, sample_rate: int): + """Yield PCM as the model produces it, one chunk per piece of the reply.""" + with _synthesis_lock: + resampler = None + try: + for chunk in engine.stream(text): + samples = np.asarray(chunk.samples, dtype=np.float32) + if resampler is None: + resampler = LinearResampler(chunk.sample_rate, sample_rate) + pcm = to_pcm16(resampler.process(samples)) + if pcm: + yield pcm + finally: + # A caller that hangs up mid-sentence (barge-in does exactly that) + # closes this generator early, leaving a generation in flight that + # would reject the next request. + if engine.is_streaming: + engine.cancel_stream() + + +app = FastAPI(title="Moonshine TTS") + + +@app.post("/v1/speak") +async def speak( + request: Request, + voice: str = Query(default=None), + language: str = Query(default=None), + encoding: str = Query(default="linear16"), + sample_rate: int = Query(default=TARGET_SAMPLE_RATE), + container: str = Query(default="none"), +): + if encoding != "linear16": + raise HTTPException(400, f"unsupported encoding {encoding!r}; only linear16") + if container != "none": + raise HTTPException(400, f"unsupported container {container!r}; only none") + if sample_rate <= 0: + raise HTTPException(400, "sample_rate must be positive") + + text = await read_text(request) + if not text.strip(): + raise HTTPException(400, "text must not be empty") + + language = language or DEFAULT_LANGUAGE + voice = voice or DEFAULT_VOICE + try: + engine = get_engine(language, voice) + except MoonshineError as e: + raise HTTPException(400, str(e)) + + logger.info(f"Synthesizing {len(text)} chars with {voice!r}") + return StreamingResponse( + synthesize_stream(engine, text, sample_rate), + media_type="audio/pcm", + ) + + +async def read_text(request: Request) -> str: + """Accept either a text/plain body or Deepgram's {"text": "..."} JSON.""" + body = await request.body() + if request.headers.get("content-type", "").startswith("application/json"): + try: + return json.loads(body).get("text", "") + except (json.JSONDecodeError, AttributeError): + raise HTTPException(400, "body is not valid JSON") + return body.decode("utf-8", errors="replace") + + +@app.get("/v1/voices") +async def voices(language: str = Query(default=None)): + language = language or DEFAULT_LANGUAGE + try: + by_availability = list_tts_voices(language) + except MoonshineError as e: + raise HTTPException(400, str(e)) + return { + "language": language, + "present": by_availability.get("present", []), + "downloadable": by_availability.get("downloadable", []), + } + + +@app.get("/health") +async def health(): + return {"status": "ok"} + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description="Moonshine TTS HTTP Server") + parser.add_argument("--host", default="127.0.0.1", help="Bind host") + parser.add_argument("--port", type=int, default=8310, help="Bind port") + parser.add_argument( + "--language", default=DEFAULT_LANGUAGE, help="Default synthesis language" + ) + parser.add_argument( + "--voice", + default=DEFAULT_VOICE, + help="Default voice id. The prefix picks the vocoder: kokoro_, piper_, " + "or zipvoice_", + ) + parser.add_argument( + "--preload", + action="store_true", + help="Download and load the default voice at startup rather than on the " + "first request", + ) + parser.add_argument("--log-level", default="INFO") + args = parser.parse_args() + + logging.basicConfig( + level=getattr(logging, args.log_level.upper()), + format="%(asctime)s [%(name)s] %(levelname)s: %(message)s", + ) + + DEFAULT_LANGUAGE = args.language + DEFAULT_VOICE = args.voice + + if args.preload: + get_engine(DEFAULT_LANGUAGE, DEFAULT_VOICE) + + uvicorn.run(app, host=args.host, port=args.port, log_level=args.log_level.lower()) diff --git a/external/vibeVoice/vibeVoiceAsr/README.md b/external/vibeVoice/vibeVoiceAsr/README.md index 4acb89c..026ec21 100644 --- a/external/vibeVoice/vibeVoiceAsr/README.md +++ b/external/vibeVoice/vibeVoiceAsr/README.md @@ -9,7 +9,9 @@ Live streaming speech-to-text server using Microsoft VibeVoice-ASR. Accepts raw | Platform | Model | Backend | |----------|-------|---------| | Apple Silicon | `mlx-community/VibeVoice-ASR-4bit` | mlx-audio | -| Linux / CUDA | `microsoft/VibeVoice-ASR` | PyTorch + transformers | +| Linux / CUDA | `microsoft/VibeVoice-ASR-HF` | PyTorch + transformers ≥ 5.3 | + +The PyTorch backend loads the model with transformers' native `VibeVoiceAsrForConditionalGeneration`, so it needs the `-HF` checkpoint. The original `microsoft/VibeVoice-ASR` checkpoint only loads through Microsoft's `vibevoice` package and won't work here. ## Install @@ -19,7 +21,7 @@ pip install -r requirements.txt # Then install one backend: pip install mlx-audio # Apple Silicon # OR -pip install torch transformers librosa # PyTorch +pip install torch "transformers>=5.3.0" accelerate librosa # PyTorch ``` ## Run @@ -29,7 +31,7 @@ python server.py # ws://127.0.0.1:8200 python server.py --port 9000 --model mlx-community/VibeVoice-ASR-bf16 -python server.py --silence-timeout 1.0 --energy-threshold 400 +python server.py --silence-timeout 1.0 --vad-threshold 0.6 ``` ## Protocol @@ -42,7 +44,7 @@ python server.py --silence-timeout 1.0 --energy-threshold 400 {"text": "hello how are you doing", "is_final": true} ``` -The server buffers incoming audio, detects speech boundaries via energy-based VAD, and transcribes when silence is detected (~800 ms default). Partial results are emitted every ~3 seconds during long utterances. +The server buffers incoming audio, detects speech boundaries with Silero VAD, and transcribes when silence is detected (~800 ms default). Only final results are emitted. ## Options @@ -52,5 +54,5 @@ The server buffers incoming audio, detects speech boundaries via energy-based VA | `--port` | `8200` | Bind port | | `--model` | auto (MLX 4-bit or PyTorch) | HuggingFace model name | | `--silence-timeout` | `0.8` | Seconds of silence before final result | -| `--energy-threshold` | `500` | RMS energy threshold for speech detection | +| `--vad-threshold` | `0.5` | Silero VAD speech probability threshold (0.0-1.0) | | `--log-level` | `INFO` | Logging level | diff --git a/external/vibeVoice/vibeVoiceAsr/README.zh-CN.md b/external/vibeVoice/vibeVoiceAsr/README.zh-CN.md index 8d3b6f4..7a07c3d 100644 --- a/external/vibeVoice/vibeVoiceAsr/README.zh-CN.md +++ b/external/vibeVoice/vibeVoiceAsr/README.zh-CN.md @@ -9,7 +9,9 @@ | 平台 | 模型 | 后端 | |----------|-------|---------| | Apple Silicon | `mlx-community/VibeVoice-ASR-4bit` | mlx-audio | -| Linux / CUDA | `microsoft/VibeVoice-ASR` | PyTorch + transformers | +| Linux / CUDA | `microsoft/VibeVoice-ASR-HF` | PyTorch + transformers ≥ 5.3 | + +PyTorch 后端通过 transformers 原生的 `VibeVoiceAsrForConditionalGeneration` 加载模型,因此需要 `-HF` 版本的权重。原始的 `microsoft/VibeVoice-ASR` 只能通过微软的 `vibevoice` 包加载,在这里无法使用。 ## 安装 @@ -19,7 +21,7 @@ pip install -r requirements.txt # Then install one backend: pip install mlx-audio # Apple Silicon # OR -pip install torch transformers librosa # PyTorch +pip install torch "transformers>=5.3.0" accelerate librosa # PyTorch ``` ## 运行 @@ -29,7 +31,7 @@ python server.py # ws://127.0.0.1:8200 python server.py --port 9000 --model mlx-community/VibeVoice-ASR-bf16 -python server.py --silence-timeout 1.0 --energy-threshold 400 +python server.py --silence-timeout 1.0 --vad-threshold 0.6 ``` ## 协议 @@ -42,7 +44,7 @@ python server.py --silence-timeout 1.0 --energy-threshold 400 {"text": "hello how are you doing", "is_final": true} ``` -服务端会缓冲进入的音频,用基于能量的 VAD 检测语音边界,并在检测到静音时(默认约 800 ms)进行转写。长句期间每约 3 秒发出一次中间结果。 +服务端会缓冲进入的音频,用 Silero VAD 检测语音边界,并在检测到静音时(默认约 800 ms)进行转写。只发出最终结果。 ## 选项 @@ -52,5 +54,5 @@ python server.py --silence-timeout 1.0 --energy-threshold 400 | `--port` | `8200` | 绑定端口 | | `--model` | 自动(MLX 4-bit 或 PyTorch) | HuggingFace 模型名 | | `--silence-timeout` | `0.8` | 出最终结果前的静音秒数 | -| `--energy-threshold` | `500` | 语音检测的 RMS 能量阈值 | +| `--vad-threshold` | `0.5` | Silero VAD 语音概率阈值(0.0-1.0) | | `--log-level` | `INFO` | 日志级别 | diff --git a/external/vibeVoice/vibeVoiceAsr/requirements.txt b/external/vibeVoice/vibeVoiceAsr/requirements.txt index 8444cb1..7437387 100644 --- a/external/vibeVoice/vibeVoiceAsr/requirements.txt +++ b/external/vibeVoice/vibeVoiceAsr/requirements.txt @@ -11,5 +11,5 @@ requests # For downloading VAD model # Default model: mlx-community/VibeVoice-ASR-4bit # PyTorch (Linux / Windows / CUDA) -# pip install torch transformers librosa -# Default model: microsoft/VibeVoice-ASR +# pip install torch "transformers>=5.3.0" accelerate librosa +# Default model: microsoft/VibeVoice-ASR-HF diff --git a/external/vibeVoice/vibeVoiceAsr/server.py b/external/vibeVoice/vibeVoiceAsr/server.py index f2e182c..f265f45 100644 --- a/external/vibeVoice/vibeVoiceAsr/server.py +++ b/external/vibeVoice/vibeVoiceAsr/server.py @@ -144,6 +144,7 @@ def load_model(model_name): try: from mlx_audio.stt.utils import load + model_name = model_name or "mlx-community/VibeVoice-ASR-4bit" logger.info("Using MLX backend") logger.info(f"Loading model: {model_name}") _model = load(model_name) @@ -153,34 +154,35 @@ def load_model(model_name): except ImportError: logger.warning("mlx-audio not installed, falling back to PyTorch") - # PyTorch fallback + # PyTorch fallback. The original microsoft/VibeVoice-ASR checkpoint ships no + # transformers remote code, so AutoModelForCausalLM can't load it; the -HF + # checkpoint works with the native class added in transformers 5.3. try: import torch - from transformers import AutoModelForCausalLM, AutoProcessor + from transformers import AutoProcessor, VibeVoiceAsrForConditionalGeneration + model_name = model_name or "microsoft/VibeVoice-ASR-HF" logger.info("Using PyTorch backend") logger.info(f"Loading model: {model_name}") + model = VibeVoiceAsrForConditionalGeneration.from_pretrained( + model_name, + device_map="auto" if torch.cuda.is_available() else None, + ) + model.eval() _model = { - "processor": AutoProcessor.from_pretrained( - model_name, trust_remote_code=True - ), - "model": AutoModelForCausalLM.from_pretrained( - model_name, - trust_remote_code=True, - torch_dtype=( - torch.float16 if torch.cuda.is_available() else torch.float32 - ), - device_map="auto" if torch.cuda.is_available() else None, - ), + "processor": AutoProcessor.from_pretrained(model_name), + "model": model, } _backend = "pytorch" - logger.info("Model loaded successfully") + logger.info(f"Model loaded on {model.device} ({model.dtype})") except Exception as e: raise RuntimeError( f"Failed to load model: {e}\n" "Install mlx-audio (Apple Silicon): pip install mlx-audio\n" - "Install PyTorch: pip install torch transformers" + 'Install PyTorch: pip install torch "transformers>=5.3.0" accelerate librosa\n' + "The PyTorch backend needs a transformers-format checkpoint such as " + "microsoft/VibeVoice-ASR-HF" ) @@ -213,23 +215,22 @@ def transcribe_audio(wav_path): import torch import librosa - audio, sr = librosa.load(wav_path, sr=16000) processor = _model["processor"] model = _model["model"] - inputs = processor( - audios=audio, - sampling_rate=sr, - return_tensors="pt", - trust_remote_code=True, + # The acoustic tokenizer runs at 24 kHz and the processor passes numpy + # arrays through without resampling, so upsample the 16 kHz capture here. + audio, _ = librosa.load(wav_path, sr=processor.feature_extractor.sampling_rate) + + inputs = processor.apply_transcription_request(audio=audio).to( + model.device, model.dtype ) - if torch.cuda.is_available(): - inputs = {k: v.cuda() for k, v in inputs.items()} with torch.no_grad(): output_ids = model.generate(**inputs, max_new_tokens=8192) - raw = processor.batch_decode(output_ids, skip_special_tokens=True)[0] + generated_ids = output_ids[:, inputs["input_ids"].shape[1] :] + raw = processor.decode(generated_ids, return_format="transcription_only")[0] text = _extract_text(raw) return "" if _is_noise_only(text) else text @@ -396,7 +397,7 @@ async def main(host: str, port: int, model_name: str): parser = argparse.ArgumentParser(description="VibeVoice ASR WebSocket Server") parser.add_argument("--host", default="127.0.0.1", help="Bind host") parser.add_argument("--port", type=int, default=8200, help="Bind port") - parser.add_argument("--model", default=None, help="Model name or path") + parser.add_argument("--model", default=None, help="Model name or path (default depends on backend)") parser.add_argument( "--silence-timeout", type=float, @@ -417,11 +418,4 @@ async def main(host: str, port: int, model_name: str): format="%(asctime)s [%(name)s] %(levelname)s: %(message)s", ) - if args.model is None: - args.model = ( - "mlx-community/VibeVoice-ASR-4bit" - if is_apple_silicon() - else "microsoft/VibeVoice-ASR" - ) - asyncio.run(main(args.host, args.port, args.model)) diff --git a/internal/config/config.go b/internal/config/config.go index 0949eb6..be0396b 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -37,6 +37,7 @@ type Config struct { Ollama OllamaConfig `toml:"ollama"` Agent AgentConfig `toml:"agent"` VibeVoice VibeVoiceConfig `toml:"vibevoice"` + Moonshine MoonshineConfig `toml:"moonshine"` Volcengine VolcengineConfig `toml:"volcengine"` Aliyun AliyunConfig `toml:"aliyun"` Cartesia CartesiaConfig `toml:"cartesia"` @@ -475,6 +476,15 @@ type VibeVoiceConfig struct { Voice string `toml:"voice"` // TTS voice name } +// MoonshineConfig points at the two local Moonshine sidecars. Language and +// model selection live on the sidecars' own flags, since both are fixed when +// the model loads rather than per request. +type MoonshineConfig struct { + STTURL string `toml:"stt_url"` // WebSocket URL for the STT server + TTSURL string `toml:"tts_url"` // HTTP URL for the TTS server + Voice string `toml:"voice"` // TTS voice id, e.g. kokoro_af_heart +} + // AliyunConfig configures Alibaba Cloud Model Studio (DashScope), which serves // both streaming ASR and streaming TTS from the same endpoint and key. type AliyunConfig struct { @@ -588,6 +598,8 @@ func Load(path string) (*Config, error) { setDefault(&cfg.Ollama.SystemPrompt, "You are a helpful AI voice assistant having a natural phone conversation. Keep responses to 1-2 sentences unless asked for detail. When interrupted (indicated by bracketed context), respond the way a patient human would: if they say 'no', address the disagreement; if they redirect, follow their lead. Never repeat what you already said, never ask 'would you like me to continue', and never mention that you were interrupted.") setDefault(&cfg.VibeVoice.ASRURL, "ws://127.0.0.1:8200") setDefault(&cfg.VibeVoice.TTSURL, "http://127.0.0.1:8300") + setDefault(&cfg.Moonshine.STTURL, "ws://127.0.0.1:8210") + setDefault(&cfg.Moonshine.TTSURL, "http://127.0.0.1:8310") setDefault(&cfg.Display.Plugin, "display-projector") if cfg.Display.Projector.TimeoutMs == 0 { @@ -600,6 +612,7 @@ func Load(path string) (*Config, error) { return nil, fmt.Errorf("[display.projector] timeout_ms and fast_path_max_chars must be positive") } setDefault(&cfg.VibeVoice.Voice, "en-Emma_woman") + setDefault(&cfg.Moonshine.Voice, "kokoro_af_heart") setDefault(&cfg.Codex.Binary, "codex") setDefault(&cfg.Codex.ModelProvider, "openai") diff --git a/internal/stt/moonshine.go b/internal/stt/moonshine.go new file mode 100644 index 0000000..aa75a59 --- /dev/null +++ b/internal/stt/moonshine.go @@ -0,0 +1,152 @@ +package stt + +import ( + "context" + "encoding/json" + "fmt" + "log" + "sync" + "time" + + "github.com/gorilla/websocket" + + msginterfaces "github.com/deepgram/deepgram-go-sdk/v3/pkg/api/listen/v1/websocket/interfaces" +) + +// moonshineClient connects to the Moonshine STT WebSocket server +// (external/moonshine/moonshineStt) and implements the STT Client interface. +// +// The sidecar answers in Deepgram's live-transcription frames, so they are +// decoded into the SDK's own types and handed to deepgramCallback instead of +// a second accumulator. Merging overlapping finals, suppressing an immediate +// repeat, and flushing on UtteranceEnd are decisions this pipeline already +// made once for Deepgram; a local model does not need them made differently. +type moonshineClient struct { + conn *websocket.Conn + cancel context.CancelFunc + cb *deepgramCallback + mu sync.Mutex + closed bool +} + +// moonshineErrorFrame is the sidecar's error event. Deepgram's own error type +// carries more fields than a local process has anything to put in. +type moonshineErrorFrame struct { + Description string `json:"description"` + Message string `json:"message"` +} + +func NewMoonshineClient(ctx context.Context, wsURL string, onResult func(TranscriptResult)) (*moonshineClient, error) { + sttCtx, cancel := context.WithCancel(ctx) + + dialer := websocket.Dialer{ + HandshakeTimeout: 10 * time.Second, + } + + conn, _, err := dialer.DialContext(sttCtx, wsURL, nil) + if err != nil { + cancel() + return nil, fmt.Errorf("moonshine stt connect %s: %w", wsURL, err) + } + + c := &moonshineClient{ + conn: conn, + cancel: cancel, + cb: &deepgramCallback{onResult: onResult}, + } + + go c.readLoop() + + log.Printf("[stt] connected to Moonshine STT at %s", wsURL) + return c, nil +} + +func (c *moonshineClient) readLoop() { + for { + _, message, err := c.conn.ReadMessage() + if err != nil { + c.mu.Lock() + closed := c.closed + c.mu.Unlock() + if !closed { + log.Printf("[stt] moonshine read error: %v", err) + } + return + } + + if err := routeMoonshineSTTFrame(message, c.cb); err != nil { + log.Printf("[stt] moonshine parse error: %v", err) + } + } +} + +// routeMoonshineSTTFrame decodes one sidecar frame and dispatches it to cb. +// An unrecognised type is dropped rather than reported: the sidecar sends +// Metadata the pipeline has no use for, and the Deepgram SDK likewise ignores +// events with no callback behind them. +func routeMoonshineSTTFrame(data []byte, cb *deepgramCallback) error { + var header msginterfaces.MessageType + if err := json.Unmarshal(data, &header); err != nil { + return fmt.Errorf("frame is not JSON: %w", err) + } + + switch header.Type { + case "Results": + var msg msginterfaces.MessageResponse + if err := json.Unmarshal(data, &msg); err != nil { + return fmt.Errorf("parse Results: %w", err) + } + return cb.Message(&msg) + + case "SpeechStarted": + var msg msginterfaces.SpeechStartedResponse + if err := json.Unmarshal(data, &msg); err != nil { + return fmt.Errorf("parse SpeechStarted: %w", err) + } + return cb.SpeechStarted(&msg) + + case "UtteranceEnd": + var msg msginterfaces.UtteranceEndResponse + if err := json.Unmarshal(data, &msg); err != nil { + return fmt.Errorf("parse UtteranceEnd: %w", err) + } + return cb.UtteranceEnd(&msg) + + case "Error": + var msg moonshineErrorFrame + if err := json.Unmarshal(data, &msg); err != nil { + return fmt.Errorf("parse Error: %w", err) + } + log.Printf("[stt] moonshine sidecar error: %s (%s)", msg.Description, msg.Message) + } + + return nil +} + +// SendAudio sends raw linear16 PCM bytes to the Moonshine STT server. +func (c *moonshineClient) SendAudio(data []byte) error { + c.mu.Lock() + defer c.mu.Unlock() + if c.closed { + return fmt.Errorf("moonshine: connection closed") + } + return c.conn.WriteMessage(websocket.BinaryMessage, data) +} + +// Close shuts down the WebSocket connection. +func (c *moonshineClient) Close() { + c.mu.Lock() + defer c.mu.Unlock() + if c.closed { + return + } + c.closed = true + + c.conn.WriteMessage( + websocket.CloseMessage, + websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""), + ) + c.conn.Close() + c.cancel() + log.Println("[stt] moonshine closed") +} diff --git a/internal/stt/moonshine_test.go b/internal/stt/moonshine_test.go new file mode 100644 index 0000000..eb07538 --- /dev/null +++ b/internal/stt/moonshine_test.go @@ -0,0 +1,229 @@ +package stt + +import ( + "context" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/gorilla/websocket" +) + +// collectMoonshineFrames routes frames through a callback wired to a channel, +// the way the read loop does, so the routing can be exercised without a socket. +func collectMoonshineFrames(t *testing.T, frames ...string) []TranscriptResult { + t.Helper() + + var got []TranscriptResult + cb := &deepgramCallback{onResult: func(r TranscriptResult) { + got = append(got, r) + }} + + for _, frame := range frames { + if err := routeMoonshineSTTFrame([]byte(frame), cb); err != nil { + t.Fatalf("frame %q rejected: %v", frame, err) + } + } + return got +} + +// The sidecar completes a line only at the end of an utterance, so its final +// frame carries both is_final and speech_final. That is the combination the +// Deepgram callback treats as a finished turn. +func TestMoonshineFinalFrameEmitsOneFinal(t *testing.T) { + frame := `{"type":"Results","channel_index":[0,1],"start":0.5,"duration":1.8,` + + `"is_final":true,"speech_final":true,` + + `"channel":{"alternatives":[{"transcript":"what is the weather today"}]}}` + + got := collectMoonshineFrames(t, frame) + + if len(got) != 1 { + t.Fatalf("emitted %d results, want 1: %+v", len(got), got) + } + if !got[0].IsFinal { + t.Error("result reported partial, want final") + } + if got[0].Text != "what is the weather today" { + t.Errorf("text = %q, want the full line", got[0].Text) + } +} + +// Moonshine reports no confidence, so the field is absent from every frame. +// Zero is what the pipeline reads as "unknown"; anything else would be a +// number the model never produced. +func TestMoonshineMissingConfidenceIsUnknown(t *testing.T) { + frame := `{"type":"Results","is_final":true,"speech_final":true,` + + `"channel":{"alternatives":[{"transcript":"hello"}]}}` + + got := collectMoonshineFrames(t, frame) + + if len(got) != 1 { + t.Fatalf("emitted %d results, want 1", len(got)) + } + if got[0].Confidence != 0 { + t.Errorf("confidence = %v, want 0 (absent maps to unknown)", got[0].Confidence) + } +} + +// Partials are what barge-in classifies a short burst with; losing them +// degrades interruption to VAD alone. +func TestMoonshineInterimFrameEmitsPartial(t *testing.T) { + frame := `{"type":"Results","is_final":false,"speech_final":false,` + + `"channel":{"alternatives":[{"transcript":"what is the"}]}}` + + got := collectMoonshineFrames(t, frame) + + if len(got) != 1 { + t.Fatalf("emitted %d results, want 1", len(got)) + } + if got[0].IsFinal { + t.Error("interim frame reported final, want partial") + } + if got[0].Text != "what is the" { + t.Errorf("text = %q, want %q", got[0].Text, "what is the") + } +} + +// Routing the frames through the Deepgram callback is the whole point of the +// wire format: a final that never carries speech_final still has to reach the +// caller, flushed by UtteranceEnd rather than sitting in the accumulator. +func TestMoonshineUtteranceEndFlushesBufferedFinal(t *testing.T) { + got := collectMoonshineFrames(t, + `{"type":"Results","is_final":true,"speech_final":false,`+ + `"channel":{"alternatives":[{"transcript":"book a table"}]}}`, + `{"type":"UtteranceEnd","channel":[0,1],"last_word_end":2.3}`, + ) + + var finals []TranscriptResult + for _, r := range got { + if r.IsFinal { + finals = append(finals, r) + } + } + if len(finals) != 1 { + t.Fatalf("emitted %d finals, want 1: %+v", len(finals), got) + } + if finals[0].Text != "book a table" { + t.Errorf("flushed text = %q, want %q", finals[0].Text, "book a table") + } +} + +// Metadata and SpeechStarted carry nothing the pipeline consumes. Dropping +// them quietly is what the Deepgram SDK does; treating them as errors would +// fill the log on every session. +func TestMoonshineNonTranscriptFramesEmitNothing(t *testing.T) { + got := collectMoonshineFrames(t, + `{"type":"SpeechStarted","channel":[0,1],"timestamp":0.5}`, + `{"type":"Metadata","request_id":"abc","channels":1,"duration":4.2}`, + `{"type":"Something-New"}`, + ) + + if len(got) != 0 { + t.Fatalf("emitted %d results, want none: %+v", len(got), got) + } +} + +func TestMoonshineRejectsMalformedJSON(t *testing.T) { + cb := &deepgramCallback{onResult: func(TranscriptResult) {}} + + if err := routeMoonshineSTTFrame([]byte("not json"), cb); err == nil { + t.Fatal("malformed frame accepted, want an error") + } +} + +// fakeMoonshineSTT stands in for the sidecar: it upgrades the connection, +// writes a scripted set of frames, and keeps the socket open so the client's +// read loop is exercised end to end. +func newFakeMoonshineSTT(t *testing.T, frames ...string) string { + t.Helper() + + upgrader := websocket.Upgrader{} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + defer conn.Close() + + // Wait for audio before answering, so the test covers the send path too. + if _, _, err := conn.ReadMessage(); err != nil { + return + } + for _, frame := range frames { + if err := conn.WriteMessage(websocket.TextMessage, []byte(frame)); err != nil { + return + } + } + // Hold the connection until the client tears it down. + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + })) + t.Cleanup(server.Close) + + return "ws" + strings.TrimPrefix(server.URL, "http") +} + +func TestMoonshineClientStreamsPartialThenFinal(t *testing.T) { + results := make(chan TranscriptResult, 8) + + url := newFakeMoonshineSTT(t, + `{"type":"SpeechStarted","channel":[0,1],"timestamp":0.1}`, + `{"type":"Results","is_final":false,"speech_final":false,`+ + `"channel":{"alternatives":[{"transcript":"turn on"}]}}`, + `{"type":"Results","is_final":true,"speech_final":true,`+ + `"channel":{"alternatives":[{"transcript":"turn on the lights"}]}}`, + ) + + client, err := NewMoonshineClient(context.Background(), url, func(r TranscriptResult) { + results <- r + }) + if err != nil { + t.Fatalf("dial failed: %v", err) + } + defer client.Close() + + if err := client.SendAudio(make([]byte, 640)); err != nil { + t.Fatalf("SendAudio failed: %v", err) + } + + partial := waitResult(t, results) + if partial.IsFinal || partial.Text != "turn on" { + t.Errorf("first result = %+v, want the partial %q", partial, "turn on") + } + + final := waitResult(t, results) + if !final.IsFinal || final.Text != "turn on the lights" { + t.Errorf("second result = %+v, want the final %q", final, "turn on the lights") + } +} + +func TestMoonshineClientSendAfterCloseFails(t *testing.T) { + url := newFakeMoonshineSTT(t) + + client, err := NewMoonshineClient(context.Background(), url, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("dial failed: %v", err) + } + + client.Close() + client.Close() // idempotent: the pipeline can tear down twice on hangup + + if err := client.SendAudio(make([]byte, 640)); err == nil { + t.Fatal("SendAudio succeeded after Close, want an error") + } +} + +func TestMoonshineClientDialFailureIsReported(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + // Port 1 on loopback refuses immediately rather than hanging. + if _, err := NewMoonshineClient(ctx, "ws://127.0.0.1:1", func(TranscriptResult) {}); err == nil { + t.Fatal("dial to a dead port succeeded, want an error") + } +} diff --git a/internal/stt/stt.go b/internal/stt/stt.go index b8eb9f7..fd8fc20 100644 --- a/internal/stt/stt.go +++ b/internal/stt/stt.go @@ -76,7 +76,9 @@ func NewClient(ctx context.Context, cfg *config.Config, onResult func(Transcript return NewTelnyxClient(ctx, cfg.Telnyx, onResult) case "vibevoice": return NewVibeVoiceClient(ctx, cfg.VibeVoice.ASRURL, onResult) + case "moonshine": + return NewMoonshineClient(ctx, cfg.Moonshine.STTURL, onResult) default: - return nil, fmt.Errorf("unknown stt provider %q (supported: aliyun, assemblyai, deepgram, openai, telnyx, vibevoice, volcengine)", cfg.STT.Provider) + return nil, fmt.Errorf("unknown stt provider %q (supported: aliyun, assemblyai, deepgram, moonshine, openai, telnyx, vibevoice, volcengine)", cfg.STT.Provider) } } diff --git a/internal/tts/client.go b/internal/tts/client.go index 5bd6cdd..1b97275 100644 --- a/internal/tts/client.go +++ b/internal/tts/client.go @@ -100,8 +100,10 @@ func newProviderClient(cfg *config.Config) (Client, error) { cfg.Volcengine.TTSResourceID, cfg.Volcengine.TTSURL), nil case "vibevoice": return NewVibeVoiceClient(cfg.VibeVoice.TTSURL, cfg.VibeVoice.Voice), nil + case "moonshine": + return NewMoonshineClient(cfg.Moonshine.TTSURL, cfg.Moonshine.Voice), nil default: - return nil, fmt.Errorf("unknown tts provider %q (supported: aliyun, cartesia, deepgram, elevenlabs, mimo, minimax, speechify, telnyx, vibevoice, volcengine)", cfg.TTS.Provider) + return nil, fmt.Errorf("unknown tts provider %q (supported: aliyun, cartesia, deepgram, elevenlabs, mimo, minimax, moonshine, speechify, telnyx, vibevoice, volcengine)", cfg.TTS.Provider) } } diff --git a/internal/tts/moonshine.go b/internal/tts/moonshine.go new file mode 100644 index 0000000..6a4ac08 --- /dev/null +++ b/internal/tts/moonshine.go @@ -0,0 +1,110 @@ +package tts + +import ( + "context" + "fmt" + "io" + "log" + "net/http" + "net/url" + "strings" +) + +// defaultMoonshineVoice is the catalog voice the sidecar loads when none is +// configured. The prefix picks the vocoder: kokoro_, piper_, or zipvoice_. +const defaultMoonshineVoice = "kokoro_af_heart" + +const defaultMoonshineTTSURL = "http://127.0.0.1:8310" + +// moonshineTTSClient connects to the Moonshine TTS HTTP server +// (external/moonshine/moonshineTts) and implements the TTS Client interface. +// +// The sidecar mirrors Deepgram's /v1/speak: text in the body, encoding and +// sample rate as query parameters, raw headerless PCM written back as it is +// synthesized. So this streams the response the same way the Deepgram client +// does, and the first clause plays while the rest is still being generated. +type moonshineTTSClient struct { + baseURL string + voice string + httpClient *http.Client +} + +// NewMoonshineClient creates a client that talks to the Moonshine TTS server. +// baseURL is the HTTP address (e.g. http://127.0.0.1:8310); an empty voice +// falls back to defaultMoonshineVoice. +func NewMoonshineClient(baseURL, voice string) Client { + if baseURL == "" { + baseURL = defaultMoonshineTTSURL + } + if voice == "" { + voice = defaultMoonshineVoice + } + return &moonshineTTSClient{ + baseURL: strings.TrimSuffix(baseURL, "/"), + voice: voice, + httpClient: newPooledHTTPClient(), + } +} + +func (c *moonshineTTSClient) buildRequest(ctx context.Context, text string) (*http.Request, error) { + params := url.Values{} + params.Set("voice", c.voice) + params.Set("encoding", "linear16") + params.Set("sample_rate", "16000") + params.Set("container", "none") + + u := fmt.Sprintf("%s/v1/speak?%s", c.baseURL, params.Encode()) + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, u, strings.NewReader(text)) + if err != nil { + return nil, fmt.Errorf("moonshine tts create request: %w", err) + } + req.Header.Set("Content-Type", "text/plain") + return req, nil +} + +// SynthesizeStream reads the response body as it arrives. The sidecar writes +// one chunk per piece of the reply, so bytes are playable the moment they land. +func (c *moonshineTTSClient) SynthesizeStream(ctx context.Context, text string) (<-chan StreamChunk, error) { + req, err := c.buildRequest(ctx, text) + if err != nil { + return nil, err + } + + resp, err := c.httpClient.Do(req) + if err != nil { + return nil, fmt.Errorf("moonshine tts request: %w", err) + } + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + resp.Body.Close() + return nil, fmt.Errorf("moonshine tts error %d: %s", resp.StatusCode, string(body)) + } + return streamHTTPResponse(ctx, resp), nil +} + +func (c *moonshineTTSClient) Synthesize(ctx context.Context, text string) ([]byte, error) { + req, err := c.buildRequest(ctx, text) + if err != nil { + return nil, err + } + + resp, err := c.httpClient.Do(req) + if err != nil { + return nil, fmt.Errorf("moonshine tts request: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + return nil, fmt.Errorf("moonshine tts error %d: %s", resp.StatusCode, string(body)) + } + + data, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("moonshine tts read response: %w", err) + } + + log.Printf("[tts:moonshine] synthesized %d bytes for %d chars of text", len(data), len(text)) + return data, nil +} diff --git a/internal/tts/moonshine_test.go b/internal/tts/moonshine_test.go new file mode 100644 index 0000000..60af3ce --- /dev/null +++ b/internal/tts/moonshine_test.go @@ -0,0 +1,166 @@ +package tts + +import ( + "context" + "io" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "testing" + "time" +) + +// fakeMoonshineTTS records what the client asked for and replies with body. +type fakeMoonshineTTS struct { + server *httptest.Server + path string + query url.Values + body string + header string +} + +func newFakeMoonshineTTS(t *testing.T, handler http.HandlerFunc) *fakeMoonshineTTS { + t.Helper() + + f := &fakeMoonshineTTS{} + f.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + f.path = r.URL.Path + f.query = r.URL.Query() + f.body = string(body) + f.header = r.Header.Get("Content-Type") + handler(w, r) + })) + t.Cleanup(f.server.Close) + return f +} + +func writePCM(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + w.Write([]byte{0x01, 0x02, 0x03, 0x04}) +} + +// The sidecar mirrors Deepgram's /v1/speak, so the request has to be the +// shape that endpoint takes: text in the body, audio format in the query. +func TestMoonshineRequestMatchesDeepgramSpeak(t *testing.T) { + f := newFakeMoonshineTTS(t, writePCM) + + client := NewMoonshineClient(f.server.URL, "kokoro_af_bella") + if _, err := client.Synthesize(context.Background(), "hello there"); err != nil { + t.Fatalf("Synthesize failed: %v", err) + } + + if f.path != "/v1/speak" { + t.Errorf("path = %q, want /v1/speak", f.path) + } + if f.body != "hello there" { + t.Errorf("body = %q, want the raw text", f.body) + } + if f.header != "text/plain" { + t.Errorf("Content-Type = %q, want text/plain", f.header) + } + for param, want := range map[string]string{ + "voice": "kokoro_af_bella", + "encoding": "linear16", + "sample_rate": "16000", + "container": "none", + } { + if got := f.query.Get(param); got != want { + t.Errorf("%s = %q, want %q", param, got, want) + } + } +} + +func TestMoonshineDefaults(t *testing.T) { + f := newFakeMoonshineTTS(t, writePCM) + + // A trailing slash on the configured URL must not produce //v1/speak. + client := NewMoonshineClient(f.server.URL+"/", "") + if _, err := client.Synthesize(context.Background(), "hi"); err != nil { + t.Fatalf("Synthesize failed: %v", err) + } + + if f.path != "/v1/speak" { + t.Errorf("path = %q, want /v1/speak", f.path) + } + if got := f.query.Get("voice"); got != defaultMoonshineVoice { + t.Errorf("voice = %q, want the default %q", got, defaultMoonshineVoice) + } +} + +// Playback starts on the first chunk, so the client must hand bytes over as +// they arrive rather than after the sidecar has finished the sentence. +func TestMoonshineStreamsChunksAsTheyArrive(t *testing.T) { + released := make(chan struct{}) + + f := newFakeMoonshineTTS(t, func(w http.ResponseWriter, _ *http.Request) { + flusher, ok := w.(http.Flusher) + if !ok { + t.Error("test server cannot flush") + return + } + w.Write([]byte{0x01, 0x02}) + flusher.Flush() + <-released + w.Write([]byte{0x03, 0x04}) + flusher.Flush() + }) + + client := NewMoonshineClient(f.server.URL, "") + ch, err := client.SynthesizeStream(context.Background(), "two chunks") + if err != nil { + t.Fatalf("SynthesizeStream failed: %v", err) + } + + first := waitChunk(t, ch) + if string(first.PCM) != "\x01\x02" { + t.Errorf("first chunk = % x, want 01 02", first.PCM) + } + + close(released) + + second := waitChunk(t, ch) + if string(second.PCM) != "\x03\x04" { + t.Errorf("second chunk = % x, want 03 04", second.PCM) + } +} + +// A sidecar that has not downloaded the voice answers 400 with the reason. +// Swallowing it would leave the caller with silence and no explanation. +func TestMoonshineSurfacesErrorBody(t *testing.T) { + f := newFakeMoonshineTTS(t, func(w http.ResponseWriter, _ *http.Request) { + http.Error(w, `{"detail":"unknown voice 'nope'"}`, http.StatusBadRequest) + }) + + client := NewMoonshineClient(f.server.URL, "nope") + + _, err := client.Synthesize(context.Background(), "hi") + if err == nil { + t.Fatal("Synthesize succeeded on a 400, want an error") + } + if !strings.Contains(err.Error(), "unknown voice") { + t.Errorf("error = %v, want the sidecar's reason in it", err) + } + + if _, err := client.SynthesizeStream(context.Background(), "hi"); err == nil { + t.Fatal("SynthesizeStream succeeded on a 400, want an error") + } +} + +func waitChunk(t *testing.T, ch <-chan StreamChunk) StreamChunk { + t.Helper() + select { + case chunk, ok := <-ch: + if !ok { + t.Fatal("stream closed before a chunk arrived") + } + if chunk.Err != nil { + t.Fatalf("stream error: %v", chunk.Err) + } + return chunk + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for a chunk") + return StreamChunk{} + } +}