Add native Qwen3 TTS streaming for Hermes

This commit is contained in:
Mikei386 committed 2026-09-05 16:06:27 +02:00
1 parent 5f3d064bb4
commit 61aa20cb52
7 files changed
+382 -4

No files matched your search

+116 -2
View File
@@ -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")