From 78bd71ff89a93a6dc4c24bfadebae166aa99fdfc Mon Sep 17 00:00:00 2001 From: Bifang <915779419@qq.com> Date: Wed, 23 Sep 2026 10:30:16 +0800 Subject: [PATCH] Split FunASR frontend and backend services --- .env.example | 72 +- .env.funasr.example | 31 + .env.qwen_legacy.example | 50 ++ FUNASR_README.md | 53 ++ README.md | 160 +---- README_QWEN_LEGACY.md | 151 ++++ pyproject.toml | 18 +- pyproject_qwen_legacy.toml | 20 + .../{FIXES.md => FIXES_QWEN_LEGACY.md} | 0 .../{README.md => README_QWEN_LEGACY.md} | 0 realtime_websocket/__init__.py | 1 + realtime_websocket/funasr_engine.py | 393 +++++++++++ realtime_websocket/funasr_server.py | 642 ++++++++++++++++++ ...ervice.py => model_service_qwen_legacy.py} | 0 realtime_websocket/requirements.txt | 4 + .../{server.py => server_qwen_legacy.py} | 0 realtime_websocket/static/app.js | 14 +- realtime_websocket/static/index.html | 13 +- ...e.py => qwen_legacy_model_service_test.py} | 0 ...t_server.py => qwen_legacy_server_test.py} | 0 .../tests/test_funasr_engine.py | 75 ++ .../tests/test_funasr_server.py | 112 +++ .../tests/test_speaker_assembler.py | 5 +- requirements-auxiliary.txt | 2 +- requirements-funasr.txt | 8 + ...deploy.txt => requirements-qwen-legacy.txt | 0 scripts/auxiliary_server.py | 43 +- ...dels.py => download_models_qwen_legacy.py} | 0 scripts/run_backend.py | 162 +++++ scripts/run_frontend.py | 57 ++ scripts/run_funasr_demo.py | 17 + scripts/{serve.py => serve_qwen_legacy.py} | 0 ....py => qwen_legacy_model_manifest_test.py} | 0 ...est_serve.py => qwen_legacy_serve_test.py} | 0 34 files changed, 1870 insertions(+), 233 deletions(-) create mode 100644 .env.funasr.example create mode 100644 .env.qwen_legacy.example create mode 100644 FUNASR_README.md create mode 100644 README_QWEN_LEGACY.md create mode 100644 pyproject_qwen_legacy.toml rename realtime_websocket/{FIXES.md => FIXES_QWEN_LEGACY.md} (100%) rename realtime_websocket/{README.md => README_QWEN_LEGACY.md} (100%) create mode 100644 realtime_websocket/__init__.py create mode 100644 realtime_websocket/funasr_engine.py create mode 100644 realtime_websocket/funasr_server.py rename realtime_websocket/{model_service.py => model_service_qwen_legacy.py} (100%) rename realtime_websocket/{server.py => server_qwen_legacy.py} (100%) rename realtime_websocket/tests/{test_model_service.py => qwen_legacy_model_service_test.py} (100%) rename realtime_websocket/tests/{test_server.py => qwen_legacy_server_test.py} (100%) create mode 100644 realtime_websocket/tests/test_funasr_engine.py create mode 100644 realtime_websocket/tests/test_funasr_server.py create mode 100644 requirements-funasr.txt rename requirements-deploy.txt => requirements-qwen-legacy.txt (100%) rename scripts/{download_models.py => download_models_qwen_legacy.py} (100%) create mode 100644 scripts/run_backend.py create mode 100644 scripts/run_frontend.py create mode 100644 scripts/run_funasr_demo.py rename scripts/{serve.py => serve_qwen_legacy.py} (100%) rename tests/{test_model_manifest.py => qwen_legacy_model_manifest_test.py} (100%) rename tests/{test_serve.py => qwen_legacy_serve_test.py} (100%) diff --git a/.env.example b/.env.example index 5901c32..287172a 100644 --- a/.env.example +++ b/.env.example @@ -1,50 +1,28 @@ -# 选择已经下载到 demo/models 目录中的 ASR 模型;必须与实际模型目录保持一致。 -# 默认使用轻量的 0.6B 模型。如下载的是 1.7B,请改为 1.7b,不能混用。 -QWEN3_ASR_MODEL=0.6b -# 本地模型根目录。下载脚本、VLLM 启动脚本和辅助服务都会从这里查找模型。 -MODEL_DIR=D:/github-project/ASR/Qwen-Asr/demo/models +# Local model assets; the backend never downloads models at startup. +MODEL_DIR=models +# Set this if CAM++ is not at a model_manifest.json path. +# CAM_MODEL_PATH=iic/speech_campplus_sv_zh-cn_16k-common -# VLLM 宿主机服务绑定地址;0.0.0.0 表示允许服务器网卡接收外部请求。 -VLLM_HOST=0.0.0.0 -# VLLM 服务端口;启动器、健康检查和前端连接地址必须使用同一个端口。 -VLLM_PORT=9950 -# VLLM 可执行文件名称;新版环境通常为 vllm,启动器会自动执行 vllm serve。 -VLLM_EXECUTABLE=vllm -# 启动成功提示中显示的地址,只影响日志和使用说明,不改变实际监听地址。 -VLLM_DISPLAY_HOST=127.0.0.1 -# 启动轮询使用的地址;如果 VLLM 部署在本机,通常保持 127.0.0.1 即可。 -VLLM_PROBE_HOST=127.0.0.1 -# VLLM 启动阶段最多检查多少次健康状态,超过次数仍未就绪则启动失败。 -VLLM_STARTUP_CHECK_LOOPS=60 -# 两次健康检查之间的等待秒数;模型加载较慢时可以适当增大检查次数。 -VLLM_STARTUP_CHECK_INTERVAL_SECONDS=2 -# VLLM 使用的显存比例。显存还要留给 VAD、聚类和声纹辅助模型时,建议预留余量。 -VLLM_GPU_MEMORY_UTILIZATION=0.3 -# VLLM 的最大上下文长度;数值越大通常占用越多显存,请结合显卡容量调整。 -VLLM_MAX_MODEL_LEN=16384 -# VLLM 同时处理的最大序列数;实时单路验证可保持默认值,多路并发时再调大。 -VLLM_MAX_NUM_SEQS=16 -# 张量并行 GPU 数量;单卡部署为 1,多卡部署时填写参与并行的 GPU 数量。 -VLLM_TENSOR_PARALLEL_SIZE=1 -# 是否启用 eager 模式。true 通常更容易启动和排查,false 可能获得更高性能。 -VLLM_ENFORCE_EAGER=true +# FunASR realtime engine +FUNASR_ASR_MODEL=paraformer-zh-streaming +FUNASR_VAD_MODEL=fsmn-vad +FUNASR_DEVICE=cuda:0 +FUNASR_VAD_DEVICE=cpu +FUNASR_CHUNK_SIZE=0,10,5 +FUNASR_ENCODER_LOOK_BACK=4 +FUNASR_DECODER_LOOK_BACK=1 +FUNASR_VAD_CHUNK_MS=200 +FUNASR_MAX_SEGMENT_SEC=30 -# GB10 需要使用 CUDA 13 工具链中的 ptxas;serve.py 会自动加载该配置并传给 VLLM。 -TRITON_PTXAS_PATH=/usr/local/cuda/bin/ptxas - -# 可选:为外部客户端设置稳定的公开模型名称。留空时默认使用完整 ModelScope ID。 -# VLLM_SERVED_MODEL_NAME=Qwen/Qwen3-ASR-0.6B - -# 辅助模型服务地址;它独立加载 VAD、CAM++ 聚类和声纹模型,WebSocket -# 服务只通过 HTTP 调用,不会把这些 GPU 模型重复加载到 WebSocket 进程。 +# CAM++ is required and started by the backend launcher. AUXILIARY_SERVICE_URL=http://127.0.0.1:8010 -# WebSocket 调用已独立部署的 vLLM;跨服务器时改成模型服务器的实际地址。 -MODEL_SERVICE_URL=http://127.0.0.1:9950/v1 -# 辅助服务监听的 GPU;单卡服务器保持 cuda:0,多卡时可改成指定卡号。 -AUXILIARY_DEVICE=cuda:0 -# 仅保留兼容旧配置;VAD/CAM++ 核心模型缺失时始终拒绝启动,不会静默降级。 -# 可选 diarization/aligner 缺失不会阻断实时服务。 -AUXILIARY_ALLOW_MISSING=false -# 启动时强制加载 VAD + CAM++ speaker_verification;完整 diarization 首次调用时按需加载。 -# 如确实需要启动时额外预加载,可追加:vad,speaker_verification,diarization -AUXILIARY_PRELOAD_KINDS=vad,speaker_verification + +# Standalone frontend and browser-facing backend origin. +FRONTEND_HOST=127.0.0.1 +FRONTEND_PORT=8080 +FRONTEND_ORIGIN=http://127.0.0.1:8080 +BACKEND_PUBLIC_URL=http://127.0.0.1:8082 + +WEB_HOST=0.0.0.0 +WEB_PORT=8082 +WEB_DISPLAY_HOST=127.0.0.1 diff --git a/.env.funasr.example b/.env.funasr.example new file mode 100644 index 0000000..29b643f --- /dev/null +++ b/.env.funasr.example @@ -0,0 +1,31 @@ +# Local model assets; the backend never downloads models at startup. +MODEL_DIR=models +# Set this if CAM++ is not at a model_manifest.json path. +# CAM_MODEL_PATH=iic/speech_campplus_sv_zh-cn_16k-common + +# FunASR streaming ASR model. The backend launcher requires it to exist locally. +FUNASR_ASR_MODEL=paraformer-zh-streaming +FUNASR_VAD_MODEL=fsmn-vad + +# ASR and VAD may use different devices to limit peak GPU memory. +FUNASR_DEVICE=cuda:0 +FUNASR_VAD_DEVICE=cpu + +# [left context, current chunk, right lookahead]; 10 * 60ms = 600ms. +FUNASR_CHUNK_SIZE=0,10,5 +FUNASR_ENCODER_LOOK_BACK=4 +FUNASR_DECODER_LOOK_BACK=1 +FUNASR_VAD_CHUNK_MS=200 +FUNASR_MAX_SEGMENT_SEC=30 + +# CAM++ is required. The backend launcher starts this service itself. +AUXILIARY_SERVICE_URL=http://127.0.0.1:8010 + +# Standalone frontend and browser-facing backend origin. +FRONTEND_HOST=127.0.0.1 +FRONTEND_PORT=8080 +FRONTEND_ORIGIN=http://127.0.0.1:8080 +BACKEND_PUBLIC_URL=http://127.0.0.1:8082 +WEB_HOST=0.0.0.0 +WEB_PORT=8082 +WEB_DISPLAY_HOST=127.0.0.1 diff --git a/.env.qwen_legacy.example b/.env.qwen_legacy.example new file mode 100644 index 0000000..5901c32 --- /dev/null +++ b/.env.qwen_legacy.example @@ -0,0 +1,50 @@ +# 选择已经下载到 demo/models 目录中的 ASR 模型;必须与实际模型目录保持一致。 +# 默认使用轻量的 0.6B 模型。如下载的是 1.7B,请改为 1.7b,不能混用。 +QWEN3_ASR_MODEL=0.6b +# 本地模型根目录。下载脚本、VLLM 启动脚本和辅助服务都会从这里查找模型。 +MODEL_DIR=D:/github-project/ASR/Qwen-Asr/demo/models + +# VLLM 宿主机服务绑定地址;0.0.0.0 表示允许服务器网卡接收外部请求。 +VLLM_HOST=0.0.0.0 +# VLLM 服务端口;启动器、健康检查和前端连接地址必须使用同一个端口。 +VLLM_PORT=9950 +# VLLM 可执行文件名称;新版环境通常为 vllm,启动器会自动执行 vllm serve。 +VLLM_EXECUTABLE=vllm +# 启动成功提示中显示的地址,只影响日志和使用说明,不改变实际监听地址。 +VLLM_DISPLAY_HOST=127.0.0.1 +# 启动轮询使用的地址;如果 VLLM 部署在本机,通常保持 127.0.0.1 即可。 +VLLM_PROBE_HOST=127.0.0.1 +# VLLM 启动阶段最多检查多少次健康状态,超过次数仍未就绪则启动失败。 +VLLM_STARTUP_CHECK_LOOPS=60 +# 两次健康检查之间的等待秒数;模型加载较慢时可以适当增大检查次数。 +VLLM_STARTUP_CHECK_INTERVAL_SECONDS=2 +# VLLM 使用的显存比例。显存还要留给 VAD、聚类和声纹辅助模型时,建议预留余量。 +VLLM_GPU_MEMORY_UTILIZATION=0.3 +# VLLM 的最大上下文长度;数值越大通常占用越多显存,请结合显卡容量调整。 +VLLM_MAX_MODEL_LEN=16384 +# VLLM 同时处理的最大序列数;实时单路验证可保持默认值,多路并发时再调大。 +VLLM_MAX_NUM_SEQS=16 +# 张量并行 GPU 数量;单卡部署为 1,多卡部署时填写参与并行的 GPU 数量。 +VLLM_TENSOR_PARALLEL_SIZE=1 +# 是否启用 eager 模式。true 通常更容易启动和排查,false 可能获得更高性能。 +VLLM_ENFORCE_EAGER=true + +# GB10 需要使用 CUDA 13 工具链中的 ptxas;serve.py 会自动加载该配置并传给 VLLM。 +TRITON_PTXAS_PATH=/usr/local/cuda/bin/ptxas + +# 可选:为外部客户端设置稳定的公开模型名称。留空时默认使用完整 ModelScope ID。 +# VLLM_SERVED_MODEL_NAME=Qwen/Qwen3-ASR-0.6B + +# 辅助模型服务地址;它独立加载 VAD、CAM++ 聚类和声纹模型,WebSocket +# 服务只通过 HTTP 调用,不会把这些 GPU 模型重复加载到 WebSocket 进程。 +AUXILIARY_SERVICE_URL=http://127.0.0.1:8010 +# WebSocket 调用已独立部署的 vLLM;跨服务器时改成模型服务器的实际地址。 +MODEL_SERVICE_URL=http://127.0.0.1:9950/v1 +# 辅助服务监听的 GPU;单卡服务器保持 cuda:0,多卡时可改成指定卡号。 +AUXILIARY_DEVICE=cuda:0 +# 仅保留兼容旧配置;VAD/CAM++ 核心模型缺失时始终拒绝启动,不会静默降级。 +# 可选 diarization/aligner 缺失不会阻断实时服务。 +AUXILIARY_ALLOW_MISSING=false +# 启动时强制加载 VAD + CAM++ speaker_verification;完整 diarization 首次调用时按需加载。 +# 如确实需要启动时额外预加载,可追加:vad,speaker_verification,diarization +AUXILIARY_PRELOAD_KINDS=vad,speaker_verification diff --git a/FUNASR_README.md b/FUNASR_README.md new file mode 100644 index 0000000..115f6de --- /dev/null +++ b/FUNASR_README.md @@ -0,0 +1,53 @@ +# FunASR realtime browser demo + +The frontend and backend start separately. The backend launcher loads three +required local assets: streaming ASR, FSMN VAD, and CAM++ speaker verification. +It starts the CAM++ model service and the WebSocket service and stops both +together. No model is downloaded by the launcher. + +## Model directories + +Put the assets under models/ or set MODEL_DIR in .env. The default +FunASR names resolve to these local directories: + +- models/iic/speech_paraformer-large_asr_nat-zh-cn-16k-common-vocab8404-online +- models/iic/speech_fsmn_vad_zh-cn-16k-common-pytorch (or the damo/ variant) +- models/iic/speech_campplus_sv_zh-cn_16k-common (or the damo/ variant) + +If your directories have different names, set FUNASR_ASR_MODEL and +FUNASR_VAD_MODEL to their paths. The CAM++ directories follow +model_manifest.json; set CAM_MODEL_PATH for another location. Backend startup reports every checked path when a +required asset is missing. + +## Start + +First install a torch/torchaudio build suitable for the host CPU or CUDA, +then install the project dependencies: + +~~~powershell +cd D:\github-project\ASR\Asr-demo +python -m pip install -r requirements-funasr.txt +python -m pip install -r requirements-auxiliary.txt +if (-not (Test-Path .env)) { Copy-Item .env.funasr.example .env } +~~~ + +In one terminal start the backend: + +~~~powershell +python scripts\run_backend.py +~~~ + +In another terminal start the frontend: + +~~~powershell +python scripts\run_frontend.py +~~~ + +Open http://127.0.0.1:8080/. The backend WebSocket listens on port 8082 and +the CAM++ service on port 8010. BACKEND_PUBLIC_URL must be reachable from +the browser. Set FRONTEND_ORIGIN to the exact frontend origin if it differs +from the default. + +The browser sends start, PCM16 frames, and stop/eof over WebSocket. Speaker +labels are required. Short or unusable speech may still receive an unknown +speaker label, but missing CAM++ prevents backend startup. diff --git a/README.md b/README.md index 0140208..d1e2c28 100644 --- a/README.md +++ b/README.md @@ -1,151 +1,25 @@ -# Qwen3-ASR VLLM 独立部署项目 +# FunASR realtime ASR demo -本目录是后续实时 ASR 功能验证使用的独立模型服务项目。 +This branch runs streaming ASR, VAD, and CAM++ from local model assets. +The browser UI starts as a separate service. -它不导入、不启动、也不调用仓库根目录下原项目的 `app/` 代码。模型下载、VLLM 启动、配置和服务验证都在本目录内完成。后续验证 demo 只需要调用这里提供的 VLLM OpenAI 兼容接口。 +## Start -## 当前下载范围 +Install a torch/torchaudio build for the host, then install +requirements-funasr.txt and requirements-auxiliary.txt. Copy +.env.funasr.example to .env if no .env exists, and check MODEL_DIR. -默认下载一个 ASR 模型和独立辅助模型运行服务所需的全部模型资产: +Start the backend (CAM++ model service plus WebSocket) in one terminal: -```text -Qwen/Qwen3-ASR-0.6B -``` +~~~powershell +python scripts\run_backend.py +~~~ -ASR 如需使用大模型,可显式选择 1.7B: +Start the frontend in another terminal: -```text -Qwen/Qwen3-ASR-1.7B -``` +~~~powershell +python scripts\run_frontend.py +~~~ -辅助资产包括 VAD、CAM++ 分离、配置声纹、实时声纹、CAM++ Transformer 和 Qwen3 ForcedAligner。它们不会由 `qwen-asr-serve` 启动,而是由独立 Python 辅助模型服务加载。 - -不会下载另一个未选择的 ASR 模型;辅助模型资产会随默认部署包下载,供独立 Python 运行服务预加载。 - -## 1. 下载模型 - -模型直接下载到宿主机的 `demo/models`。ASR 由 VLLM 启动,辅助模型由独立 Python 运行服务启动。 - -```powershell -cd D:\github-project\ASR\Qwen-Asr\demo -python -m venv .venv -.\.venv\Scripts\Activate.ps1 -pip install -r requirements-download.txt -python scripts\download_models.py -``` - -上述命令会下载默认 `0.6B` ASR 以及全部辅助模型,并在下载完成后把 CAM++ 配置中的依赖模型 ID 改为 `demo/models` 下的本地路径,保证辅助服务可以离线启动。只下载 ASR 时使用: - -```powershell -python scripts\download_models.py --skip-auxiliary -``` - -只下载辅助模型时使用: - -```powershell -python scripts\download_models.py --auxiliary-only -``` - -选择 1.7B: - -```powershell -python scripts\download_models.py --model 1.7b -``` - -检查模型是否完整但不下载: - -```powershell -python scripts\download_models.py --check-only -``` - -ModelScope 下载也可以通过环境变量调整缓存目录: - -```powershell -$env:MODELSCOPE_CACHE = 'D:\modelscope-cache' -python scripts\download_models.py -``` - -## 2. 宿主机启动 VLLM 服务 - -需要宿主机具备与 VLLM 兼容的 Python、CUDA 和 NVIDIA 驱动环境。安装部署依赖: - -```bash -python -m pip install -r requirements-deploy.txt -``` - -先复制并按服务器实际路径修改 `.env`,启动器会自动读取该文件: - -```bash -cp .env.example .env -``` - -默认启动 `Qwen/Qwen3-ASR-0.6B`,监听地址为 `0.0.0.0:9950`: - -```bash -python scripts/serve.py -``` - -模型和启动检查循环可以通过命令行或环境变量传入;端口统一在 `scripts/serve.py` 的 `SERVER_PORT` 变量中维护: - -```bash -QWEN3_ASR_MODEL=0.6b VLLM_STARTUP_CHECK_LOOPS=120 \ -VLLM_STARTUP_CHECK_INTERVAL_SECONDS=2 python scripts/serve.py -``` - -如果使用 `1.7B`,下载和启动必须指定同一个模型: - -```powershell -python scripts\download_models.py --model 1.7b -python -m scripts.serve --model 1.7b -``` - -启动器会按 `VLLM_STARTUP_CHECK_LOOPS` 次数轮询 `/health`,每次间隔由 `VLLM_STARTUP_CHECK_INTERVAL_SECONDS` 指定。服务端口、健康检查端口和就绪提示统一使用 `scripts/serve.py` 中的 `SERVER_PORT`。 - -服务启动后可检查: - -```bash -curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/health" -curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/v1/models" -``` - -## 3. 调用转写接口 - -VLLM 服务提供 OpenAI 兼容的音频转写接口: - -```bash -curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/v1/audio/transcriptions" \ - -H "Authorization: Bearer EMPTY" \ - -F "file=@./audio/sample.wav" \ - -F "model=Qwen/Qwen3-ASR-0.6B" -``` - -## 4. 启动辅助模型服务 - -另开一个终端,在同一台服务器启动 VAD、CAM++ 和声纹模型运行服务: - -```powershell -cd D:\github-project\ASR\Qwen-Asr\demo -pip install -r requirements-auxiliary.txt -python scripts\auxiliary_server.py -``` - -辅助服务默认监听 `0.0.0.0:8010`。实时链路启动时严格加载 VAD 和 CAM++ `speaker_verification` 声纹模型,用于每个 turn 的特征提取与在线聚类;完整 CAM++ 分离、Transformer 和 ForcedAligner 不阻断核心服务,完整分离模型会在调用 `/v1/diarization` 时按需加载。检查状态: - -```bash -curl http://127.0.0.1:8010/health -``` - -WebSocket demo 默认连接 `9950` 的 ASR VLLM,辅助服务使用 `8010`。`/health` 的 `ready` 要求 `vad_ready` 与 `speaker_embedding_ready` 同时为 true;完整 diarization 资产缺失不会影响实时 `/v1/speaker/resolve`。 - -实时链路中的职责是:WebSocket 用 RMS 帧门控快速检测停顿;辅助服务用 FunASR VAD 提供 `/v1/vad`,并加载 CAM++ `speaker_verification` 提取 turn embedding,再由服务端在线聚类。`speech_campplus_speaker-diarization_common` 是完整音频分离接口的额外 pipeline,不是实时 turn 聚类的唯一入口。 - -也可以使用多模态 Chat Completions 接口,后续实时验证项目将以此服务边界为准。 - -## 项目边界 - -- `scripts/download_models.py`:下载选定 ASR 和全部辅助模型资产。 -- `scripts/serve.py`:读取 `.env`,解析模型选择、宿主机参数并启动新版 `vllm serve`。 -- `requirements-deploy.txt`:安装宿主机部署所需的官方 Qwen3-ASR VLLM 依赖。 -- `tests/`:只验证本项目自己的模型清单和选择逻辑,不依赖原项目。 - -模型服务就绪后,新的实时 ASR demo 放在同级 `demo` 项目中继续开发,但不得通过 Python import 或 HTTP/WebSocket 调用原项目服务。 +Open http://127.0.0.1:8080/. Model directory layout and address settings +are described in FUNASR_README.md. diff --git a/README_QWEN_LEGACY.md b/README_QWEN_LEGACY.md new file mode 100644 index 0000000..0140208 --- /dev/null +++ b/README_QWEN_LEGACY.md @@ -0,0 +1,151 @@ +# Qwen3-ASR VLLM 独立部署项目 + +本目录是后续实时 ASR 功能验证使用的独立模型服务项目。 + +它不导入、不启动、也不调用仓库根目录下原项目的 `app/` 代码。模型下载、VLLM 启动、配置和服务验证都在本目录内完成。后续验证 demo 只需要调用这里提供的 VLLM OpenAI 兼容接口。 + +## 当前下载范围 + +默认下载一个 ASR 模型和独立辅助模型运行服务所需的全部模型资产: + +```text +Qwen/Qwen3-ASR-0.6B +``` + +ASR 如需使用大模型,可显式选择 1.7B: + +```text +Qwen/Qwen3-ASR-1.7B +``` + +辅助资产包括 VAD、CAM++ 分离、配置声纹、实时声纹、CAM++ Transformer 和 Qwen3 ForcedAligner。它们不会由 `qwen-asr-serve` 启动,而是由独立 Python 辅助模型服务加载。 + +不会下载另一个未选择的 ASR 模型;辅助模型资产会随默认部署包下载,供独立 Python 运行服务预加载。 + +## 1. 下载模型 + +模型直接下载到宿主机的 `demo/models`。ASR 由 VLLM 启动,辅助模型由独立 Python 运行服务启动。 + +```powershell +cd D:\github-project\ASR\Qwen-Asr\demo +python -m venv .venv +.\.venv\Scripts\Activate.ps1 +pip install -r requirements-download.txt +python scripts\download_models.py +``` + +上述命令会下载默认 `0.6B` ASR 以及全部辅助模型,并在下载完成后把 CAM++ 配置中的依赖模型 ID 改为 `demo/models` 下的本地路径,保证辅助服务可以离线启动。只下载 ASR 时使用: + +```powershell +python scripts\download_models.py --skip-auxiliary +``` + +只下载辅助模型时使用: + +```powershell +python scripts\download_models.py --auxiliary-only +``` + +选择 1.7B: + +```powershell +python scripts\download_models.py --model 1.7b +``` + +检查模型是否完整但不下载: + +```powershell +python scripts\download_models.py --check-only +``` + +ModelScope 下载也可以通过环境变量调整缓存目录: + +```powershell +$env:MODELSCOPE_CACHE = 'D:\modelscope-cache' +python scripts\download_models.py +``` + +## 2. 宿主机启动 VLLM 服务 + +需要宿主机具备与 VLLM 兼容的 Python、CUDA 和 NVIDIA 驱动环境。安装部署依赖: + +```bash +python -m pip install -r requirements-deploy.txt +``` + +先复制并按服务器实际路径修改 `.env`,启动器会自动读取该文件: + +```bash +cp .env.example .env +``` + +默认启动 `Qwen/Qwen3-ASR-0.6B`,监听地址为 `0.0.0.0:9950`: + +```bash +python scripts/serve.py +``` + +模型和启动检查循环可以通过命令行或环境变量传入;端口统一在 `scripts/serve.py` 的 `SERVER_PORT` 变量中维护: + +```bash +QWEN3_ASR_MODEL=0.6b VLLM_STARTUP_CHECK_LOOPS=120 \ +VLLM_STARTUP_CHECK_INTERVAL_SECONDS=2 python scripts/serve.py +``` + +如果使用 `1.7B`,下载和启动必须指定同一个模型: + +```powershell +python scripts\download_models.py --model 1.7b +python -m scripts.serve --model 1.7b +``` + +启动器会按 `VLLM_STARTUP_CHECK_LOOPS` 次数轮询 `/health`,每次间隔由 `VLLM_STARTUP_CHECK_INTERVAL_SECONDS` 指定。服务端口、健康检查端口和就绪提示统一使用 `scripts/serve.py` 中的 `SERVER_PORT`。 + +服务启动后可检查: + +```bash +curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/health" +curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/v1/models" +``` + +## 3. 调用转写接口 + +VLLM 服务提供 OpenAI 兼容的音频转写接口: + +```bash +curl "http://${VLLM_DISPLAY_HOST:-127.0.0.1}:9950/v1/audio/transcriptions" \ + -H "Authorization: Bearer EMPTY" \ + -F "file=@./audio/sample.wav" \ + -F "model=Qwen/Qwen3-ASR-0.6B" +``` + +## 4. 启动辅助模型服务 + +另开一个终端,在同一台服务器启动 VAD、CAM++ 和声纹模型运行服务: + +```powershell +cd D:\github-project\ASR\Qwen-Asr\demo +pip install -r requirements-auxiliary.txt +python scripts\auxiliary_server.py +``` + +辅助服务默认监听 `0.0.0.0:8010`。实时链路启动时严格加载 VAD 和 CAM++ `speaker_verification` 声纹模型,用于每个 turn 的特征提取与在线聚类;完整 CAM++ 分离、Transformer 和 ForcedAligner 不阻断核心服务,完整分离模型会在调用 `/v1/diarization` 时按需加载。检查状态: + +```bash +curl http://127.0.0.1:8010/health +``` + +WebSocket demo 默认连接 `9950` 的 ASR VLLM,辅助服务使用 `8010`。`/health` 的 `ready` 要求 `vad_ready` 与 `speaker_embedding_ready` 同时为 true;完整 diarization 资产缺失不会影响实时 `/v1/speaker/resolve`。 + +实时链路中的职责是:WebSocket 用 RMS 帧门控快速检测停顿;辅助服务用 FunASR VAD 提供 `/v1/vad`,并加载 CAM++ `speaker_verification` 提取 turn embedding,再由服务端在线聚类。`speech_campplus_speaker-diarization_common` 是完整音频分离接口的额外 pipeline,不是实时 turn 聚类的唯一入口。 + +也可以使用多模态 Chat Completions 接口,后续实时验证项目将以此服务边界为准。 + +## 项目边界 + +- `scripts/download_models.py`:下载选定 ASR 和全部辅助模型资产。 +- `scripts/serve.py`:读取 `.env`,解析模型选择、宿主机参数并启动新版 `vllm serve`。 +- `requirements-deploy.txt`:安装宿主机部署所需的官方 Qwen3-ASR VLLM 依赖。 +- `tests/`:只验证本项目自己的模型清单和选择逻辑,不依赖原项目。 + +模型服务就绪后,新的实时 ASR demo 放在同级 `demo` 项目中继续开发,但不得通过 Python import 或 HTTP/WebSocket 调用原项目服务。 diff --git a/pyproject.toml b/pyproject.toml index 30ed923..f0c63ad 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -3,18 +3,22 @@ requires = ["setuptools>=68", "wheel"] build-backend = "setuptools.build_meta" [project] -name = "qwen3-asr-vllm-deployment" +name = "funasr-realtime-asr-demo" version = "0.1.0" -description = "Standalone Qwen3-ASR model downloader and VLLM deployment" +description = "FunASR streaming ASR browser demo" requires-python = ">=3.10,<3.14" dependencies = [ - "modelscope==1.34.0", - "qwen-asr[vllm]==0.0.6", + "aiohttp==3.11.11", + "python-dotenv>=1.0", + "funasr==1.4.16", + "modelscope[framework]==1.34.0", + "soundfile==0.13.1", + "librosa==0.11.0", + "numpy>=1.24", ] [project.scripts] -qwen3-asr-download = "scripts.download_models:main" -qwen3-asr-serve = "scripts.serve:main" +funasr-realtime-demo = "scripts.run_funasr_demo:main" [tool.setuptools] -packages = ["scripts"] +packages = ["scripts", "realtime_websocket"] diff --git a/pyproject_qwen_legacy.toml b/pyproject_qwen_legacy.toml new file mode 100644 index 0000000..30ed923 --- /dev/null +++ b/pyproject_qwen_legacy.toml @@ -0,0 +1,20 @@ +[build-system] +requires = ["setuptools>=68", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "qwen3-asr-vllm-deployment" +version = "0.1.0" +description = "Standalone Qwen3-ASR model downloader and VLLM deployment" +requires-python = ">=3.10,<3.14" +dependencies = [ + "modelscope==1.34.0", + "qwen-asr[vllm]==0.0.6", +] + +[project.scripts] +qwen3-asr-download = "scripts.download_models:main" +qwen3-asr-serve = "scripts.serve:main" + +[tool.setuptools] +packages = ["scripts"] diff --git a/realtime_websocket/FIXES.md b/realtime_websocket/FIXES_QWEN_LEGACY.md similarity index 100% rename from realtime_websocket/FIXES.md rename to realtime_websocket/FIXES_QWEN_LEGACY.md diff --git a/realtime_websocket/README.md b/realtime_websocket/README_QWEN_LEGACY.md similarity index 100% rename from realtime_websocket/README.md rename to realtime_websocket/README_QWEN_LEGACY.md diff --git a/realtime_websocket/__init__.py b/realtime_websocket/__init__.py new file mode 100644 index 0000000..c77a329 --- /dev/null +++ b/realtime_websocket/__init__.py @@ -0,0 +1 @@ +"""FunASR realtime browser demo package.""" diff --git a/realtime_websocket/funasr_engine.py b/realtime_websocket/funasr_engine.py new file mode 100644 index 0000000..77abbee --- /dev/null +++ b/realtime_websocket/funasr_engine.py @@ -0,0 +1,393 @@ +"""FunASR realtime engine migrated into the demo project. + +This module keeps the browser-facing project independent from the original +Qwen/VLLM service. It follows FunASR's streaming cache lifecycle: +one cache per ASR session and is_final=True only for the last chunk. +""" + +from __future__ import annotations + +import asyncio +import logging +import os +from dataclasses import dataclass +from typing import Any + +LOGGER = logging.getLogger(__name__) + + +@dataclass(frozen=True) +class FunASRServiceConfig: + """Configuration for the in-process FunASR models.""" + + model: str = "paraformer-zh-streaming" + vad_model: str = "fsmn-vad" + device: str = "cuda:0" + vad_device: str = "cpu" + sample_rate: int = 16000 + vad_chunk_ms: int = 200 + chunk_size: tuple[int, int, int] = (0, 10, 5) + encoder_chunk_look_back: int = 4 + decoder_chunk_look_back: int = 1 + max_segment_sec: float = 30.0 + + @classmethod + def from_env(cls) -> "FunASRServiceConfig": + """Read model selection from environment without changing WS fields.""" + chunk_text = os.getenv("FUNASR_CHUNK_SIZE", "0,10,5") + try: + values = tuple(int(value.strip()) for value in chunk_text.split(",")) + chunk_size = values if len(values) == 3 else cls.chunk_size + except ValueError: + chunk_size = cls.chunk_size + return cls( + model=os.getenv("FUNASR_ASR_MODEL", cls.model), + vad_model=os.getenv("FUNASR_VAD_MODEL", cls.vad_model), + device=os.getenv("FUNASR_DEVICE", cls.device), + vad_device=os.getenv("FUNASR_VAD_DEVICE", cls.vad_device), + vad_chunk_ms=max(50, int(os.getenv("FUNASR_VAD_CHUNK_MS", str(cls.vad_chunk_ms)))), + chunk_size=chunk_size, + encoder_chunk_look_back=max( + 0, int(os.getenv("FUNASR_ENCODER_LOOK_BACK", str(cls.encoder_chunk_look_back))) + ), + decoder_chunk_look_back=max( + 0, int(os.getenv("FUNASR_DECODER_LOOK_BACK", str(cls.decoder_chunk_look_back))) + ), + max_segment_sec=max( + 2.0, float(os.getenv("FUNASR_MAX_SEGMENT_SEC", str(cls.max_segment_sec))) + ), + ) + + +@dataclass(frozen=True) +class FunASRSegment: + """A partial or final event returned by one browser session.""" + + text: str + start_time_ms: float + end_time_ms: float + audio: bytes + voiced_ms: float + is_final: bool + sentence_id: int + reason: str | None = None + + +def _result_text(result: Any) -> str: + """Extract text from the list/dict result shapes used by FunASR.""" + if isinstance(result, list): + return _result_text(result[0]) if result else "" + if isinstance(result, dict): + value = result.get("text") + if value is None: + value = result.get("value") + return str(value or "").strip() + return str(result or "").strip() + + +def _vad_events(result: Any) -> list[tuple[float, float]]: + """Normalize FunASR streaming VAD output to start/end milliseconds.""" + if isinstance(result, list): + result = result[0] if result else {} + if isinstance(result, dict): + result = result.get("value") or result.get("segments") or [] + if not isinstance(result, list): + return [] + events: list[tuple[float, float]] = [] + for item in result: + if not isinstance(item, (list, tuple)) or len(item) < 2: + continue + try: + events.append((float(item[0]), float(item[1]))) + except (TypeError, ValueError): + continue + return events + + +class FunASRModelService: + """Load FunASR once and create isolated streaming state per browser.""" + + native_partial_supported = True + + def __init__(self, config: FunASRServiceConfig | None = None) -> None: + self.config = config or FunASRServiceConfig.from_env() + self.asr_model: Any | None = None + self.vad_model: Any | None = None + # Model objects are shared; each session owns its own cache. Serializing + # calls makes the first validation version predictable on one GPU. + self.inference_lock = asyncio.Lock() + + async def start(self) -> None: + """Load streaming ASR and VAD outside the event loop.""" + try: + from funasr import AutoModel + except ImportError as exc: # pragma: no cover - deployment-only branch + raise RuntimeError( + "FunASR is not installed; run pip install -r requirements-deploy.txt" + ) from exc + + def load_models() -> tuple[Any, Any]: + common = {"disable_pbar": True, "disable_log": True} + asr = AutoModel(model=self.config.model, device=self.config.device, **common) + vad = AutoModel(model=self.config.vad_model, device=self.config.vad_device, **common) + return asr, vad + + self.asr_model, self.vad_model = await asyncio.to_thread(load_models) + LOGGER.info( + "FunASR ready: model=%s vad=%s device=%s vad_device=%s", + self.config.model, + self.config.vad_model, + self.config.device, + self.config.vad_device, + ) + + async def close(self) -> None: + """Release references so CUDA memory can be reclaimed on shutdown.""" + self.asr_model = None + self.vad_model = None + + def create_session(self) -> "FunASRRealtimeSession": + """Create a session with isolated VAD and ASR caches.""" + if self.asr_model is None or self.vad_model is None: + raise RuntimeError("FunASR model service is not started") + return FunASRRealtimeSession(self) + + async def generate_asr(self, audio: Any, status: dict[str, Any]) -> str: + """Run one blocking ASR chunk while preserving its mutable cache.""" + if self.asr_model is None: + raise RuntimeError("FunASR ASR model is not loaded") + + def generate() -> str: + return _result_text(self.asr_model.generate(input=audio, **status)) + + async with self.inference_lock: + return await asyncio.to_thread(generate) + + async def generate_vad( + self, + audio: Any, + status: dict[str, Any], + chunk_ms: int, + ) -> list[tuple[float, float]]: + """Run one streaming VAD chunk and normalize endpoint events.""" + if self.vad_model is None: + raise RuntimeError("FunASR VAD model is not loaded") + + def generate() -> list[tuple[float, float]]: + result = self.vad_model.generate(input=audio, chunk_size=chunk_ms, **status) + return _vad_events(result) + + async with self.inference_lock: + return await asyncio.to_thread(generate) + + +class _StreamingASRTurn: + """One utterance using FunASR's ordered chunk/cache lifecycle.""" + + def __init__(self, service: FunASRModelService) -> None: + self.service = service + self.cache: dict[str, Any] = {} + self.pending = bytearray() + self.cumulative_text = "" + self.last_chunk_text = "" + self.pending_outputs: list[str] = [] + self.chunk_samples = max(1, service.config.chunk_size[1] * 960) + + @staticmethod + def _merge_chunk_text(current: str, chunk: str, previous_chunk: str) -> str: + """Merge chunk text while tolerating wrappers returning cumulative text.""" + if not chunk or chunk == previous_chunk: + return current + if current and chunk.startswith(current): + return chunk + return current + chunk + + async def append(self, pcm_bytes: bytes) -> list[str]: + """Decode complete chunks but hold one chunk for the final flush.""" + if not pcm_bytes: + return [] + self.pending.extend(pcm_bytes) + chunk_bytes = self.chunk_samples * 2 + # Hold one full chunk so the real last chunk receives is_final=True. + while len(self.pending) >= chunk_bytes * 2: + chunk = bytes(self.pending[:chunk_bytes]) + del self.pending[:chunk_bytes] + text = await self.service.generate_asr( + chunk, + { + "cache": self.cache, + "is_final": False, + "chunk_size": self.service.config.chunk_size, + "encoder_chunk_look_back": self.service.config.encoder_chunk_look_back, + "decoder_chunk_look_back": self.service.config.decoder_chunk_look_back, + "batch_size": 1, + }, + ) + self.cumulative_text = self._merge_chunk_text( + self.cumulative_text, text, self.last_chunk_text + ) + self.last_chunk_text = text + if text: + self.pending_outputs.append(self.cumulative_text) + outputs = self.pending_outputs + self.pending_outputs = [] + return outputs + + async def finish(self) -> str: + """Flush the last buffered chunk and return cumulative text.""" + if self.pending: + chunk = bytes(self.pending) + self.pending.clear() + text = await self.service.generate_asr( + chunk, + { + "cache": self.cache, + "is_final": True, + "chunk_size": self.service.config.chunk_size, + "encoder_chunk_look_back": self.service.config.encoder_chunk_look_back, + "decoder_chunk_look_back": self.service.config.decoder_chunk_look_back, + "batch_size": 1, + }, + ) + self.cumulative_text = self._merge_chunk_text( + self.cumulative_text, text, self.last_chunk_text + ) + self.last_chunk_text = text + return self.cumulative_text.strip() + + +class FunASRRealtimeSession: + """FunASR VAD + streaming ASR session used by one browser WebSocket.""" + + def __init__(self, service: FunASRModelService) -> None: + self.service = service + self.sample_rate = service.config.sample_rate + self.vad_chunk_bytes = service.config.vad_chunk_ms * self.sample_rate * 2 // 1000 + self.vad_buffer = bytearray() + self.vad_cache: dict[str, Any] = {} + self.pre_roll = bytearray() + self.pre_roll_max_bytes = 300 * self.sample_rate * 2 // 1000 + self.total_samples = 0 + self.speech_started = False + self.segment_id = 0 + self.segment_start_ms = 0.0 + self.segment_audio = bytearray() + self.segment_voiced_ms = 0.0 + self.asr_turn: _StreamingASRTurn | None = None + + @staticmethod + def _to_float32(pcm_bytes: bytes) -> Any: + """Convert browser PCM16 to the float waveform expected by FunASR.""" + import numpy as np + + return np.frombuffer(pcm_bytes, dtype=np.int16).astype(np.float32) / 32768.0 + + def _start_segment(self, start_ms: float) -> None: + """Open a turn and retain a short pre-roll for initial phonemes.""" + self.speech_started = True + self.segment_start_ms = max( + 0.0, + start_ms - len(self.pre_roll) / (self.sample_rate * 2) * 1000, + ) + self.segment_audio = bytearray(self.pre_roll) + self.segment_voiced_ms = 0.0 + self.asr_turn = _StreamingASRTurn(self.service) + self.pre_roll.clear() + + async def _feed_vad_chunk(self, chunk: bytes) -> list[FunASRSegment]: + """Feed VAD, then stream this audio into the active ASR turn.""" + self.total_samples += len(chunk) // 2 + events = await self.service.generate_vad( + self._to_float32(chunk), + {"cache": self.vad_cache, "is_final": False}, + self.service.config.vad_chunk_ms, + ) + starts = [start for start, _ in events if start >= 0] + ends = [end for _, end in events if end >= 0] + results: list[FunASRSegment] = [] + if not self.speech_started and starts: + self._start_segment(starts[0]) + + if self.speech_started: + self.segment_audio.extend(chunk) + self.segment_voiced_ms += len(chunk) / (self.sample_rate * 2) * 1000 + if self.asr_turn is not None: + for text in await self.asr_turn.append(chunk): + results.append(self._partial(text)) + if self._segment_duration_ms() >= self.service.config.max_segment_sec * 1000: + results.append(await self._finish_segment("max_duration")) + else: + self.pre_roll.extend(chunk) + del self.pre_roll[:-self.pre_roll_max_bytes] + + if ends and self.speech_started: + results.append(await self._finish_segment("vad_end")) + return [event for event in results if event.text or event.is_final] + + def _segment_duration_ms(self) -> float: + return len(self.segment_audio) / (self.sample_rate * 2) * 1000 + + def _partial(self, text: str) -> FunASRSegment: + """Create a partial event using the cumulative FunASR text.""" + return FunASRSegment( + text=text, + start_time_ms=self.segment_start_ms, + end_time_ms=self.segment_start_ms + self._segment_duration_ms(), + audio=b"", + voiced_ms=self.segment_voiced_ms, + is_final=False, + sentence_id=self.segment_id, + ) + + async def _finish_segment(self, reason: str) -> FunASRSegment: + """Flush the ASR cache, then release completed segment audio.""" + text = await self.asr_turn.finish() if self.asr_turn is not None else "" + result = FunASRSegment( + text=text, + start_time_ms=self.segment_start_ms, + end_time_ms=self.segment_start_ms + self._segment_duration_ms(), + audio=bytes(self.segment_audio), + voiced_ms=self.segment_voiced_ms, + is_final=True, + sentence_id=self.segment_id, + reason=reason, + ) + self.segment_id += 1 + self.segment_audio.clear() + self.asr_turn = None + self.speech_started = False + self.segment_voiced_ms = 0.0 + return result + + async def feed(self, pcm_bytes: bytes) -> list[FunASRSegment]: + """Consume PCM16 and return FunASR partial/final events.""" + if len(pcm_bytes) % 2: + raise ValueError("PCM16 音频必须包含完整的双字节采样") + self.vad_buffer.extend(pcm_bytes) + results: list[FunASRSegment] = [] + while len(self.vad_buffer) >= self.vad_chunk_bytes: + chunk = bytes(self.vad_buffer[:self.vad_chunk_bytes]) + del self.vad_buffer[:self.vad_chunk_bytes] + results.extend(await self._feed_vad_chunk(chunk)) + return results + + async def finish(self) -> list[FunASRSegment]: + """Flush buffered audio and discard all stream caches.""" + results: list[FunASRSegment] = [] + if self.vad_buffer: + chunk = bytes(self.vad_buffer) + self.vad_buffer.clear() + if not self.speech_started and chunk: + self._start_segment(self.total_samples / self.sample_rate * 1000) + self.total_samples += len(chunk) // 2 + if self.speech_started: + self.segment_audio.extend(chunk) + self.segment_voiced_ms += len(chunk) / (self.sample_rate * 2) * 1000 + if self.asr_turn is not None: + for text in await self.asr_turn.append(chunk): + results.append(self._partial(text)) + if self.speech_started: + results.append(await self._finish_segment("eof")) + self.vad_cache = {} + self.pre_roll.clear() + return [event for event in results if event.text or event.is_final] diff --git a/realtime_websocket/funasr_server.py b/realtime_websocket/funasr_server.py new file mode 100644 index 0000000..8da3a17 --- /dev/null +++ b/realtime_websocket/funasr_server.py @@ -0,0 +1,642 @@ +"""Browser-facing WebSocket adapter for the migrated FunASR engine. + +The frontend contract remains the existing start/binary/stop protocol. All +speech activity detection and streaming ASR state now come from FunASR; this +module only translates engine events into the existing sentence snapshots. +""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import logging +import os +import time +import webbrowser +from dataclasses import dataclass, replace +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from aiohttp import WSMsgType, web +from dotenv import load_dotenv + +try: + from .auxiliary_service import AuxiliaryModelService, AuxiliaryServiceConfig + from .funasr_engine import FunASRModelService, FunASRSegment, FunASRServiceConfig + from .speaker_assembler import SegmentAssembler +except ImportError: + from auxiliary_service import AuxiliaryModelService, AuxiliaryServiceConfig + from funasr_engine import FunASRModelService, FunASRSegment, FunASRServiceConfig + from speaker_assembler import SegmentAssembler + + +PROJECT_ROOT = Path(__file__).resolve().parents[1] +load_dotenv(PROJECT_ROOT / ".env") + +WEB_HOST = os.getenv("WEB_HOST", "0.0.0.0") +WEB_PORT = int(os.getenv("WEB_PORT", "8082")) +WEB_DISPLAY_HOST = os.getenv("WEB_DISPLAY_HOST", "127.0.0.1") +LOCAL_ENGINE_URL = "local://funasr" +PARTIAL_BYTES_PER_SECOND = 16000 * 2 +MIN_SPEAKER_VOICE_MS = 800 +LOGGER = logging.getLogger(__name__) + + +class EndOfStream: + """Queue marker that cannot be confused with an audio frame.""" + + +EOF = EndOfStream() +MODEL_SERVICE_KEY = web.AppKey("model_service", FunASRModelService) +AUXILIARY_SERVICE_KEY = web.AppKey("auxiliary_service", AuxiliaryModelService) + + +@dataclass(frozen=True) +class SpeakerJob: + """Completed FunASR turn waiting for optional CAM++ speaker matching.""" + + sentence_id: int + audio: bytes + start_time_ms: float + end_time_ms: float + voiced_ms: float + + +@dataclass +class SessionMetrics: + """Small runtime snapshot shown in the existing frontend.""" + + started_at: float + audio_bytes: int = 0 + input_chunks: int = 0 + partial_count: int = 0 + partial_revisions: int = 0 + first_partial_ms: float | None = None + final_ms: float | None = None + + def snapshot(self) -> dict[str, Any]: + """Return JSON-safe metrics relative to session start.""" + return { + "audio_bytes": self.audio_bytes, + "input_chunks": self.input_chunks, + "partial_count": self.partial_count, + "partial_revisions": self.partial_revisions, + "first_partial_ms": self.first_partial_ms, + "final_ms": self.final_ms, + "elapsed_ms": round((time.perf_counter() - self.started_at) * 1000, 1), + } + + +class IncrementalWavDecoder: + """Strip a streamed RIFF header before forwarding PCM16 to FunASR.""" + + def __init__(self) -> None: + self.buffer = bytearray() + self.payload_started = False + self.riff_read = False + self.format_valid = False + self.data_remaining: int | None = None + + def feed(self, chunk: bytes) -> bytes: + """Parse complete RIFF chunks without buffering the whole recording.""" + if self.payload_started: + if self.data_remaining is None: + return chunk + payload = chunk[: self.data_remaining] + self.data_remaining -= len(payload) + return payload + + self.buffer.extend(chunk) + if not self.riff_read: + if len(self.buffer) < 12: + return b"" + if self.buffer[:4] != b"RIFF" or self.buffer[8:12] != b"WAVE": + raise ValueError("文件不是有效的 RIFF/WAV 音频") + del self.buffer[:12] + self.riff_read = True + + while len(self.buffer) >= 8: + kind = bytes(self.buffer[:4]) + size = int.from_bytes(self.buffer[4:8], "little") + if size > 1024 * 1024: + raise ValueError("WAV 元数据头过大,请转换为标准 PCM WAV") + chunk_size = 8 + size + (size % 2) + if len(self.buffer) < chunk_size: + return b"" + body = self.buffer[8 : 8 + size] + if kind == b"fmt ": + if size < 16: + raise ValueError("WAV fmt 区块不完整") + fields = ( + int.from_bytes(body[0:2], "little"), + int.from_bytes(body[2:4], "little"), + int.from_bytes(body[4:8], "little"), + int.from_bytes(body[14:16], "little"), + ) + if fields != (1, 1, 16000, 16): + raise ValueError("WAV 必须为 16kHz、单声道、PCM16") + self.format_valid = True + if kind == b"data": + if not self.format_valid or size % 2: + raise ValueError("WAV 必须为 16kHz、单声道、PCM16") + self.payload_started = True + self.data_remaining = size + del self.buffer[:8] + payload = bytes(self.buffer[:size]) + del self.buffer[: min(size, len(self.buffer))] + self.data_remaining -= len(payload) + return payload + del self.buffer[:chunk_size] + return b"" + + def finish(self) -> None: + """Reject a truncated or header-only WAV before final inference.""" + if not self.payload_started: + raise ValueError("WAV 文件不完整,未收到 data 音频区块") + if self.data_remaining not in (None, 0): + raise ValueError("WAV 文件不完整,未收到全部音频数据") + + +class RealtimeSession: + """Translate FunASR events into the existing sentence/display protocol.""" + + def __init__( + self, + ws: web.WebSocketResponse, + model_service: FunASRModelService, + auxiliary_service: AuxiliaryModelService | None, + start: dict[str, Any], + ) -> None: + self.ws = ws + self.model_service = model_service + self.auxiliary_service = auxiliary_service + self.start = start + self.engine = model_service.create_session() + self.session_id = uuid4().hex + self.send_lock = asyncio.Lock() + self.audio_queue: asyncio.Queue[bytes | EndOfStream] = asyncio.Queue(maxsize=128) + self.speaker_queue: asyncio.Queue[SpeakerJob | EndOfStream] = asyncio.Queue(maxsize=64) + self.assembler = SegmentAssembler() + self.metrics = SessionMetrics(time.perf_counter()) + self.source = str(start.get("source") or "mic") + self.file_name = str(start.get("file_name") or "audio.pcm") + self.wav_decoder = ( + IncrementalWavDecoder() + if self.source == "file" and Path(self.file_name).suffix.lower() == ".wav" + else None + ) + self.merge_adjacent = self._parse_flag(start.get("display_merge"), True) + # Speaker labels are required for every browser session. + self.speaker_enabled = True + self.speaker_warning_sent = False + self.input_stopped = False + + @staticmethod + def _parse_flag(value: Any, default: bool) -> bool: + """Accept booleans, 0/1 and string flags from old frontend clients.""" + if value is None: + return default + if isinstance(value, str): + return value.strip().lower() not in {"", "0", "false", "no", "off"} + return bool(value) + + async def emit(self, payload: dict[str, Any]) -> None: + """Send ordered JSON while the browser connection is still alive.""" + async with self.send_lock: + if not self.ws.closed: + await self.ws.send_json(payload) + + async def emit_state(self, sentence: dict[str, Any] | None = None) -> None: + """Send both the legacy sentence event and the full display snapshot.""" + state = { + "type": "display_state", + "raw_segments": self.assembler.raw_snapshot(), + "display_blocks": self.assembler.display_blocks(self.merge_adjacent), + "metrics": self.metrics.snapshot(), + } + if sentence is not None: + await self.emit( + { + "type": "sentences", + "sentences": [sentence], + "metrics": self.metrics.snapshot(), + } + ) + await self.emit(state) + + async def warn_speaker(self, message: str) -> None: + """Expose speaker-service failures without interrupting ASR.""" + if self.speaker_warning_sent: + return + await self.emit( + { + "type": "speaker_warning", + "session_id": self.session_id, + "speaker_service_url": getattr( + getattr(self.auxiliary_service, "config", None), "base_url", None + ), + "message": message, + } + ) + self.speaker_warning_sent = True + + async def _emit_engine_segment(self, segment: FunASRSegment) -> None: + """Map one FunASR partial/final event to a stable sentence ID.""" + if not segment.text: + if segment.is_final and self.assembler.segments.pop(segment.sentence_id, None) is not None: + await self.emit_state() + return + + sentence_type = 1 if segment.is_final else 0 + sentence = self.assembler.apply_sentence( + { + "sentence_id": segment.sentence_id, + "sentence": segment.text, + "sentence_type": sentence_type, + "start_time": segment.start_time_ms, + "end_time": segment.end_time_ms, + "speaker_id": -1, + "speaker_name": "", + "speaker_evidence": "pending", + "speaker_confidence": 0.0, + "speaker_strategy": "funasr_pending", + "commit_reason": segment.reason, + "speaker_status": ( + "queued" if segment.is_final else "waiting_final" + ) + if self.speaker_enabled + else "disabled", + "speaker_reason": ( + "等待 CAM++ 声纹处理" + if segment.is_final + else "FunASR 流式结果,等待片段结束" + ) + if self.speaker_enabled + else "说话人分离已关闭", + } + ) + if sentence_type == 0: + self.metrics.partial_count += 1 + if sentence["revision_count"] > 0: + self.metrics.partial_revisions += 1 + if self.metrics.first_partial_ms is None: + self.metrics.first_partial_ms = round( + (time.perf_counter() - self.metrics.started_at) * 1000, 1 + ) + else: + self.metrics.final_ms = round( + (time.perf_counter() - self.metrics.started_at) * 1000, 1 + ) + await self.emit_state(sentence) + + if segment.is_final and self.speaker_enabled: + await self.speaker_queue.put( + SpeakerJob( + sentence_id=segment.sentence_id, + audio=segment.audio, + start_time_ms=segment.start_time_ms, + end_time_ms=segment.end_time_ms, + voiced_ms=segment.voiced_ms, + ) + ) + + async def _resolve_speaker(self, job: SpeakerJob) -> None: + """Keep the existing optional CAM++ display integration.""" + if not self.speaker_enabled: + return + + async def update_status(status: str, reason: str) -> None: + updated = self.assembler.apply_speaker_update( + { + "sentence_id": job.sentence_id, + "speaker_id": -1, + "speaker_evidence": "pending", + "speaker_confidence": 0.0, + "speaker_status": status, + "speaker_reason": reason, + } + ) + await self.emit_state(updated) + + if job.voiced_ms < MIN_SPEAKER_VOICE_MS: + await update_status( + "insufficient_audio", + f"有效语音不足 {MIN_SPEAKER_VOICE_MS}ms,不继承上一位说话人", + ) + return + if self.auxiliary_service is None: + await update_status("service_unavailable", "未配置说话人辅助模型服务") + await self.warn_speaker("未配置辅助模型服务,无法执行说话人分离") + return + + await update_status("processing", "正在提取 CAM++ 声纹并匹配说话人") + try: + speaker = await self.auxiliary_service.resolve_speaker( + job.audio, + self.session_id, + job.start_time_ms, + job.end_time_ms, + ) + except Exception as exc: + LOGGER.exception("speaker resolve failed: session=%s sentence=%s", self.session_id, job.sentence_id) + await update_status("service_error", str(exc)) + await self.warn_speaker(str(exc)) + return + + if not speaker: + await update_status("no_embedding", "辅助服务未返回可用声纹结果") + return + update = dict(speaker) + update["sentence_id"] = job.sentence_id + update["speaker_name"] = str(update.get("speaker_name") or "") + updated = self.assembler.apply_speaker_update(update) + if updated is not None: + await self.emit_state(updated) + + async def process_speakers(self) -> None: + """Process completed turns in order so speaker clusters stay stable.""" + while True: + item = await self.speaker_queue.get() + if isinstance(item, EndOfStream): + return + await self._resolve_speaker(item) + + async def _feed(self, chunk: bytes) -> None: + """Forward raw PCM to FunASR and publish all returned events.""" + if not chunk: + return + self.metrics.audio_bytes += len(chunk) + for segment in await self.engine.feed(chunk): + await self._emit_engine_segment(segment) + + async def process_audio(self) -> None: + """Consume audio until EOF, with no local RMS/VLLM segmentation path.""" + while True: + item = await self.audio_queue.get() + if isinstance(item, EndOfStream): + break + self.metrics.input_chunks += 1 + chunk = self.wav_decoder.feed(item) if self.wav_decoder is not None else item + await self._feed(chunk) + + if self.wav_decoder is not None: + self.wav_decoder.finish() + for segment in await self.engine.finish(): + await self._emit_engine_segment(segment) + await self.emit({"type": "metrics", "metrics": self.metrics.snapshot()}) + + async def close(self) -> None: + """Drop per-session FunASR caches and temporary state.""" + self.engine = None # type: ignore[assignment] + + +async def index_handler(_: web.Request) -> web.FileResponse: + """Serve the browser page without stale-cache surprises.""" + return web.FileResponse( + Path(__file__).parent / "static" / "index.html", + headers={"Cache-Control": "no-store"}, + ) + + +async def config_handler(request: web.Request) -> web.Response: + """Expose local FunASR settings using the old response field names.""" + model_service: FunASRModelService = request.app[MODEL_SERVICE_KEY] + auxiliary = request.app.get(AUXILIARY_SERVICE_KEY) + response = web.json_response( + { + "model_service_url": LOCAL_ENGINE_URL, + "model": model_service.config.model, + "engine": "funasr", + "speaker_service_url": getattr( + getattr(auxiliary, "config", None), "base_url", None + ), + } + ) + # The standalone frontend reads this API from its own origin. + response.headers["Access-Control-Allow-Origin"] = os.getenv( + "FRONTEND_ORIGIN", "http://127.0.0.1:8080" + ) + return response + + +async def websocket_handler(request: web.Request) -> web.WebSocketResponse: + """Keep the old browser protocol while using FunASR internally.""" + ws = web.WebSocketResponse(max_msg_size=64 * 1024 * 1024) + await ws.prepare(request) + model_service: FunASRModelService = request.app[MODEL_SERVICE_KEY] + auxiliary_service = request.app.get(AUXILIARY_SERVICE_KEY) + processing: asyncio.Task[None] | None = None + speaker_processing: asyncio.Task[None] | None = None + session: RealtimeSession | None = None + input_finished = False + + try: + first = await ws.receive() + if first.type != WSMsgType.TEXT: + await ws.send_json({"type": "error", "message": "first message must be JSON start"}) + return ws + try: + start = json.loads(first.data) + except json.JSONDecodeError: + await ws.send_json({"type": "error", "message": "invalid start JSON"}) + return ws + if not isinstance(start, dict) or start.get("type") != "start": + await ws.send_json({"type": "error", "message": "first message must have type=start"}) + return ws + + source = str(start.get("source") or "mic") + suffix = Path(str(start.get("file_name") or "")).suffix.lower() + if source == "file" and suffix not in {".pcm", ".wav"}: + await ws.send_json( + { + "type": "error", + "message": "实时流式测试的文件模式只支持 PCM 或 WAV,请改用麦克风、PCM 或 WAV", + } + ) + return ws + + session = RealtimeSession(ws, model_service, auxiliary_service, start) + speaker_health: dict[str, Any] | None = None + speaker_health_error: str | None = None + if session.speaker_enabled and auxiliary_service is not None: + try: + speaker_health = await asyncio.wait_for(auxiliary_service.health(), timeout=5) + if speaker_health.get("speaker_embedding_ready", speaker_health.get("ready")) is False: + speaker_health_error = "辅助模型服务未就绪,请检查 /health 返回的 models 状态" + except Exception as exc: + speaker_health_error = f"说话人辅助服务不可用:{exc}" + + await session.emit( + { + "type": "start", + "model_service_url": LOCAL_ENGINE_URL, + "model": model_service.config.model, + "engine": "funasr", + "session_id": session.session_id, + "enable_native_partial_stream": True, + "native_partial_supported": True, + "partial_mode": "funasr_streaming_cache", + "speaker_diarization_enabled": session.speaker_enabled, + "speaker_service_url": getattr( + getattr(auxiliary_service, "config", None), "base_url", None + ), + "speaker_service_health": speaker_health, + "speaker_gap_enabled": False, + "sentence_strategy": start.get("sentence_strategy", 0), + "display_state_supported": True, + } + ) + if speaker_health_error: + await session.warn_speaker(speaker_health_error) + + processing = asyncio.create_task(session.process_audio()) + if session.speaker_enabled: + speaker_processing = asyncio.create_task(session.process_speakers()) + + async def receive_or_raise() -> Any: + """Wake up immediately when the engine worker fails.""" + receive_task = asyncio.create_task(ws.receive()) + workers = [task for task in (processing, speaker_processing) if task is not None] + done, _ = await asyncio.wait( + [receive_task, *workers], + return_when=asyncio.FIRST_COMPLETED, + ) + if receive_task in done: + return await receive_task + receive_task.cancel() + await asyncio.gather(receive_task, return_exceptions=True) + for worker in workers: + if worker in done: + await worker + raise RuntimeError("FunASR 实时处理任务意外结束") + + while not ws.closed: + message = await receive_or_raise() + if message.type == WSMsgType.BINARY: + await session.audio_queue.put(bytes(message.data)) + continue + if message.type == WSMsgType.TEXT: + try: + control = json.loads(message.data) + except json.JSONDecodeError: + continue + if not isinstance(control, dict): + continue + if control.get("type") in {"eof", "stop"}: + input_finished = True + session.input_stopped = control.get("type") == "stop" + await session.emit({"type": "draining", "message": "FunASR 正在完成最终识别"}) + await session.audio_queue.put(EOF) + break + if control.get("type") == "abort": + return ws + if message.type in {WSMsgType.ERROR, WSMsgType.CLOSE, WSMsgType.CLOSED}: + return ws + + if input_finished: + await processing + if speaker_processing is not None: + await session.speaker_queue.put(EOF) + await speaker_processing + await session.emit_state() + await session.emit( + { + "type": "end", + "metrics": session.metrics.snapshot(), + "sentences": session.assembler.raw_snapshot(), + "display_blocks": session.assembler.display_blocks(session.merge_adjacent), + } + ) + except asyncio.CancelledError: + raise + except Exception as exc: + LOGGER.exception("FunASR WebSocket session failed") + if not ws.closed: + await ws.send_json({"type": "error", "message": str(exc)}) + finally: + for task in (processing, speaker_processing): + if task is not None and not task.done(): + task.cancel() + await asyncio.gather( + *(task for task in (processing, speaker_processing) if task is not None), + return_exceptions=True, + ) + if session is not None: + reset = getattr(session.auxiliary_service, "reset_speaker_session", None) + await session.close() + if reset is not None: + try: + await asyncio.wait_for(reset(session.session_id), timeout=5) + except Exception: + LOGGER.warning("speaker session cleanup failed: %s", session.session_id, exc_info=True) + if not ws.closed: + await ws.close() + return ws + + +async def start_app( + model: str | None = None, + device: str | None = None, +) -> web.Application: + """Create the FunASR-backed frontend application.""" + config = FunASRServiceConfig.from_env() + if model: + config = replace(config, model=model) + if device: + config = replace(config, device=device) + + app = web.Application() + app[MODEL_SERVICE_KEY] = FunASRModelService(config) + app[AUXILIARY_SERVICE_KEY] = AuxiliaryModelService( + AuxiliaryServiceConfig( + base_url=os.getenv("AUXILIARY_SERVICE_URL", "http://127.0.0.1:8010") + ) + ) + + async def lifecycle(application: web.Application): + await application[MODEL_SERVICE_KEY].start() + await application[AUXILIARY_SERVICE_KEY].start() + try: + health = await application[AUXILIARY_SERVICE_KEY].health() + if not health.get("speaker_embedding_ready"): + raise RuntimeError("CAM++ speaker service is not ready") + except Exception: + await application[AUXILIARY_SERVICE_KEY].close() + await application[MODEL_SERVICE_KEY].close() + raise + yield + await application[AUXILIARY_SERVICE_KEY].close() + await application[MODEL_SERVICE_KEY].close() + + app.cleanup_ctx.append(lifecycle) + app.router.add_get("/", index_handler) + app.router.add_get("/api/config", config_handler) + app.router.add_static("/static/", Path(__file__).parent / "static") + app.router.add_get("/ws", websocket_handler) + app.router.add_static("/", Path(__file__).parent / "static", show_index=False) + return app + + +def main() -> None: + """Start the FunASR-backed browser demo.""" + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--model", default=os.getenv("FUNASR_ASR_MODEL")) + parser.add_argument("--device", default=os.getenv("FUNASR_DEVICE")) + parser.add_argument("--no-browser", action="store_true") + args = parser.parse_args() + logging.basicConfig(level=logging.INFO) + if not args.no_browser: + webbrowser.open(f"http://{WEB_DISPLAY_HOST}:{WEB_PORT}/") + print(f"FunASR demo: http://{WEB_DISPLAY_HOST}:{WEB_PORT}/", flush=True) + print(f"FunASR model: {args.model or FunASRServiceConfig.from_env().model}", flush=True) + web.run_app( + start_app(model=args.model, device=args.device), + host=WEB_HOST, + port=WEB_PORT, + ) + + +if __name__ == "__main__": + main() diff --git a/realtime_websocket/model_service.py b/realtime_websocket/model_service_qwen_legacy.py similarity index 100% rename from realtime_websocket/model_service.py rename to realtime_websocket/model_service_qwen_legacy.py diff --git a/realtime_websocket/requirements.txt b/realtime_websocket/requirements.txt index ed0d0a1..53492e9 100644 --- a/realtime_websocket/requirements.txt +++ b/realtime_websocket/requirements.txt @@ -1,2 +1,6 @@ aiohttp==3.11.11 python-dotenv>=1.0 +funasr==1.4.16 +modelscope[framework]==1.34.0 +soundfile==0.13.1 +librosa==0.11.0 diff --git a/realtime_websocket/server.py b/realtime_websocket/server_qwen_legacy.py similarity index 100% rename from realtime_websocket/server.py rename to realtime_websocket/server_qwen_legacy.py diff --git a/realtime_websocket/static/app.js b/realtime_websocket/static/app.js index 944df74..055d076 100644 --- a/realtime_websocket/static/app.js +++ b/realtime_websocket/static/app.js @@ -3,8 +3,6 @@ const elEngineModel = document.getElementById('engineModel'); const elModelServiceUrl = document.getElementById('modelServiceUrl'); const elSpeakerStatus = document.getElementById('speakerStatus'); const elDisplayMerge = document.getElementById('displayMerge'); -const elSpeakerDiarization = document.getElementById('speakerDiarization'); -const elDiarizationLabel = document.getElementById('diarizationLabel'); const elSentenceStrategy = document.getElementById('sentenceStrategy'); const elBtnStart = document.getElementById('btnStart'); const elBtnStop = document.getElementById('btnStop'); @@ -450,7 +448,7 @@ async function startRecognition() { sending = true; const currentSession = ++sessionId; - const useSpeaker = elSpeakerDiarization.checked; + const useSpeaker = true; // Speaker labels are mandatory for this deployment. let receivedTerminal = false; // 构造 WebSocket 首条 start 消息。 @@ -475,8 +473,10 @@ async function startRecognition() { speed_factor: speedFactor }; - const protocol = location.protocol === 'https:' ? 'wss:' : 'ws:'; - ws = new WebSocket(`${protocol}//${location.host}/ws`); + const backendUrl = (window.ASR_BACKEND_URL || location.origin).replace(/\/$/, ''); + const websocketUrl = new URL(backendUrl + '/ws'); + websocketUrl.protocol = backendUrl.startsWith('https:') ? 'wss:' : 'ws:'; + ws = new WebSocket(websocketUrl); ws.binaryType = 'arraybuffer'; ws.onopen = () => { @@ -725,7 +725,9 @@ function stopMicCapture() { } // 展示实际部署端点及模型,避免沿用旧 SDK 的无效引擎配置。 -fetch('/api/config').then(response => response.json()).then(config => { +// The static frontend reads API and WS addresses from runtime-config.js. +const backendUrl = (window.ASR_BACKEND_URL || location.origin).replace(/\/$/, ''); +fetch(backendUrl + '/api/config').then(response => response.json()).then(config => { if (!ws) { elEngineModel.value = config.model; elModelServiceUrl.value = config.model_service_url; diff --git a/realtime_websocket/static/index.html b/realtime_websocket/static/index.html index d34384d..71282c6 100644 --- a/realtime_websocket/static/index.html +++ b/realtime_websocket/static/index.html @@ -20,11 +20,11 @@