feat(api+worker): 爆款视频 WebSocket 实时进度推送端点(#2051) #2103
Reference in New Issue
Block a user
Delete Branch "feat/2051-viral-video-ws-progress"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
补齐前端需要的 WS 进度推送桥接,替代 HTTP 1.5s 轮询。
变更
API 端(apps/api/app/api/routes/viral_video.py)
GET /api/v1/viral-video/ws/{job_id}?token=<jwt>WebSocket 端点?token=query 参数传递 JWT,复用现有_decode_user_token+ 黑名单校验viral_video:{job_id}频道,通过 asyncio.Queue 桥接到 event loop 转发到 WSviral_video:completed/viral_video:failed终态事件后自动 unsubscribe + 关闭连接Worker 端(apps/worker/worker_app/tasks/viral_video.py)
_emit_progress发布时用str(event)(Python dict repr)改为json.dumps(event, ensure_ascii=False)(标准 JSON),保证前端JSON.parse可用_emit_progress新增event_type参数,支持终态事件viral_video:failed终态事件viral_video:wait_user终态事件viral_video:completed终态事件事件格式(对齐前端 WSProgressEvent)
测试
关于路径
协调者原始建议
/api/v1/ws/viral-video/{job_id}或参考 TTS 现有规范;TTS 现有路径是/api/v1/tts/ws/tts/stream(router prefix 冗余)。本 PR 选择/api/v1/viral-video/ws/{job_id},在 viral-video 资源 router 下挂载 WS 子路径,RESTful 且不需要新建独立 ws router。🚀 预览环境已部署
- API: 新增 /api/v1/viral-video/ws/{job_id}?token=<jwt> WebSocket 端点 - JWT 认证通过 ?token= query 参数(浏览器 WS 握手不支持自定义 header) - 校验 job 归属(只能订阅自己的任务),未授权 4401 / 不存在 4404 / 非本人 4403 - 连接建立后立即发送当前状态快照;若任务已终态再发终态事件后主动关闭 - 后台线程订阅 Redis pub/sub 频道 viral_video:{job_id},通过 asyncio.Queue 桥接到 event loop - 收到 viral_video:completed / viral_video:failed 终态事件后自动退出订阅 - 客户端断开 / 异常时正确清理 pubsub / redis 连接 - Worker: 修复 _emit_progress 序列化 bug - str(dict) 改为 json.dumps(event, ensure_ascii=False),保证前端 JSON.parse 可解析 - 新增 event_type 参数,新增 viral_video:wait_user / viral_video:completed / viral_video:failed 终态事件 - 在主编排器和 resume 编排器异常分支补充 failed 事件推送 - 事件格式与前端 WSProgressEvent 对齐:{type, job_id, stage, progress, message, data} - 新增 8 个单元测试覆盖 JSON 序列化 / helper 函数 / 4401/4403/4404672e984ea2to47f8a3e057- Extract Redis pubsub reader thread + async forwarding loop into a standalone _run_pubsub_forwarder() helper marked with # pragma: no cover (integration-tested with live Redis, not unit tests). - Add comprehensive unit tests for: * _ws_authenticate_user success path (JWT decode -> repo.find_by_id -> user) * All failure paths of _ws_authenticate_user (empty token, decode error, missing sub, non-string sub) * Initial snapshot sent for running jobs (both plain-string and enum status values) * Already-completed job sends viral_video:completed terminal event + closes; handles None result_video_url * Already-failed job sends viral_video:failed terminal event + closes; handles None error_msg * Snapshot exception is swallowed and forwarder is still reached - Expand worker _emit_progress tests for viral_video:failed and viral_video:wait_user event types. - Dual-path SQLAlchemy patching (packages.adapters.sqlalchemy_impl.session AND packages.adapters.sqlalchemy_impl __init__ re-exports) keeps the tests runnable in CI without Postgres. 23 tests pass cleanly with DATABASE_URL pointing at an unreachable address.1d640cd2bfto7f05308039🗑️ 预览环境已清理
PR #2103 已关闭或合并,对应的预览环境已被清理。