From 61aa20cb5276cf2168493c3d8d538792ddc770d2 Mon Sep 17 00:00:00 2001 From: Mikei386 <44135113+Mikei386@users.noreply.github.com> Date: Sat, 5 Sep 2026 16:06:27 +0200 Subject: [PATCH] Add native Qwen3 TTS streaming for Hermes --- ATHENA.md | 4 +- README.md | 8 +- integrations/hermes-qwen3-stream/README.md | 23 +++ integrations/hermes-qwen3-stream/__init__.py | 153 +++++++++++++++++++ integrations/hermes-qwen3-stream/plugin.yaml | 5 + platform/docker/tts-gateway/tts_gateway.py | 118 +++++++++++++- router/ai_profile_router.py | 75 +++++++++ 7 files changed, 382 insertions(+), 4 deletions(-) create mode 100644 integrations/hermes-qwen3-stream/README.md create mode 100644 integrations/hermes-qwen3-stream/__init__.py create mode 100644 integrations/hermes-qwen3-stream/plugin.yaml diff --git a/ATHENA.md b/ATHENA.md index f2f966f..6057178 100644 --- a/ATHENA.md +++ b/ATHENA.md @@ -58,7 +58,9 @@ Qwen-Profil wird vom Profile Controller verwaltet. - Uncensored: separates lokales Profil - FLUX.2-klein-4B: Bildgenerierung und Editing; Qwen wird dafür kurz entladen und danach automatisch wiederhergestellt -- Qwen3-TTS 1.7B: RTX 3060; Piper bleibt CPU-Fallback +- Qwen3-TTS 1.7B: RTX 3060; Piper bleibt CPU-Fallback. Der Router reicht + zusätzlich natives 24-kHz-PCM für den optionalen Hermes-Streaming-Adapter + unter `integrations/hermes-qwen3-stream` durch. Die verbindlichen Werte stehen in `config/profile-matrix.json` und `docs/STANDARD_PROFILE_MATRIX.md`. diff --git a/README.md b/README.md index 6c2a579..f92e93b 100644 --- a/README.md +++ b/README.md @@ -101,7 +101,9 @@ Standard, bis Hermes' Sitzungsfehler behoben ist. - Athena-Dashboard: `http://192.168.1.212:8099` Der Router stellt Sprache OpenAI-kompatibel bereit: Sprachausgabe über -`/v1/audio/speech` und Spracherkennung über `/v1/audio/transcriptions`. Das +`/v1/audio/speech`, natives Qwen-PCM-Streaming über +`/v1/audio/speech/pcm-stream` und Spracherkennung über +`/v1/audio/transcriptions`. Das Whisper-Modell liegt persistent im Docker-Volume `whisper-data`; Audiodaten werden lokal auf Athena verarbeitet. Für OpenClaw Talk liegt der lokale Realtime-Provider unter @@ -110,6 +112,10 @@ verbindet Mikrofon → Athena Whisper → normalen OpenClaw-Agenten → aktives Athena-TTS, sodass Modell, Werkzeuge und Memory auch im Sprachmodus erhalten bleiben. Die Installation landet in OpenClaws persistentem Datenverzeichnis und bleibt deshalb bei normalen Container-Updates bestehen. + +Für Hermes liegt unter `integrations/hermes-qwen3-stream` ein optionales, +persistentes Backend-Plugin. Es nutzt den nativen PCM-Strom und verkürzt den +Beginn der Sprachausgabe, ohne den Modellrouter oder die Textprofile zu ändern. - Portainer: `https://192.168.1.212:9443` - Hermes-Dashboard auf Unraid: `http://192.168.1.2:9119` diff --git a/integrations/hermes-qwen3-stream/README.md b/integrations/hermes-qwen3-stream/README.md new file mode 100644 index 0000000..617ca38 --- /dev/null +++ b/integrations/hermes-qwen3-stream/README.md @@ -0,0 +1,23 @@ +# Hermes Qwen3-TTS PCM streaming adapter + +This optional Hermes backend plugin uses Athena's native +`/v1/audio/speech/pcm-stream` route. It starts playback while Qwen3-TTS is +still synthesizing the current sentence instead of waiting for a complete +audio file. + +Install this directory as `${HERMES_HOME}/plugins/qwen3-stream`, enable the +plugin and set: + +```yaml +tts: + provider: qwen3-stream + streaming: + provider: qwen3-stream +``` + +The adapter reuses `tts.openai.base_url`, `tts.openai.api_key`, model, voice +and language unless an explicit `tts.qwen3-stream` section overrides them. +This avoids copying the Athena credential into another file. + +Rollback is immediate: restore `tts.provider` and `tts.streaming.provider` to +`openai`, disable the plugin and restart the Hermes gateway. diff --git a/integrations/hermes-qwen3-stream/__init__.py b/integrations/hermes-qwen3-stream/__init__.py new file mode 100644 index 0000000..836630a --- /dev/null +++ b/integrations/hermes-qwen3-stream/__init__.py @@ -0,0 +1,153 @@ +from __future__ import annotations + +from pathlib import Path +from typing import Any, Dict, Iterator, List, Optional + +import requests + +from agent.tts_provider import TTSProvider +from tools.tool_backend_helpers import resolve_openai_audio_api_key +from tools.tts_streaming import StreamingTTSProvider, register as register_streamer +from tools.tts_tool import _load_tts_config + + +NAME = "qwen3-stream" +SAMPLE_RATE = 24000 + + +def _settings() -> Dict[str, Any]: + config = _load_tts_config() + own = dict(config.get(NAME) or {}) + fallback = dict(config.get("openai") or {}) + own.setdefault("base_url", fallback.get("base_url", "")) + own.setdefault("api_key", fallback.get("api_key", "")) + own.setdefault("model", fallback.get("model", "tts-1")) + own.setdefault("voice", fallback.get("voice", "alloy")) + own.setdefault("language", fallback.get("language", "German")) + own.setdefault("chunk_size", 4) + return own + + +def _url(path: str, section: Optional[Dict[str, Any]] = None) -> str: + cfg = section or _settings() + base = str(cfg.get("base_url") or "").rstrip("/") + if not base: + raise RuntimeError("tts.qwen3-stream.base_url is not configured") + if not base.endswith("/v1"): + base += "/v1" + return base + path + + +def _headers(section: Optional[Dict[str, Any]] = None) -> Dict[str, str]: + cfg = section or _settings() + key = str(cfg.get("api_key") or resolve_openai_audio_api_key() or "").strip() + headers = {"Accept": "application/octet-stream"} + if key: + headers["Authorization"] = f"Bearer {key}" + return headers + + +def _payload(text: str, section: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: + cfg = section or _settings() + payload: Dict[str, Any] = { + "input": text, + "model": cfg.get("model") or "tts-1", + "voice": cfg.get("voice") or "alloy", + "language": cfg.get("language") or "German", + } + instruct = str(cfg.get("instruct") or "").strip() + if instruct: + payload["instruct"] = instruct + return payload + + +class Qwen3PCMStreamer(StreamingTTSProvider): + sample_rate = SAMPLE_RATE + channels = 1 + sample_width = 2 + + @staticmethod + def available() -> bool: + try: + return bool(_settings().get("base_url")) + except Exception: + return False + + def stream(self, text: str) -> Iterator[bytes]: + cfg = dict(_settings()) + cfg.update(self.section or {}) + payload = _payload(text, cfg) + payload["chunk_size"] = max(1, int(cfg.get("chunk_size", 4))) + with requests.post( + _url("/audio/speech/pcm-stream", cfg), + json=payload, + headers=_headers(cfg), + stream=True, + timeout=(5, 120), + ) as response: + response.raise_for_status() + pending = b"" + # An explicit read size prevents urllib3 from buffering the + # unknown-length response until connection close. + for chunk in response.iter_content(chunk_size=4096): + if not chunk: + continue + data = pending + chunk + even = len(data) & ~1 + if even: + yield data[:even] + pending = data[even:] + + +class Qwen3TTSProvider(TTSProvider): + @property + def name(self) -> str: + return NAME + + @property + def display_name(self) -> str: + return "Athena Qwen3-TTS Streaming" + + def is_available(self) -> bool: + return Qwen3PCMStreamer.available() + + def list_voices(self) -> List[Dict[str, Any]]: + voice = str(_settings().get("voice") or "alloy") + return [{"id": voice, "display": voice, "language": "de"}] + + def synthesize( + self, + text: str, + output_path: str, + *, + voice: Optional[str] = None, + model: Optional[str] = None, + speed: Optional[float] = None, + format: str = "mp3", + **extra: Any, + ) -> str: + cfg = _settings() + payload = _payload(text, cfg) + payload["response_format"] = format + if voice: + payload["voice"] = voice + if model: + payload["model"] = model + if speed is not None: + payload["speed"] = speed + response = requests.post( + _url("/audio/speech", cfg), + json=payload, + headers=_headers(cfg), + timeout=(5, 120), + ) + response.raise_for_status() + Path(output_path).write_bytes(response.content) + return output_path + + +register_streamer(NAME)(Qwen3PCMStreamer) + + +def register(ctx) -> None: + ctx.register_tts_provider(Qwen3TTSProvider()) diff --git a/integrations/hermes-qwen3-stream/plugin.yaml b/integrations/hermes-qwen3-stream/plugin.yaml new file mode 100644 index 0000000..3d7703d --- /dev/null +++ b/integrations/hermes-qwen3-stream/plugin.yaml @@ -0,0 +1,5 @@ +name: qwen3-stream +version: 0.1.1 +description: Native PCM streaming adapter for the local Athena Qwen3-TTS service +author: Mike AI local stack +kind: backend diff --git a/platform/docker/tts-gateway/tts_gateway.py b/platform/docker/tts-gateway/tts_gateway.py index 04315e1..244cb74 100644 --- a/platform/docker/tts-gateway/tts_gateway.py +++ b/platform/docker/tts-gateway/tts_gateway.py @@ -8,6 +8,7 @@ by the profile router. Request text is never logged or persisted. from __future__ import annotations import io +import http.client import json import os import re @@ -17,6 +18,7 @@ import time import unicodedata import urllib.error import urllib.request +import urllib.parse import wave from array import array from http import HTTPStatus @@ -728,6 +730,47 @@ def synthesize_qwen(text: str, output_format: str, return _convert(audio, "pcm", 1.0) if output_format == "pcm" else (audio, content_type) +def open_qwen_pcm_stream(text: str, chunk_size: int = 4) \ + -> tuple[http.client.HTTPConnection, http.client.HTTPResponse]: + """Open Qwen's native token-level PCM stream without buffering it. + + The upstream emits headerless 24 kHz mono signed 16-bit little-endian + PCM. Keeping this response streaming is what lets playback begin while + the remainder of the sentence is still being synthesized. + """ + parsed = urllib.parse.urlparse(QWEN_TTS_URL) + if parsed.scheme != "http" or not parsed.hostname: + raise RuntimeError("QWEN_TTS_URL must be an http URL") + port = parsed.port or 80 + prefix = parsed.path.rstrip("/") + payload = json.dumps({ + "model": QWEN_TTS_MODEL, + "input": prepare_for_qwen_speech(text), + "voice": QWEN_TTS_VOICE, + "language": QWEN_TTS_LANGUAGE, + "chunk_size": chunk_size, + }, separators=(",", ":")).encode() + connection = http.client.HTTPConnection( + parsed.hostname, port, timeout=QWEN_TTS_TIMEOUT) + try: + connection.request( + "POST", + f"{prefix}/v1/audio/speech/pcm-stream", + body=payload, + headers={"Content-Type": "application/json", + "Accept": "application/octet-stream"}, + ) + response = connection.getresponse() + if response.status != HTTPStatus.OK: + message = response.read(512).decode(errors="replace") + raise RuntimeError( + f"Qwen PCM stream failed ({response.status}): {message}") + return connection, response + except Exception: + connection.close() + raise + + def synthesize(text: str, output_format: str, speed: float) -> tuple[bytes, str]: acquired = SYNTHESIS_LOCK.acquire(timeout=QUEUE_TIMEOUT) if acquired: @@ -797,7 +840,7 @@ class Handler(BaseHTTPRequestHandler): ) def do_POST(self) -> None: # noqa: N802 - if self.path != "/tts": + if self.path not in {"/tts", "/tts/pcm-stream"}: self.send_json(HTTPStatus.NOT_FOUND, {"error": "not found"}) return try: @@ -810,7 +853,7 @@ class Handler(BaseHTTPRequestHandler): return try: request = json.loads(self.rfile.read(length)) - text = request.get("text", "") + text = request.get("input", request.get("text", "")) voice = request.get("voice", VOICE_ALIAS) output_format = request.get("format", "mp3") speed = float(request.get("speed", 1.0)) @@ -829,6 +872,9 @@ class Handler(BaseHTTPRequestHandler): if not 0.5 <= speed <= 2.0: self.send_json(HTTPStatus.BAD_REQUEST, {"error": "invalid speed"}) return + if self.path == "/tts/pcm-stream": + self._stream_qwen_pcm(text.strip(), request) + return started = time.monotonic() try: audio, content_type = synthesize(text.strip(), output_format, speed) @@ -842,6 +888,74 @@ class Handler(BaseHTTPRequestHandler): f"{time.monotonic() - started:.2f}s") self.send_bytes(HTTPStatus.OK, audio, content_type) + def _stream_qwen_pcm(self, text: str, request: dict) -> None: + """Unframe Qwen's PCM frames and relay their audio immediately.""" + try: + chunk_size = max(1, min(32, int(request.get("chunk_size", 4)))) + except (TypeError, ValueError): + self.send_json(HTTPStatus.BAD_REQUEST, + {"error": "invalid chunk_size"}) + return + acquired = SYNTHESIS_LOCK.acquire(timeout=QUEUE_TIMEOUT) + if not acquired: + self.send_json(HTTPStatus.SERVICE_UNAVAILABLE, + {"error": "speech queue timeout"}) + return + connection = None + started = time.monotonic() + headers_sent = False + try: + connection, response = open_qwen_pcm_stream(text, chunk_size) + self.send_response(HTTPStatus.OK) + self.send_header("Content-Type", "application/octet-stream") + self.send_header("Cache-Control", "no-store") + self.send_header("Connection", "close") + self.end_headers() + headers_sent = True + first = True + while True: + frame_header = response.read(4) + if not frame_header: + break + if len(frame_header) != 4: + raise RuntimeError("truncated Qwen PCM frame header") + frame_length = int.from_bytes(frame_header, "big") + if frame_length == 0: + break + if frame_length > MAX_AUDIO_BYTES: + raise RuntimeError("Qwen PCM frame is too large") + remaining = frame_length + while remaining: + chunk = response.read(min(16384, remaining)) + if not chunk: + raise RuntimeError("truncated Qwen PCM frame") + if first: + print("tts-gateway: first Qwen PCM chunk in " + f"{time.monotonic() - started:.2f}s") + first = False + self.wfile.write(chunk) + self.wfile.flush() + remaining -= len(chunk) + with STATE_LOCK: + STATE["last_backend"] = "qwen3-tts-1.7b-stream" + STATE["last_error"] = None + except Exception as exc: + with STATE_LOCK: + STATE["last_error"] = type(exc).__name__ + # Once PCM started, simply close the truncated response. Sending + # JSON into the audio stream would produce loud corrupt samples. + if not headers_sent and not self.wfile.closed: + try: + self.send_json(HTTPStatus.SERVICE_UNAVAILABLE, + {"error": "local PCM stream failed"}) + except (OSError, BrokenPipeError): + pass + finally: + if connection is not None: + connection.close() + SYNTHESIS_LOCK.release() + self.close_connection = True + if __name__ == "__main__": print(f"TTS gateway ready on {HOST}:{PORT}; primary={QWEN_TTS_VOICE}; fallback=Piper") diff --git a/router/ai_profile_router.py b/router/ai_profile_router.py index 4277d2d..fdcb9af 100755 --- a/router/ai_profile_router.py +++ b/router/ai_profile_router.py @@ -1519,6 +1519,13 @@ class Handler(BaseHTTPRequestHandler): else: self._send_error(503, "Sprachausgabe ist nicht installiert", "server_error", "feature_disabled") + elif (path == "/v1/audio/speech/pcm-stream" + and self.command == "POST"): + if ENABLE_TTS: + self._speech_pcm_stream() + else: + self._send_error(503, "Sprachausgabe ist nicht installiert", + "server_error", "feature_disabled") elif path == "/v1/audio/transcriptions" and self.command == "POST": if ENABLE_STT: self._transcribe() @@ -2062,6 +2069,74 @@ class Handler(BaseHTTPRequestHandler): self.end_headers() self.wfile.write(audio) + def _speech_pcm_stream(self) -> None: + """Pass through Qwen's native 24 kHz PCM stream without buffering.""" + try: + body = self._read_body() + data = json.loads(body) + except ValueError as exc: + self._send_error(400, str(exc) or "ungültiges JSON", + "invalid_request_error", "invalid_body") + return + if not isinstance(data, dict): + self._send_error(400, "Request muss ein JSON-Objekt sein", + "invalid_request_error", "invalid_request") + return + text = data.get("input", data.get("text")) + if not isinstance(text, str) or not text.strip() or len(text) > 8000: + self._send_error(400, "'input' fehlt, ist leer oder zu lang", + "invalid_request_error", "invalid_input") + return + voice = data.get("voice", TTS_DEFAULT_VOICE) + if voice not in TTS_VOICES: + self._send_error(400, f"ungültige Stimme: {voice!r}", + "invalid_request_error", "invalid_voice") + return + + hostport = TTS_WORKER_URL.split("://", 1)[-1] + host, _, port = hostport.partition(":") + connection = None + headers_sent = False + try: + connection = http.client.HTTPConnection( + host, int(port) if port else 80, timeout=TTS_CONNECT_TIMEOUT) + connection.request( + "POST", "/tts/pcm-stream", body=body, + headers={"Content-Type": "application/json", + "Accept": "application/octet-stream"}) + connection.sock.settimeout(TTS_TIMEOUT) + response = connection.getresponse() + if response.status != 200: + message = response.read(512).decode(errors="replace") + self._send_error(503, + f"TTS-Stream fehlgeschlagen: {message}", + "server_error", "tts_failed") + return + self._last_code = 200 + self.send_response(200) + self.send_header("Content-Type", "application/octet-stream") + self.send_header("Cache-Control", "no-store") + self.send_header("Connection", "close") + self.end_headers() + headers_sent = True + while True: + chunk = response.read1(16384) + if not chunk: + break + self.wfile.write(chunk) + self.wfile.flush() + except (OSError, http.client.HTTPException) as exc: + if not headers_sent: + try: + self._send_error(502, f"TTS-Worker nicht erreichbar: {exc}", + "server_error", "tts_unavailable") + except (OSError, BrokenPipeError): + pass + finally: + if connection is not None: + connection.close() + self.close_connection = True + # ---------- Audio-Discovery ---------- def _audio_models_payload(self) -> dict: