"""DashScope 客户端(阿里云百炼 Wan 3.0 等非方舟模型)。 #2159: 新增 Wan 3.0 视频生成支持。DashScope 异步协议: - POST {base_url}/services/aigc/video-generation/video-synthesis (X-DashScope-Async: enable) → 返回 output.task_id - GET {base_url}/tasks/{task_id} 轮询状态 → SUCCEEDED 时 output.video_url 可下载 认证:Authorization: Bearer {DASHSCOPE_API_KEY} """ from __future__ import annotations import logging import os import time from pathlib import Path from typing import Any from urllib.parse import urlparse import httpx from packages.shared.config import get_shared_settings logger = logging.getLogger(__name__) _DASHSCOPE_CLIENT_SINGLETON: "DashScopeClient | None" = None class DashScopeClient: """阿里云 DashScope 异步 API 客户端(Wan 3.0 等视频生成)。""" def __init__(self) -> None: settings = get_shared_settings() self.api_key: str = getattr(settings, "dashscope_api_key", "") or os.getenv("DASHSCOPE_API_KEY", "") self.base_url: str = ( getattr(settings, "dashscope_base_url", "") or "https://dashscope.aliyuncs.com/api/v1" ).rstrip("/") self.poll_interval: int = int(getattr(settings, "dashscope_video_poll_interval", 10) or 10) self.total_timeout: int = int(getattr(settings, "dashscope_video_timeout", 900) or 900) self.max_retries: int = 2 @property def is_available(self) -> bool: return bool(self.api_key) def video_generation( self, prompt: str, *, image_url: str | None = None, duration: int = 5, ratio: str | None = "9:16", resolution: str = "720p", watermark: bool = False, output_dir: str | None = None, model: str = "wan3.0-video", ) -> dict | None: """调用 DashScope 异步视频合成接口,轮询完成后下载到本地。 返回 {"video_path": str, "usage": dict | None};失败返回 None。 """ if not self.is_available: logger.error("[dashscope] API key 未配置,无法调用视频生成") return None if not prompt or not prompt.strip(): return None # DashScope 分辨率参数:720P / 1080P / 480P(大写 P) res_upper = (resolution or "720p").upper().replace("P", "P") if res_upper == "480P": ds_res = "480P" elif res_upper == "1080P": ds_res = "1080P" else: ds_res = "720P" # 构造 input+parameters input_obj: dict[str, Any] = {"prompt": prompt.strip()} if image_url: input_obj["img_url"] = image_url params: dict[str, Any] = { "resolution": ds_res, "duration": str(float(duration)), "watermark": bool(watermark), } # 比例透传:Wan 支持 "9:16" / "16:9" / "1:1" 等 if ratio and ratio != "adaptive": params["aspect_ratio"] = ratio payload: dict[str, Any] = { "model": model, "input": input_obj, "parameters": params, } headers = { "Authorization": f"Bearer {self.api_key}", "Content-Type": "application/json", "X-DashScope-Async": "enable", } create_url = f"{self.base_url}/services/aigc/video-generation/video-synthesis" logger.info( "[dashscope] 创建任务: model=%s dur=%ds ratio=%s res=%s img=%s", model, duration, ratio, ds_res, bool(image_url), ) # 创建任务 task_id: str | None = None last_err: Exception | None = None for attempt in range(self.max_retries + 1): try: resp = httpx.post(create_url, headers=headers, json=payload, timeout=60) sc = int(getattr(resp, "status_code", 0) or 0) body_text = (getattr(resp, "text", "") or "")[:1500] if sc >= 400: logger.error("[dashscope] 创建任务 HTTP %d: %s", sc, body_text) resp.raise_for_status() data = resp.json() tid = (data.get("output") or {}).get("task_id") if tid: task_id = tid break # 部分情况下 code != 错误 code = data.get("code") if code and code != "": last_err = RuntimeError(f"dashscope create failed: {body_text[:300]}") else: last_err = RuntimeError(f"create ok but no task_id: {str(data)[:300]}") except Exception as e: last_err = e if attempt < self.max_retries: time.sleep(0.5 * (2**attempt)) continue logger.error("[dashscope] 创建任务最终失败: %s", last_err) return None if not task_id: return None # 轮询任务 poll_url = f"{self.base_url}/tasks/{task_id}" deadline = time.time() + self.total_timeout video_url: str | None = None usage: dict | None = None while time.time() < deadline: try: r = httpx.get(poll_url, headers=headers, timeout=30) if int(getattr(r, "status_code", 0) or 0) >= 400: logger.warning("[dashscope] 轮询 HTTP %d", r.status_code) time.sleep(self.poll_interval) continue d = r.json() out = d.get("output") or {} task_status = out.get("task_status") or d.get("task_status") or "" if task_status == "SUCCEEDED": video_url = out.get("video_url") or "" usage = d.get("usage") if not video_url: # 结果在 results 数组 results = out.get("results") or [] if results and isinstance(results, list): video_url = results[0].get("url") or results[0].get("video_url") if video_url: logger.info("[dashscope] 任务 %s 完成: %s", task_id, video_url[:120]) break logger.error("[dashscope] 任务 %s SUCCEEDED 但无 video_url", task_id) return None if task_status in ("FAILED", "FAILED_WITH_ERROR", "ERROR"): msg = out.get("message") or d.get("message") or "unknown error" logger.error("[dashscope] 任务 %s 失败: %s", task_id, msg) return None if task_status in ("CANCELED", "CANCELLED"): logger.warning("[dashscope] 任务 %s 被取消", task_id) return None # PENDING / RUNNING / SUSPENDED → 继续轮询 logger.debug("[dashscope] 任务 %s 状态 %s,继续轮询", task_id, task_status) except Exception as e: logger.warning("[dashscope] 轮询异常: %s", e) time.sleep(self.poll_interval) if not video_url: logger.error("[dashscope] 任务 %s 轮询超时(%ds)", task_id, self.total_timeout) return None # 下载视频 out_dir = output_dir or os.path.join(os.getcwd(), "seedance_outputs") os.makedirs(out_dir, exist_ok=True) suffix = Path(urlparse(video_url).path).suffix or ".mp4" if suffix.lower() not in (".mp4", ".mov", ".webm"): suffix = ".mp4" safe_tid = "".join(c if c.isalnum() or c in "-_" else "_" for c in task_id)[:40] out_path = os.path.join(out_dir, f"wan_{safe_tid}{suffix}") try: with httpx.stream("GET", video_url, timeout=300, follow_redirects=True) as resp: if int(getattr(resp, "status_code", 0) or 0) >= 400: logger.error("[dashscope] 下载 HTTP %d", resp.status_code) return None with open(out_path, "wb") as f: for chunk in resp.iter_bytes(chunk_size=1024 * 256): if chunk: f.write(chunk) except Exception as e: logger.error("[dashscope] 下载视频失败: %s", e) return None size = os.path.getsize(out_path) if os.path.exists(out_path) else 0 if size < 1024: logger.error("[dashscope] 下载文件过小: %d bytes", size) return None logger.info("[dashscope] 视频已下载: %s (%d bytes)", out_path, size) return {"video_path": out_path, "usage": usage} def get_dashscope_client() -> DashScopeClient | None: """返回 DashScope 客户端单例;未配置 API key 时返回 None。""" global _DASHSCOPE_CLIENT_SINGLETON if _DASHSCOPE_CLIENT_SINGLETON is None: _DASHSCOPE_CLIENT_SINGLETON = DashScopeClient() if not _DASHSCOPE_CLIENT_SINGLETON.is_available: return None return _DASHSCOPE_CLIENT_SINGLETON