first commit
This commit is contained in:
@@ -0,0 +1,25 @@
|
||||
"""Custom integration for Qwen ASR Speech-to-Text."""
|
||||
from __future__ import annotations
|
||||
|
||||
from homeassistant.config_entries import ConfigEntry
|
||||
from homeassistant.core import HomeAssistant
|
||||
|
||||
from .const import DOMAIN
|
||||
|
||||
PLATFORMS = ["stt"]
|
||||
|
||||
async def async_setup_entry(hass: HomeAssistant, entry: ConfigEntry) -> bool:
|
||||
"""Set up Qwen ASR from a config entry."""
|
||||
hass.data.setdefault(DOMAIN, {})
|
||||
hass.data[DOMAIN][entry.entry_id] = entry.data
|
||||
|
||||
await hass.config_entries.async_forward_entry_setups(entry, PLATFORMS)
|
||||
return True
|
||||
|
||||
async def async_unload_entry(hass: HomeAssistant, entry: ConfigEntry) -> bool:
|
||||
"""Unload a config entry."""
|
||||
unload_ok = await hass.config_entries.async_unload_platforms(entry, PLATFORMS)
|
||||
if unload_ok:
|
||||
hass.data[DOMAIN].pop(entry.entry_id)
|
||||
|
||||
return unload_ok
|
||||
@@ -0,0 +1,176 @@
|
||||
"""Config flow for Qwen ASR Speech-to-Text integration."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Any
|
||||
|
||||
from aiohttp import ClientError, WSMsgType, WSServerHandshakeError
|
||||
import voluptuous as vol
|
||||
|
||||
from homeassistant import config_entries
|
||||
from homeassistant.data_entry_flow import FlowResult
|
||||
from homeassistant.helpers.aiohttp_client import async_get_clientsession
|
||||
|
||||
from .const import (
|
||||
CONF_API_KEY,
|
||||
CONF_API_URL,
|
||||
CONF_ENABLE_SERVER_VAD,
|
||||
CONF_MODEL,
|
||||
CONF_REGION,
|
||||
CONF_TIMEOUT,
|
||||
CONF_VAD_SILENCE_DURATION_MS,
|
||||
CONF_VAD_THRESHOLD,
|
||||
DEFAULT_API_URL_BEIJING,
|
||||
DEFAULT_API_URL_SINGAPORE,
|
||||
DEFAULT_MODEL,
|
||||
DEFAULT_REGION,
|
||||
DOMAIN,
|
||||
SUPPORTED_MODELS,
|
||||
SUPPORTED_REGIONS,
|
||||
)
|
||||
|
||||
_LOGGER = logging.getLogger(__name__)
|
||||
|
||||
STEP_USER_DATA_SCHEMA = vol.Schema(
|
||||
{
|
||||
vol.Required(CONF_API_KEY): str,
|
||||
vol.Optional(CONF_MODEL, default=DEFAULT_MODEL): vol.In(SUPPORTED_MODELS),
|
||||
vol.Optional(CONF_REGION, default=DEFAULT_REGION): vol.In(SUPPORTED_REGIONS),
|
||||
}
|
||||
)
|
||||
|
||||
STEP_OPTIONS_DATA_SCHEMA = vol.Schema(
|
||||
{
|
||||
vol.Optional(CONF_API_URL): str,
|
||||
vol.Optional(CONF_ENABLE_SERVER_VAD, default=False): bool,
|
||||
vol.Optional(CONF_VAD_THRESHOLD, default=0.0): vol.Coerce(float),
|
||||
vol.Optional(CONF_VAD_SILENCE_DURATION_MS, default=400): vol.Coerce(int),
|
||||
vol.Optional(CONF_TIMEOUT, default=30): vol.Coerce(int),
|
||||
}
|
||||
)
|
||||
|
||||
class QwenSTTConfigFlow(config_entries.ConfigFlow, domain=DOMAIN):
|
||||
"""Handle a config flow for Qwen ASR Speech-to-Text."""
|
||||
|
||||
VERSION = 1
|
||||
|
||||
|
||||
async def _async_validate(self, user_input: dict[str, Any]) -> None:
|
||||
api_key = user_input[CONF_API_KEY]
|
||||
model = user_input.get(CONF_MODEL, DEFAULT_MODEL)
|
||||
region = user_input.get(CONF_REGION, DEFAULT_REGION)
|
||||
api_url = (
|
||||
DEFAULT_API_URL_SINGAPORE if region == "singapore" else DEFAULT_API_URL_BEIJING
|
||||
)
|
||||
uri = f"{api_url}/realtime?model={model}"
|
||||
headers = {
|
||||
"Authorization": f"Bearer {api_key}",
|
||||
"OpenAI-Beta": "realtime=v1",
|
||||
}
|
||||
|
||||
session = async_get_clientsession(self.hass)
|
||||
async with session.ws_connect(uri, headers=headers, heartbeat=30) as ws:
|
||||
await ws.send_json(
|
||||
{
|
||||
"event_id": "event_validate",
|
||||
"type": "session.update",
|
||||
"session": {
|
||||
"modalities": ["text"],
|
||||
"input_audio_format": "pcm",
|
||||
"sample_rate": 16000,
|
||||
"input_audio_transcription": {"language": "zh"},
|
||||
"turn_detection": None,
|
||||
},
|
||||
}
|
||||
)
|
||||
async with asyncio.timeout(5):
|
||||
async for msg in ws:
|
||||
if msg.type != WSMsgType.TEXT:
|
||||
continue
|
||||
try:
|
||||
data = msg.json()
|
||||
except Exception:
|
||||
continue
|
||||
msg_type = data.get("type")
|
||||
if msg_type in ("session.created", "session.updated"):
|
||||
return
|
||||
if msg_type == "error":
|
||||
raise ClientError(str(data.get("error", {})))
|
||||
|
||||
async def async_step_user(
|
||||
self, user_input: dict[str, Any] | None = None
|
||||
) -> FlowResult:
|
||||
"""Handle the initial step."""
|
||||
if user_input is None:
|
||||
return self.async_show_form(
|
||||
step_id="user", data_schema=STEP_USER_DATA_SCHEMA
|
||||
)
|
||||
|
||||
errors: dict[str, str] = {}
|
||||
|
||||
# Check if already configured
|
||||
if self._async_current_entries():
|
||||
return self.async_abort(reason="already_configured")
|
||||
|
||||
try:
|
||||
await self._async_validate(user_input)
|
||||
except WSServerHandshakeError as err:
|
||||
if err.status in (401, 403):
|
||||
errors["base"] = "invalid_auth"
|
||||
else:
|
||||
errors["base"] = "cannot_connect"
|
||||
except TimeoutError:
|
||||
errors["base"] = "cannot_connect"
|
||||
except ClientError:
|
||||
errors["base"] = "cannot_connect"
|
||||
except Exception:
|
||||
_LOGGER.exception("Unexpected error validating Qwen ASR configuration")
|
||||
errors["base"] = "cannot_connect"
|
||||
|
||||
if errors:
|
||||
return self.async_show_form(
|
||||
step_id="user",
|
||||
data_schema=STEP_USER_DATA_SCHEMA,
|
||||
errors=errors,
|
||||
)
|
||||
|
||||
return self.async_create_entry(title="Qwen ASR", data=user_input)
|
||||
|
||||
@staticmethod
|
||||
def async_get_options_flow(config_entry: config_entries.ConfigEntry) -> config_entries.OptionsFlow:
|
||||
return QwenSTTOptionsFlow(config_entry)
|
||||
|
||||
|
||||
class QwenSTTOptionsFlow(config_entries.OptionsFlow):
|
||||
def __init__(self, config_entry: config_entries.ConfigEntry) -> None:
|
||||
super().__init__(config_entry)
|
||||
|
||||
async def async_step_init(self, user_input: dict[str, Any] | None = None) -> FlowResult:
|
||||
if user_input is not None:
|
||||
return self.async_create_entry(title="", data=user_input)
|
||||
|
||||
opts = self.config_entry.options
|
||||
schema = vol.Schema(
|
||||
{
|
||||
vol.Optional(CONF_API_URL, default=opts.get(CONF_API_URL, "")): str,
|
||||
vol.Optional(
|
||||
CONF_ENABLE_SERVER_VAD,
|
||||
default=opts.get(CONF_ENABLE_SERVER_VAD, False),
|
||||
): bool,
|
||||
vol.Optional(
|
||||
CONF_VAD_THRESHOLD,
|
||||
default=opts.get(CONF_VAD_THRESHOLD, 0.0),
|
||||
): vol.Coerce(float),
|
||||
vol.Optional(
|
||||
CONF_VAD_SILENCE_DURATION_MS,
|
||||
default=opts.get(CONF_VAD_SILENCE_DURATION_MS, 400),
|
||||
): vol.Coerce(int),
|
||||
vol.Optional(
|
||||
CONF_TIMEOUT,
|
||||
default=opts.get(CONF_TIMEOUT, 30),
|
||||
): vol.Coerce(int),
|
||||
}
|
||||
)
|
||||
|
||||
return self.async_show_form(step_id="init", data_schema=schema)
|
||||
@@ -0,0 +1,48 @@
|
||||
"""Constants for the Qwen ASR Speech-to-Text integration."""
|
||||
|
||||
DOMAIN = "hass_stt_qwen"
|
||||
|
||||
CONF_API_KEY = "api_key"
|
||||
CONF_API_URL = "api_url"
|
||||
CONF_MODEL = "model"
|
||||
CONF_REGION = "region"
|
||||
CONF_ENABLE_SERVER_VAD = "enable_server_vad"
|
||||
CONF_VAD_THRESHOLD = "vad_threshold"
|
||||
CONF_VAD_SILENCE_DURATION_MS = "vad_silence_duration_ms"
|
||||
CONF_TIMEOUT = "timeout"
|
||||
|
||||
DEFAULT_API_URL_BEIJING = "wss://dashscope.aliyuncs.com/api-ws/v1"
|
||||
DEFAULT_API_URL_SINGAPORE = "wss://dashscope-intl.aliyuncs.com/api-ws/v1"
|
||||
DEFAULT_MODEL = "qwen3-asr-flash-realtime"
|
||||
DEFAULT_REGION = "beijing"
|
||||
|
||||
SUPPORTED_MODELS = [
|
||||
"qwen3-asr-flash-realtime",
|
||||
"qwen3-asr-flash-realtime-2025-10-27"
|
||||
]
|
||||
|
||||
SUPPORTED_REGIONS = [
|
||||
"beijing",
|
||||
"singapore",
|
||||
]
|
||||
|
||||
# Qwen ASR supports multiple languages
|
||||
SUPPORTED_LANGUAGES = [
|
||||
"zh", # Chinese
|
||||
"en", # English
|
||||
"yue", # Cantonese
|
||||
"ja", # Japanese
|
||||
"ko", # Korean
|
||||
"es", # Spanish
|
||||
"fr", # French
|
||||
"de", # German
|
||||
"it", # Italian
|
||||
"pt", # Portuguese
|
||||
"ru", # Russian
|
||||
"ar", # Arabic
|
||||
"tr", # Turkish
|
||||
"vi", # Vietnamese
|
||||
"th", # Thai
|
||||
"id", # Indonesian
|
||||
"ms", # Malay
|
||||
]
|
||||
@@ -0,0 +1,10 @@
|
||||
{
|
||||
"domain": "hass_stt_qwen",
|
||||
"name": "阿里通义千问语音识别",
|
||||
"codeowners": ["@CircleLiu"],
|
||||
"config_flow": true,
|
||||
"documentation": "https://bbs.nextrt.com/d/2-home-assistant-a-li-yun-tong-yi-qian-wen-stt-yu-yin-shi-bie-cha-jian",
|
||||
"iot_class": "cloud_push",
|
||||
"issue_tracker": "https://bbs.nextrt.com/d/2-home-assistant-a-li-yun-tong-yi-qian-wen-stt-yu-yin-shi-bie-cha-jian",
|
||||
"version": "0.1.0"
|
||||
}
|
||||
@@ -0,0 +1,372 @@
|
||||
"""WebSocket client for Qwen3 ASR."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import base64
|
||||
from collections.abc import AsyncIterable
|
||||
import json
|
||||
import logging
|
||||
import time
|
||||
from typing import Final
|
||||
|
||||
from aiohttp import ClientError, WSCloseCode, WSMsgType
|
||||
|
||||
from homeassistant.components.stt import SpeechMetadata, SpeechResult, SpeechResultState
|
||||
|
||||
_LOGGER = logging.getLogger(__name__)
|
||||
|
||||
# Maximum time to wait for a response (in seconds)
|
||||
WEBSOCKET_TIMEOUT: Final = 30
|
||||
|
||||
|
||||
class Qwen3AsrClientError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
class Qwen3AsrClientTimeout(Qwen3AsrClientError):
|
||||
pass
|
||||
|
||||
|
||||
class Qwen3AsrClient:
|
||||
"""WebSocket client for Qwen3 ASR API."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
client,
|
||||
api_key: str,
|
||||
api_url: str,
|
||||
model: str,
|
||||
) -> None:
|
||||
"""Initialize the WebSocket client."""
|
||||
self.client = client
|
||||
self.api_key = api_key
|
||||
self.api_url = api_url
|
||||
self.model = model
|
||||
self.start_time = 0.0
|
||||
self.ws = None
|
||||
self.timeout = WEBSOCKET_TIMEOUT
|
||||
self.enable_server_vad = False
|
||||
self.vad_threshold = 0.0
|
||||
self.vad_silence_duration_ms = 400
|
||||
|
||||
def _new_event_id(self) -> str:
|
||||
return f"event_{int(time.time() * 1000)}"
|
||||
|
||||
def _normalize_language(self, language: str | None) -> str:
|
||||
if not language:
|
||||
return "zh"
|
||||
lower = language.lower()
|
||||
if lower.startswith("zh"):
|
||||
return "zh"
|
||||
return lower
|
||||
|
||||
async def _iter_pcm16(self, stream: AsyncIterable[bytes]) -> AsyncIterable[bytes]:
|
||||
buf = bytearray()
|
||||
state = "probe"
|
||||
pos = 0
|
||||
data_remaining: int | None = None
|
||||
|
||||
async for chunk in stream:
|
||||
if state == "raw":
|
||||
if chunk:
|
||||
yield chunk
|
||||
continue
|
||||
|
||||
if chunk:
|
||||
buf.extend(chunk)
|
||||
|
||||
if state == "probe":
|
||||
if len(buf) < 12:
|
||||
continue
|
||||
if buf[:4] != b"RIFF" or buf[8:12] != b"WAVE":
|
||||
yield bytes(buf)
|
||||
buf.clear()
|
||||
state = "raw"
|
||||
continue
|
||||
state = "wav_chunks"
|
||||
pos = 12
|
||||
|
||||
if state == "wav_chunks":
|
||||
while True:
|
||||
if len(buf) < pos + 8:
|
||||
break
|
||||
chunk_id = bytes(buf[pos : pos + 4])
|
||||
chunk_size = int.from_bytes(buf[pos + 4 : pos + 8], "little")
|
||||
end = pos + 8 + chunk_size + (chunk_size % 2)
|
||||
if len(buf) < end:
|
||||
break
|
||||
if chunk_id == b"data":
|
||||
del buf[: pos + 8]
|
||||
data_remaining = chunk_size
|
||||
state = "wav_data"
|
||||
break
|
||||
pos = end
|
||||
if pos > 65536:
|
||||
yield bytes(buf)
|
||||
buf.clear()
|
||||
state = "raw"
|
||||
break
|
||||
|
||||
if state == "wav_data" and data_remaining is not None:
|
||||
while buf and data_remaining > 0:
|
||||
take = min(len(buf), data_remaining)
|
||||
yield bytes(buf[:take])
|
||||
del buf[:take]
|
||||
data_remaining -= take
|
||||
if data_remaining == 0:
|
||||
if buf:
|
||||
yield bytes(buf)
|
||||
buf.clear()
|
||||
state = "raw"
|
||||
|
||||
async def _send_audio_stream(self, stream: AsyncIterable[bytes]) -> None:
|
||||
"""Send audio chunks to WebSocket server."""
|
||||
try:
|
||||
chunk_count = 0
|
||||
async for chunk in self._iter_pcm16(stream):
|
||||
if not chunk or self.ws.closed:
|
||||
_LOGGER.debug("Stopping audio send: empty chunk or closed connection")
|
||||
break
|
||||
# Audio data must be base64 encoded
|
||||
b64 = base64.b64encode(chunk).decode("utf-8")
|
||||
await self.ws.send_json(
|
||||
{
|
||||
"event_id": self._new_event_id(),
|
||||
"type": "input_audio_buffer.append",
|
||||
"audio": b64,
|
||||
}
|
||||
)
|
||||
chunk_count += 1
|
||||
_LOGGER.debug("Audio chunk #%d sent (%d bytes)", chunk_count, len(chunk))
|
||||
|
||||
if not self.ws.closed:
|
||||
# Signal the end of the audio stream to the server
|
||||
_LOGGER.info("Sent %d audio chunks, sending finish", chunk_count)
|
||||
if not self.enable_server_vad:
|
||||
await self.ws.send_json(
|
||||
{
|
||||
"event_id": self._new_event_id(),
|
||||
"type": "input_audio_buffer.commit",
|
||||
}
|
||||
)
|
||||
await self.ws.send_json(
|
||||
{
|
||||
"event_id": self._new_event_id(),
|
||||
"type": "session.finish",
|
||||
}
|
||||
)
|
||||
|
||||
# Set start time after sending all audio data
|
||||
self.start_time = time.perf_counter()
|
||||
|
||||
except asyncio.CancelledError:
|
||||
_LOGGER.debug("send_audio() was cancelled")
|
||||
raise
|
||||
except Exception:
|
||||
_LOGGER.exception("Error sending audio")
|
||||
if not self.ws.closed:
|
||||
await self.ws.close(
|
||||
code=WSCloseCode.INTERNAL_ERROR,
|
||||
message=b"Error sending audio",
|
||||
)
|
||||
raise
|
||||
|
||||
async def _receive_transcription(self, send_task: asyncio.Task) -> str:
|
||||
"""Receive transcription results from WebSocket server."""
|
||||
final_text: str | None = None
|
||||
try:
|
||||
async with asyncio.timeout(self.timeout):
|
||||
async for msg in self.ws:
|
||||
if msg.type == WSMsgType.TEXT:
|
||||
data = json.loads(msg.data)
|
||||
msg_type = data.get("type")
|
||||
_LOGGER.debug("Received message type: %s", msg_type)
|
||||
|
||||
if msg_type == "session.created":
|
||||
_LOGGER.info("Session created: %s", data.get("session", {}).get("id"))
|
||||
elif msg_type == "session.updated":
|
||||
_LOGGER.info("Session updated")
|
||||
elif msg_type == "conversation.item.input_audio_transcription.text":
|
||||
text = data.get("text") or data.get("stash") or ""
|
||||
if text:
|
||||
_LOGGER.debug('Intermediate transcription: "%s"', text)
|
||||
elif msg_type == "conversation.item.input_audio_transcription.completed":
|
||||
# Get final transcription
|
||||
final_text = (data.get("transcript") or data.get("text") or "").strip()
|
||||
if self.start_time > 0:
|
||||
duration = time.perf_counter() - self.start_time
|
||||
_LOGGER.info(
|
||||
"Transcription processing duration: %.2f seconds",
|
||||
duration,
|
||||
)
|
||||
_LOGGER.info('Final transcription received: "%s"', final_text)
|
||||
return final_text or ""
|
||||
elif msg_type == "input_audio_buffer.speech_started":
|
||||
_LOGGER.info("Speech started")
|
||||
elif msg_type == "input_audio_buffer.speech_stopped":
|
||||
_LOGGER.info("Speech stopped")
|
||||
elif msg_type == "input_audio_buffer.committed":
|
||||
_LOGGER.info("Audio buffer committed")
|
||||
elif msg_type == "session.finished":
|
||||
final_text = (data.get("transcript") or final_text or "").strip()
|
||||
_LOGGER.info('Session finished: "%s"', final_text)
|
||||
return final_text or ""
|
||||
elif msg_type == "error":
|
||||
error_msg = data.get("error", {})
|
||||
raise Qwen3AsrClientError(str(error_msg))
|
||||
else:
|
||||
_LOGGER.debug("Unhandled message type: %s, data: %s", msg_type, data)
|
||||
elif msg.type == WSMsgType.BINARY:
|
||||
_LOGGER.warning("Received unexpected binary message (%d bytes)", len(msg.data))
|
||||
elif msg.type == WSMsgType.ERROR:
|
||||
raise Qwen3AsrClientError(str(self.ws.exception()))
|
||||
elif msg.type == WSMsgType.CLOSED:
|
||||
_LOGGER.info("WebSocket closed by server")
|
||||
break
|
||||
else:
|
||||
_LOGGER.debug("Received message of type: %s", msg.type)
|
||||
except TimeoutError:
|
||||
raise Qwen3AsrClientTimeout("Timeout waiting for transcription response")
|
||||
except asyncio.CancelledError:
|
||||
_LOGGER.debug("receive_transcription() was cancelled")
|
||||
raise
|
||||
except Qwen3AsrClientError:
|
||||
raise
|
||||
except Exception as err:
|
||||
raise Qwen3AsrClientError(str(err)) from err
|
||||
finally:
|
||||
if not send_task.done():
|
||||
send_task.cancel()
|
||||
|
||||
if final_text is not None:
|
||||
return final_text
|
||||
raise Qwen3AsrClientError("WebSocket closed before transcription completed")
|
||||
|
||||
def _create_session_config(self, metadata: SpeechMetadata) -> dict:
|
||||
"""Create configuration for the transcription session."""
|
||||
language = self._normalize_language(metadata.language)
|
||||
if self.enable_server_vad:
|
||||
turn_detection = {
|
||||
"type": "server_vad",
|
||||
"threshold": self.vad_threshold,
|
||||
"silence_duration_ms": self.vad_silence_duration_ms,
|
||||
}
|
||||
else:
|
||||
turn_detection = None
|
||||
config = {
|
||||
"event_id": self._new_event_id(),
|
||||
"type": "session.update",
|
||||
"session": {
|
||||
"modalities": ["text"],
|
||||
"input_audio_format": "pcm",
|
||||
"sample_rate": 16000,
|
||||
"input_audio_transcription": {
|
||||
"language": language,
|
||||
},
|
||||
"turn_detection": turn_detection,
|
||||
},
|
||||
}
|
||||
return config
|
||||
|
||||
async def _handle_tasks(
|
||||
self, send_task: asyncio.Task, recv_task: asyncio.Task
|
||||
) -> None:
|
||||
"""Handle task completion and cancellation logic."""
|
||||
try:
|
||||
done, pending = await asyncio.wait(
|
||||
[send_task, recv_task],
|
||||
return_when=asyncio.FIRST_COMPLETED,
|
||||
)
|
||||
|
||||
# Handle task completion and cancellation
|
||||
if recv_task in done:
|
||||
_LOGGER.debug("Transcription finished - cancelling audio task")
|
||||
if not send_task.done():
|
||||
send_task.cancel()
|
||||
await asyncio.gather(send_task, return_exceptions=True)
|
||||
elif send_task in done:
|
||||
_LOGGER.debug("Audio finished - waiting for final transcription")
|
||||
await recv_task
|
||||
else:
|
||||
_LOGGER.warning(
|
||||
"Unexpected state in task completion, ensuring tasks are awaited/cancelled"
|
||||
)
|
||||
for task in pending:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
await asyncio.gather(*pending, return_exceptions=True)
|
||||
|
||||
for task in done:
|
||||
if task.exception():
|
||||
_LOGGER.error("Task completed with exception: %s", task.exception())
|
||||
finally:
|
||||
if not self.ws.closed:
|
||||
try:
|
||||
await self.ws.close()
|
||||
_LOGGER.debug("WebSocket closed cleanly")
|
||||
except Exception:
|
||||
_LOGGER.exception("Error closing WebSocket connection")
|
||||
|
||||
async def async_process_audio_stream(
|
||||
self, metadata: SpeechMetadata, stream: AsyncIterable[bytes]
|
||||
) -> SpeechResult:
|
||||
"""Process audio stream via WebSocket to Qwen3 ASR Realtime API."""
|
||||
|
||||
# Construct WebSocket URL for realtime API
|
||||
uri = f"{self.api_url}/realtime?model={self.model}"
|
||||
headers = {
|
||||
"Authorization": f"Bearer {self.api_key}",
|
||||
"OpenAI-Beta": "realtime=v1",
|
||||
}
|
||||
|
||||
_LOGGER.info(
|
||||
"Starting Qwen3 ASR transcription: language=%s, format=%s, sample_rate=%s, channels=%s",
|
||||
metadata.language,
|
||||
metadata.format,
|
||||
metadata.sample_rate,
|
||||
metadata.channel,
|
||||
)
|
||||
|
||||
try:
|
||||
_LOGGER.debug("Opening WebSocket connection to %s", uri)
|
||||
async with self.client.ws_connect(uri, headers=headers, heartbeat=30) as ws:
|
||||
self.ws = ws
|
||||
self.start_time = 0 # Reset start_time
|
||||
|
||||
config = self._create_session_config(metadata)
|
||||
_LOGGER.debug("Sending session configuration: %s", config)
|
||||
await ws.send_json(config)
|
||||
|
||||
# Create and manage concurrent tasks
|
||||
send_task = asyncio.create_task(self._send_audio_stream(stream))
|
||||
recv_task = asyncio.create_task(self._receive_transcription(send_task))
|
||||
|
||||
# Handle tasks completion
|
||||
await self._handle_tasks(send_task, recv_task)
|
||||
|
||||
# Process final result
|
||||
if not recv_task.done() or recv_task.cancelled():
|
||||
return SpeechResult("", SpeechResultState.ERROR)
|
||||
|
||||
exc = recv_task.exception()
|
||||
if exc:
|
||||
_LOGGER.error("Transcription failed: %s", exc)
|
||||
return SpeechResult("", SpeechResultState.ERROR)
|
||||
|
||||
final_text = recv_task.result().strip()
|
||||
|
||||
_LOGGER.info('Transcription completed successfully: "%s"', final_text)
|
||||
|
||||
if not final_text:
|
||||
_LOGGER.warning("Qwen3 ASR transcription resulted in empty text")
|
||||
return SpeechResult("", SpeechResultState.SUCCESS)
|
||||
|
||||
return SpeechResult(final_text, SpeechResultState.SUCCESS)
|
||||
|
||||
except ClientError as err:
|
||||
_LOGGER.error("WebSocket connection error: %s", err)
|
||||
return SpeechResult("", SpeechResultState.ERROR)
|
||||
except Exception:
|
||||
_LOGGER.exception("Unexpected error in WebSocket communication")
|
||||
return SpeechResult("", SpeechResultState.ERROR)
|
||||
@@ -0,0 +1,181 @@
|
||||
"""Setting up QwenSTTProvider."""
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import AsyncIterable
|
||||
import logging
|
||||
|
||||
from homeassistant.components.stt import (
|
||||
AudioBitRates,
|
||||
AudioChannels,
|
||||
AudioCodecs,
|
||||
AudioFormats,
|
||||
AudioSampleRates,
|
||||
SpeechMetadata,
|
||||
SpeechResult,
|
||||
SpeechResultState,
|
||||
)
|
||||
from homeassistant.components.stt import SpeechToTextEntity
|
||||
|
||||
from homeassistant.config_entries import ConfigEntry
|
||||
from homeassistant.core import HomeAssistant
|
||||
from homeassistant.helpers.aiohttp_client import async_get_clientsession
|
||||
from homeassistant.helpers.entity_platform import AddEntitiesCallback
|
||||
|
||||
from .const import (
|
||||
CONF_API_KEY,
|
||||
CONF_API_URL,
|
||||
CONF_ENABLE_SERVER_VAD,
|
||||
CONF_MODEL,
|
||||
CONF_REGION,
|
||||
CONF_TIMEOUT,
|
||||
CONF_VAD_SILENCE_DURATION_MS,
|
||||
CONF_VAD_THRESHOLD,
|
||||
DEFAULT_API_URL_BEIJING,
|
||||
DEFAULT_API_URL_SINGAPORE,
|
||||
DEFAULT_MODEL,
|
||||
DEFAULT_REGION,
|
||||
DOMAIN,
|
||||
SUPPORTED_LANGUAGES,
|
||||
)
|
||||
from .qwen3_asr_client import Qwen3AsrClient
|
||||
|
||||
_LOGGER = logging.getLogger(__name__)
|
||||
|
||||
async def async_setup_entry(
|
||||
hass: HomeAssistant,
|
||||
config_entry: ConfigEntry,
|
||||
async_add_entities: AddEntitiesCallback,
|
||||
) -> None:
|
||||
"""Set up Qwen STT from a config entry."""
|
||||
data = config_entry.data
|
||||
options = config_entry.options
|
||||
|
||||
api_key = data[CONF_API_KEY]
|
||||
model = data.get(CONF_MODEL, DEFAULT_MODEL)
|
||||
region = data.get(CONF_REGION, DEFAULT_REGION)
|
||||
|
||||
api_url = options.get(CONF_API_URL) or DEFAULT_API_URL_SINGAPORE if region == "singapore" else DEFAULT_API_URL_BEIJING
|
||||
enable_server_vad = options.get(CONF_ENABLE_SERVER_VAD, False)
|
||||
vad_threshold = options.get(CONF_VAD_THRESHOLD, 0.0)
|
||||
vad_silence_duration_ms = options.get(CONF_VAD_SILENCE_DURATION_MS, 400)
|
||||
timeout = options.get(CONF_TIMEOUT, 30)
|
||||
|
||||
async_add_entities(
|
||||
[
|
||||
QwenSTTEntity(
|
||||
hass,
|
||||
config_entry,
|
||||
api_key,
|
||||
api_url,
|
||||
model,
|
||||
enable_server_vad=enable_server_vad,
|
||||
vad_threshold=vad_threshold,
|
||||
vad_silence_duration_ms=vad_silence_duration_ms,
|
||||
timeout=timeout,
|
||||
)
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
class QwenSTTEntity(SpeechToTextEntity):
|
||||
"""The Qwen STT entity."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
hass: HomeAssistant,
|
||||
config_entry: ConfigEntry,
|
||||
api_key: str,
|
||||
api_url: str,
|
||||
model: str,
|
||||
*,
|
||||
enable_server_vad: bool,
|
||||
vad_threshold: float,
|
||||
vad_silence_duration_ms: int,
|
||||
timeout: int,
|
||||
) -> None:
|
||||
"""Init Qwen STT service."""
|
||||
self.hass = hass
|
||||
self._config_entry = config_entry
|
||||
self._api_key = api_key
|
||||
self._api_url = api_url
|
||||
self._model = model
|
||||
self._enable_server_vad = enable_server_vad
|
||||
self._vad_threshold = vad_threshold
|
||||
self._vad_silence_duration_ms = vad_silence_duration_ms
|
||||
self._timeout = timeout
|
||||
self._client = self._create_client()
|
||||
|
||||
self._attr_name = f"Qwen ASR ({model})"
|
||||
self._attr_unique_id = f"{config_entry.entry_id}_{model}"
|
||||
|
||||
@property
|
||||
def supported_languages(self) -> list[str]:
|
||||
"""Return a list of supported languages."""
|
||||
return SUPPORTED_LANGUAGES
|
||||
|
||||
@property
|
||||
def supported_formats(self) -> list[AudioFormats]:
|
||||
"""Return a list of supported formats."""
|
||||
return [AudioFormats.WAV]
|
||||
|
||||
@property
|
||||
def supported_codecs(self) -> list[AudioCodecs]:
|
||||
"""Return a list of supported codecs."""
|
||||
return [AudioCodecs.PCM]
|
||||
|
||||
@property
|
||||
def supported_bit_rates(self) -> list[AudioBitRates]:
|
||||
"""Return a list of supported bitrates."""
|
||||
return [AudioBitRates.BITRATE_16]
|
||||
|
||||
@property
|
||||
def supported_sample_rates(self) -> list[AudioSampleRates]:
|
||||
"""Return a list of supported samplerates."""
|
||||
return [AudioSampleRates.SAMPLERATE_16000]
|
||||
|
||||
@property
|
||||
def supported_channels(self) -> list[AudioChannels]:
|
||||
"""Return a list of supported channels."""
|
||||
return [AudioChannels.CHANNEL_MONO]
|
||||
|
||||
def _create_client(self):
|
||||
"""Create and return the appropriate client based on model."""
|
||||
session = async_get_clientsession(self.hass)
|
||||
|
||||
if self._model.startswith("qwen3-asr"):
|
||||
client = Qwen3AsrClient(
|
||||
session,
|
||||
self._api_key,
|
||||
self._api_url,
|
||||
self._model,
|
||||
)
|
||||
client.enable_server_vad = self._enable_server_vad
|
||||
client.vad_threshold = self._vad_threshold
|
||||
client.vad_silence_duration_ms = self._vad_silence_duration_ms
|
||||
client.timeout = self._timeout
|
||||
return client
|
||||
else:
|
||||
_LOGGER.error("Unsupported model: %s", self._model)
|
||||
raise ValueError(f"Unsupported model: {self._model}")
|
||||
|
||||
async def async_process_audio_stream(
|
||||
self, metadata: SpeechMetadata, stream: AsyncIterable[bytes]
|
||||
) -> SpeechResult:
|
||||
"""Process audio stream using the configured client."""
|
||||
if (
|
||||
metadata.format not in self.supported_formats
|
||||
or metadata.codec not in self.supported_codecs
|
||||
or metadata.bit_rate not in self.supported_bit_rates
|
||||
or metadata.sample_rate not in self.supported_sample_rates
|
||||
or metadata.channel not in self.supported_channels
|
||||
):
|
||||
_LOGGER.error(
|
||||
"Unsupported audio metadata: format=%s codec=%s bit_rate=%s sample_rate=%s channel=%s",
|
||||
metadata.format,
|
||||
metadata.codec,
|
||||
metadata.bit_rate,
|
||||
metadata.sample_rate,
|
||||
metadata.channel,
|
||||
)
|
||||
return SpeechResult("", state=SpeechResultState.ERROR)
|
||||
return await self._client.async_process_audio_stream(metadata, stream)
|
||||
@@ -0,0 +1,36 @@
|
||||
{
|
||||
"config": {
|
||||
"step": {
|
||||
"user": {
|
||||
"title": "Connect to Qwen ASR",
|
||||
"description": "Enter your Alibaba Cloud DashScope API details.",
|
||||
"data": {
|
||||
"api_key": "API Key",
|
||||
"model": "Model",
|
||||
"region": "Region"
|
||||
}
|
||||
}
|
||||
},
|
||||
"error": {
|
||||
"cannot_connect": "Failed to connect",
|
||||
"invalid_auth": "Invalid authentication"
|
||||
},
|
||||
"abort": {
|
||||
"already_configured": "Service is already configured"
|
||||
}
|
||||
},
|
||||
"options": {
|
||||
"step": {
|
||||
"init": {
|
||||
"title": "Qwen ASR options",
|
||||
"data": {
|
||||
"api_url": "WebSocket endpoint (optional)",
|
||||
"enable_server_vad": "Enable Server VAD",
|
||||
"vad_threshold": "VAD threshold",
|
||||
"vad_silence_duration_ms": "VAD silence duration (ms)",
|
||||
"timeout": "Timeout (seconds)"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
{
|
||||
"config": {
|
||||
"step": {
|
||||
"user": {
|
||||
"title": "连接到阿里云 Qwen ASR",
|
||||
"description": "请输入您的阿里云 DashScope API 信息。",
|
||||
"data": {
|
||||
"api_key": "API 密钥",
|
||||
"model": "模型",
|
||||
"region": "区域"
|
||||
}
|
||||
}
|
||||
},
|
||||
"error": {
|
||||
"cannot_connect": "连接失败",
|
||||
"invalid_auth": "无效的身份验证"
|
||||
},
|
||||
"abort": {
|
||||
"already_configured": "服务已配置"
|
||||
}
|
||||
},
|
||||
"options": {
|
||||
"step": {
|
||||
"init": {
|
||||
"title": "Qwen ASR 选项",
|
||||
"data": {
|
||||
"api_url": "WebSocket 接入点(可选)",
|
||||
"enable_server_vad": "启用 Server VAD",
|
||||
"vad_threshold": "VAD 阈值",
|
||||
"vad_silence_duration_ms": "VAD 静默时长(毫秒)",
|
||||
"timeout": "超时(秒)"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user