064b06ecf0
迁入 app/park:llm(对话)、rag(双路向量知识库)、tools/ai_tools(智能体工具)、 asr(语音识别)、s2s_bridge(实时语音桥)、vision_yolo/vision_llm(人脸/多模态)、 knowledge/*.md、vendor/s2s-cloud(s2s 云化栈);routers 补 /api/ai|kb|asr|vision| s2s|tools 端点。智能体提示词/企业名录改读 park_config 主数据源。pyproject 加重依赖 (dashscope/numpy/openai/torch/transformers/ultralytics/websockets/soundfile/scipy/ nltk/jinja2)。TestClient 冒烟:ai/chat(无 Key 走本地规则)、kb、s2s、tools、display 均 200。
144 lines
6.4 KiB
Python
144 lines
6.4 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""s2s 实时语音栈(VAD 本地 + 云 ASR/LLM/TTS)桥接
|
|
|
|
在 DPM 进程内以后台线程运行 s2s 的 RealtimeServer(默认 :8765 /v1/realtime)。
|
|
- 复用 live-avatar 云化栈(speech-to-speech v0.2.11 + cloud-migration 补丁)。
|
|
- 依赖已并入 backend/pyproject.toml;本模块负责按 DPM Settings 组装并启动。
|
|
- 前端语音对话页通过 WS 直连 ws://<host>:<S2S_PORT>/v1/realtime。
|
|
"""
|
|
|
|
import logging
|
|
import threading
|
|
import time
|
|
|
|
from .config import settings
|
|
|
|
log = logging.getLogger("dpm.s2s")
|
|
|
|
|
|
def _build_pool(stop_event: threading.Event):
|
|
"""按 DPM 配置构建 s2s 管线池(含 VAD/silero 加载 + LLM warmup,耗时约 30~60s)。"""
|
|
import sys
|
|
|
|
from speech_to_speech.api.openai_realtime.websocket_router import create_app
|
|
from speech_to_speech.s2s_pipeline import (
|
|
_build_realtime_pipeline_unit,
|
|
parse_arguments,
|
|
prepare_all_args,
|
|
)
|
|
|
|
# 与 start_backend.sh 等价的参数(api_key 留空 → handler 从 env DASHSCOPE_API_KEY 读取)
|
|
sys.argv = [
|
|
"dpm-s2s",
|
|
"--mode", "realtime",
|
|
"--stt", "dashscope-asr",
|
|
"--llm_backend", "chat-completions",
|
|
"--tts", "qwen3-cloud",
|
|
"--model_name", settings.S2S_LLM_MODEL,
|
|
"--responses_api_base_url", settings.LLM_BASE_URL,
|
|
"--responses_api_api_key", settings.DASHSCOPE_API_KEY,
|
|
"--dashscope_asr_model", settings.S2S_STT_MODEL,
|
|
"--qwen3_cloud_model", settings.S2S_TTS_MODEL,
|
|
"--qwen3_cloud_voice", settings.S2S_TTS_VOICE,
|
|
"--enable_live_transcription", "false",
|
|
"--thresh", "0.6",
|
|
"--min_silence_ms", "500",
|
|
"--min_speech_ms", "500",
|
|
"--log_level", settings.S2S_LOG_LEVEL,
|
|
]
|
|
args = parse_arguments()
|
|
prepare_all_args(
|
|
args.module_kwargs,
|
|
args.whisper_stt_handler_kwargs,
|
|
args.paraformer_stt_handler_kwargs,
|
|
args.faster_whisper_stt_handler_kwargs,
|
|
args.mlx_audio_whisper_stt_handler_kwargs,
|
|
args.parakeet_tdt_stt_handler_kwargs,
|
|
args.dashscope_asr_stt_handler_kwargs,
|
|
args.language_model_handler_kwargs,
|
|
args.responses_api_language_model_handler_kwargs,
|
|
args.chat_tts_handler_kwargs,
|
|
args.facebook_mms_tts_handler_kwargs,
|
|
args.pocket_tts_handler_kwargs,
|
|
args.kokoro_tts_handler_kwargs,
|
|
args.qwen3_tts_handler_kwargs,
|
|
args.qwen3_cloud_tts_handler_kwargs,
|
|
)
|
|
|
|
pool = [
|
|
_build_realtime_pipeline_unit(
|
|
index=i,
|
|
stop_event=stop_event,
|
|
module_kwargs=args.module_kwargs,
|
|
vad_handler_kwargs=args.vad_handler_kwargs,
|
|
whisper_stt_handler_kwargs=args.whisper_stt_handler_kwargs,
|
|
faster_whisper_stt_handler_kwargs=args.faster_whisper_stt_handler_kwargs,
|
|
paraformer_stt_handler_kwargs=args.paraformer_stt_handler_kwargs,
|
|
mlx_audio_whisper_stt_handler_kwargs=args.mlx_audio_whisper_stt_handler_kwargs,
|
|
parakeet_tdt_stt_handler_kwargs=args.parakeet_tdt_stt_handler_kwargs,
|
|
dashscope_asr_stt_handler_kwargs=args.dashscope_asr_stt_handler_kwargs,
|
|
language_model_handler_kwargs=args.language_model_handler_kwargs,
|
|
responses_api_language_model_handler_kwargs=args.responses_api_language_model_handler_kwargs,
|
|
chat_tts_handler_kwargs=args.chat_tts_handler_kwargs,
|
|
facebook_mms_tts_handler_kwargs=args.facebook_mms_tts_handler_kwargs,
|
|
pocket_tts_handler_kwargs=args.pocket_tts_handler_kwargs,
|
|
kokoro_tts_handler_kwargs=args.kokoro_tts_handler_kwargs,
|
|
qwen3_tts_handler_kwargs=args.qwen3_tts_handler_kwargs,
|
|
qwen3_cloud_tts_handler_kwargs=args.qwen3_cloud_tts_handler_kwargs,
|
|
)
|
|
for i in range(settings.S2S_NUM_PIPELINES)
|
|
]
|
|
|
|
app = create_app(pool=pool, stop_event=stop_event)
|
|
return app, pool
|
|
|
|
|
|
def start_s2s_backend(stop_event: threading.Event) -> None:
|
|
"""后台线程构建管线并启动 RealtimeServer(阻塞调用,放线程里跑)。"""
|
|
|
|
def _run() -> None:
|
|
from speech_to_speech.api.openai_realtime.server import RealtimeServer
|
|
|
|
# LLM warmup / 云 ASR/TTS 连接在弱网下会瞬断:整个池构建重试几次
|
|
last_exc: Exception | None = None
|
|
log.info(
|
|
"s2s: 启动配置 STT=%s LLM=%s TTS=%s(%s) pipelines=%d ws://%s:%d/v1/realtime retries=%d",
|
|
settings.S2S_STT_MODEL, settings.S2S_LLM_MODEL, settings.S2S_TTS_MODEL,
|
|
settings.S2S_TTS_VOICE, settings.S2S_NUM_PIPELINES,
|
|
settings.S2S_HOST, settings.S2S_PORT, settings.S2S_BUILD_RETRIES,
|
|
)
|
|
for attempt in range(settings.S2S_BUILD_RETRIES):
|
|
try:
|
|
log.info(
|
|
"s2s: 构建语音管线池(%d 路,尝试 %d/%d,首次约 30~60s)…",
|
|
settings.S2S_NUM_PIPELINES,
|
|
attempt + 1,
|
|
settings.S2S_BUILD_RETRIES,
|
|
)
|
|
app, pool = _build_pool(stop_event)
|
|
server = RealtimeServer(
|
|
stop_event=stop_event,
|
|
pool=pool,
|
|
host=settings.S2S_HOST,
|
|
port=settings.S2S_PORT,
|
|
)
|
|
# 关键:与 s2s_pipeline.main() 的 realtime 分支一致——用 ThreadManager
|
|
# 启动 RealtimeServer + 全部 handler 线程(VAD/STT/LLM/TTS)。
|
|
# 只跑 server 不启动 handler 线程会导致管线不处理任何音频。
|
|
from speech_to_speech.utils.thread_manager import ThreadManager
|
|
|
|
all_handlers = [server] + [h for unit in pool for h in unit.handlers]
|
|
thread_manager = ThreadManager(all_handlers)
|
|
thread_manager.start()
|
|
log.info("s2s: 实时语音服务 ws://%s:%d/v1/realtime", settings.S2S_HOST, settings.S2S_PORT)
|
|
thread_manager.wait()
|
|
return
|
|
except Exception as e: # noqa: BLE001
|
|
last_exc = e
|
|
log.warning("s2s: 构建尝试 %d/%d 失败:%s", attempt + 1, settings.S2S_BUILD_RETRIES, e)
|
|
if attempt + 1 < settings.S2S_BUILD_RETRIES:
|
|
time.sleep(settings.S2S_BUILD_RETRY_DELAY)
|
|
log.error("s2s: 启动失败(DPM_S2S_ENABLED=0 可关闭):%s", last_exc)
|
|
|
|
threading.Thread(target=_run, daemon=True).start()
|