CMVR-AI-ANALYSIS/app/main.py
lixiaolong b386f003a0 feat(vision): 集成Qwen3.8视觉模型并优化音视频分析功能
- 集成Qwen3.8-27B-FP8快速模型和Qwen3.8-27B精确模型作为SGLang服务
- 添加SGLang API配置选项(SGLANG_FAST_BASE_URL、SGLANG_ACCURATE_BASE_URL等)
- 实现音频分类中的决策聚合算法(topK、nearestWeight、labelMaxDistance)
- 添加视频采样帧限制(VIDEO_SAMPLING_FRAME_LIMIT)和上下文token限制
- 更新健康检查以监控SGLang服务状态
- 实现视频分析的双模式决策策略(快速+精确)
- 添加音频参考文件导入工具(import_audio_references.py)
- 扩展音频分类标签支持FIND_VEHICLE_HORN类别
- 优化视频分析的帧采样策略,始终包含视频尾部帧
- 添加决策策略参数(tuning、decisionPolicy)支持
- 更新配置类以支持新的SGLang和视频参数
- 修改compose配置以支持Qwen3.8模型部署
- 更新音频分类测试用例验证聚合逻辑
- 重构视频测试以支持SGLang API格式和决策策略
2026-08-19 09:28:53 +08:00

146 lines
5.5 KiB
Python

import asyncio
import shutil
import time
from contextlib import asynccontextmanager
import httpx
from fastapi import Depends, FastAPI, Header, HTTPException
from app.audio import classify_audio
from app.config import settings
from app.media import download_media
from app.profile_store import profile_store
from app.schemas import AnalysisRequest, AnalysisResponse, AnalysisType
from app.video import VisionModelError, analyze_video
def authorize(authorization: str | None = Header(default=None)) -> None:
if not settings.api_key:
return
if authorization != f"Bearer {settings.api_key}":
raise HTTPException(status_code=401, detail="Invalid analysis service credential")
@asynccontextmanager
async def lifespan(_: FastAPI):
settings.jobs_dir.mkdir(parents=True, exist_ok=True)
settings.artifacts_dir.mkdir(parents=True, exist_ok=True)
profile_store.reload()
yield
app = FastAPI(title="CMVR Media Analysis Service", version="1.0.0", lifespan=lifespan)
@app.get("/health")
async def health() -> dict:
async def probe(url: str) -> str:
try:
async with httpx.AsyncClient(timeout=2) as client:
response = await client.get(url)
response.raise_for_status()
return "UP"
except Exception:
return "DOWN"
sglang_urls = {
"fast": settings.sglang_fast_base_url,
"accurate": settings.sglang_accurate_base_url,
}
checks = [probe(f"{settings.ollama_base_url.rstrip('/')}/api/version")]
configured_names = [name for name, url in sglang_urls.items() if url]
checks.extend(
probe(f"{sglang_urls[name].rstrip('/')}/models") for name in configured_names
)
statuses = await asyncio.gather(*checks)
ollama = statuses[0]
configured_statuses = dict(zip(configured_names, statuses[1:]))
sglang = {
name: configured_statuses.get(name, "NOT_CONFIGURED")
for name in sglang_urls
}
return {
"status": "UP",
"ollama": ollama,
"sglang": sglang,
"visionModel": settings.ollama_model,
"profiles": profile_store.status(),
}
@app.post(
"/api/v1/analysis/run",
response_model=AnalysisResponse,
dependencies=[Depends(authorize)],
)
async def run_analysis(request: AnalysisRequest) -> AnalysisResponse:
started = time.monotonic()
media_path = None
try:
profile = profile_store.get(request.profileCode, request.analysisType.value)
if request.analysisType == AnalysisType.AUDIO_CLASSIFICATION:
media_path = await download_media(
str(request.mediaUrl), ".audio", settings.max_audio_bytes
)
result = classify_audio(profile, media_path)
model = {"provider": "CMVR", "name": "mfcc-dtw-audio-fingerprint-v2"}
else:
media_path = await download_media(
str(request.mediaUrl), ".video", settings.max_video_bytes
)
result, duration, video_metadata = await analyze_video(
profile,
media_path,
str(request.options.get("instruction", "")),
str(request.options.get("analysisMode", "AUTO")),
request.options.get("tuning"),
str(request.options.get("decisionPolicy", "FAIL_CLOSED")),
)
evidence = result.get("evidence")
if not isinstance(evidence, dict):
evidence = {}
result["evidence"] = evidence
evidence.update(
{
"sampledFrameCount": video_metadata["sampledFrameCount"],
"sampleFps": video_metadata["effectiveSampleFps"],
"samplingFrameLimit": video_metadata["samplingFrameLimit"],
"samplingCapped": video_metadata["samplingCapped"],
"numCtx": video_metadata["numCtx"],
"maximumWidth": video_metadata["maximumWidth"],
"durationSeconds": duration,
}
)
result["analysisMode"] = video_metadata["mode"]
result["fallback"] = video_metadata["fallback"]
result["fallbackReason"] = video_metadata["fallbackReason"]
result["effectiveTuning"] = video_metadata["tuning"]
model = {
"provider": video_metadata["provider"],
"name": video_metadata["model"],
"primaryProvider": video_metadata["primaryProvider"],
"primaryModel": video_metadata["primaryModel"],
"providerFallback": video_metadata["providerFallback"],
"requestedMode": video_metadata["requestedMode"],
"usedMode": video_metadata["mode"],
"fallback": video_metadata["fallback"],
}
return AnalysisResponse(
requestId=request.requestId,
analysisType=request.analysisType,
profileCode=request.profileCode,
status="SUCCEEDED",
result=result,
model=model,
timingMs=round((time.monotonic() - started) * 1000),
)
except (KeyError, ValueError) as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
except httpx.HTTPError as exc:
raise HTTPException(status_code=502, detail=f"Remote service request failed: {exc}") from exc
except VisionModelError as exc:
raise HTTPException(status_code=502, detail=str(exc)) from exc
finally:
if media_path is not None:
media_path.unlink(missing_ok=True)