Compare commits
18 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ba3e97c986 | |||
| 3106496c12 | |||
| f79f75b863 | |||
| b7a439d319 | |||
| 4d7c80ae07 | |||
| 26c0140d79 | |||
| c507f76c14 | |||
| 8a5cfe831e | |||
| d7d3f3184b | |||
| ec9240b52a | |||
| d0e5ef1753 | |||
| 6e8199581d | |||
| 1e23a3f094 | |||
| b04a803655 | |||
| e496f127a3 | |||
| 423be1446f | |||
| 6d9d2e8179 | |||
| b233529eee |
@@ -3,7 +3,8 @@
|
||||
*/
|
||||
|
||||
/** 任务状态 */
|
||||
export type TaskStatus = "pending" | "waiting" | "running" | "completed" | "failed" | "cancelled"
|
||||
export type TaskStatus =
|
||||
"pending" | "waiting" | "running" | "awaiting_cover" | "completed" | "failed" | "cancelled"
|
||||
|
||||
/** 任务类型 */
|
||||
export type TaskType = "ingest" | "generation" | string
|
||||
|
||||
@@ -239,9 +239,14 @@ const GeneratePage: React.FC = () => {
|
||||
voiceModePerVideo,
|
||||
variantCoverUrls: previewCovers,
|
||||
selectedVariantIndexes: isBatch ? selectedVariantIds : undefined,
|
||||
onGenerationSuccess: () => {
|
||||
onGenerationSuccess: (status?: "completed" | "awaiting_cover") => {
|
||||
setPreviewTaskId(null)
|
||||
setStoredSourceEditPlanId(null)
|
||||
// #2088:渲染完成后自动跳到封面选择页(step 5),不再等用户手动点「下一步」
|
||||
// awaiting_cover 和 completed 都走封面页(completed 是旧 worker 或 finalize 后状态,仍支持选封面)
|
||||
if (status === "awaiting_cover" || status === "completed" || !status) {
|
||||
setCurrentStep(5)
|
||||
}
|
||||
},
|
||||
})
|
||||
|
||||
|
||||
@@ -41,8 +41,8 @@ export interface UseGenerateVideoProps {
|
||||
enabled: boolean
|
||||
music_id?: string
|
||||
}
|
||||
/** 生成成功后的回调(用于清除持久化的 previewTaskId 等状态) */
|
||||
onGenerationSuccess?: () => void
|
||||
/** 生成成功后的回调(用于清除持久化的 previewTaskId 等状态);status=awaiting_cover 表示需进封面选择 */
|
||||
onGenerationSuccess?: (status?: "completed" | "awaiting_cover") => void
|
||||
/* ── 批量生成(#1677)── */
|
||||
/** 生成数量(1=单条旧逻辑,>1=批量) */
|
||||
previewCount?: number
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { useRef, useCallback, useState } from "react"
|
||||
import { useRef, useCallback, useState, useEffect } from "react"
|
||||
import { message } from "antd"
|
||||
import axios from "axios"
|
||||
import { getGenerationTask, retryTask as retryGenerationTaskApi } from "@/api/tasks/tasks"
|
||||
@@ -19,7 +19,7 @@ export interface BatchTaskState {
|
||||
|
||||
interface UseGenerationPollingOptions {
|
||||
onProgress: (progress: number) => void
|
||||
onComplete: (videos: unknown[]) => void
|
||||
onComplete: (videos: unknown[], taskStatus?: "completed" | "awaiting_cover") => void
|
||||
onFailed: (errorMsg: string) => void
|
||||
/** 批量:单任务状态变化(第5步逐卡片展示) */
|
||||
onBatchTaskUpdate?: (taskId: string, patch: Partial<BatchTaskState>) => void
|
||||
@@ -31,13 +31,19 @@ const MAX_RETRYABLE_ERRORS = 10
|
||||
const MAX_RESULTS_RETRIES = 3
|
||||
|
||||
/**
|
||||
* 生成状态轮询 Hook(v4 — 批量任务独立状态 + 单任务重试)
|
||||
* 生成状态轮询 Hook(v5 — awaiting_cover 状态识别 + visibilitychange 恢复 + 状态透传)
|
||||
*
|
||||
* startPolling(taskId) 轮询单个任务;
|
||||
* startPollingBatch(tasks) 并行轮询 N 个任务:
|
||||
* - 每个任务独立进度/状态/失败,通过 onBatchTaskUpdate 实时回传
|
||||
* - 全部成功才 onComplete(聚合视频按变体顺序);任一失败不影响其他任务继续
|
||||
* - retryTask(taskId) 单独重试失败任务(重新轮询,后端任务仍在跑则直接接续)
|
||||
*
|
||||
* v5 修复(#2088):
|
||||
* 1. 单任务路径透传 taskStatus(completed / awaiting_cover)到 onComplete,外层据此区分跳转
|
||||
* 2. 监听 visibilitychange,页面从后台切回可见时立即补拉一次,解决切后台 setInterval 被浏览器
|
||||
* 降频/冻结导致进度卡在 56% 的问题
|
||||
* 3. 非 4xx/5xx 网络错误按 3s 退避重试(已有 MAX_RETRYABLE_ERRORS=10 兜底)
|
||||
*/
|
||||
export function useGenerationPolling({
|
||||
onProgress,
|
||||
@@ -49,12 +55,15 @@ export function useGenerationPolling({
|
||||
const cancelledRef = useRef(false)
|
||||
/** 批量任务上下文:taskId → 变体序号 */
|
||||
const batchContextRef = useRef<Map<string, number>>(new Map())
|
||||
/** 当前活跃的「立刻补拉一次」函数(visibilitychange 回调使用) */
|
||||
const immediateTickRef = useRef<(() => void) | null>(null)
|
||||
const [, forceTick] = useState(0)
|
||||
|
||||
const clearTimer = useCallback(() => {
|
||||
cancelledRef.current = true
|
||||
progressTimer.current.forEach((t) => clearTimeout(t))
|
||||
progressTimer.current = []
|
||||
immediateTickRef.current = null
|
||||
}, [])
|
||||
|
||||
/** 任务完成后拉取结果列表,带重试 */
|
||||
@@ -113,6 +122,7 @@ export function useGenerationPolling({
|
||||
|
||||
if (task.status === "completed" || task.status === "awaiting_cover") {
|
||||
done = true
|
||||
immediateTickRef.current = null
|
||||
const videos = await fetchResultsWithRetry(taskId)
|
||||
if (cancelledRef.current) return
|
||||
if (videos === null) {
|
||||
@@ -128,6 +138,7 @@ export function useGenerationPolling({
|
||||
|
||||
if (task.status === "failed" || task.status === "cancelled") {
|
||||
done = true
|
||||
immediateTickRef.current = null
|
||||
const rawMsg =
|
||||
task.error_info?.error_message ||
|
||||
task.error_message ||
|
||||
@@ -138,6 +149,7 @@ export function useGenerationPolling({
|
||||
return
|
||||
}
|
||||
|
||||
// running / pending / waiting:更新进度并安排下一次轮询
|
||||
const pct = Math.max(0, Math.min(99, Math.round(Number(task.progress) || 0)))
|
||||
callbacks?.onTaskProgress?.(pct)
|
||||
if (!callbacks && runId === 0) {
|
||||
@@ -149,16 +161,20 @@ export function useGenerationPolling({
|
||||
if (cancelledRef.current || done) return
|
||||
console.error("[轮询出错] taskId:", taskId, pollErr)
|
||||
const status = axios.isAxiosError(pollErr) ? pollErr.response?.status : undefined
|
||||
// 4xx 视为不可重试(任务不存在/权限问题等),直接失败
|
||||
if (status && status >= 400 && status < 500) {
|
||||
done = true
|
||||
immediateTickRef.current = null
|
||||
const msg = extractErrorMessage(pollErr, status)
|
||||
callbacks?.onTaskFailed?.(msg)
|
||||
reject(new Error(msg))
|
||||
return
|
||||
}
|
||||
// 网络错误 / 5xx:3s 退避重试,最多 MAX_RETRYABLE_ERRORS 次
|
||||
consecutiveErrors += 1
|
||||
if (consecutiveErrors >= MAX_RETRYABLE_ERRORS) {
|
||||
done = true
|
||||
immediateTickRef.current = null
|
||||
const msg = "任务状态查询连续失败,请稍后在任务列表查看结果"
|
||||
callbacks?.onTaskFailed?.(msg)
|
||||
reject(new Error(msg))
|
||||
@@ -169,6 +185,18 @@ export function useGenerationPolling({
|
||||
}
|
||||
}
|
||||
|
||||
// 注册「立刻补拉一次」回调,供 visibilitychange 恢复时调用
|
||||
// 注意:必须在 done 后清理,避免切换页面时误触发已结束任务的补拉
|
||||
immediateTickRef.current = () => {
|
||||
if (!done && !cancelledRef.current) {
|
||||
// 清除未触发的 setTimeout,立即拉一次
|
||||
progressTimer.current.forEach((t) => clearTimeout(t))
|
||||
progressTimer.current = []
|
||||
consecutiveErrors = 0
|
||||
void poll()
|
||||
}
|
||||
}
|
||||
|
||||
const timer = setTimeout(poll, 1500)
|
||||
progressTimer.current.push(timer)
|
||||
})
|
||||
@@ -181,12 +209,22 @@ export function useGenerationPolling({
|
||||
(taskId: string) => {
|
||||
cancelledRef.current = false
|
||||
batchContextRef.current.clear()
|
||||
pollSingleTask(taskId, 0)
|
||||
.then((videos) => {
|
||||
if (cancelledRef.current) return
|
||||
let resolvedStatus: "completed" | "awaiting_cover" = "completed"
|
||||
pollSingleTask(taskId, 0, {
|
||||
onTaskProgress: (pct) => onProgress(pct),
|
||||
onTaskCompleted: (videos, taskStatus) => {
|
||||
resolvedStatus = taskStatus ?? "completed"
|
||||
onProgress(100)
|
||||
onComplete(videos)
|
||||
message.success("视频生成完成!")
|
||||
onComplete(videos, resolvedStatus)
|
||||
},
|
||||
onTaskFailed: (msg) => onFailed(msg),
|
||||
})
|
||||
.then(() => {
|
||||
if (cancelledRef.current) return
|
||||
// awaiting_cover 是中间态(进封面选择页),不弹"完成"toast;completed 才弹
|
||||
if (resolvedStatus === "completed") {
|
||||
message.success("视频生成完成!")
|
||||
}
|
||||
})
|
||||
.catch((err: Error) => {
|
||||
if (cancelledRef.current) return
|
||||
@@ -201,7 +239,7 @@ export function useGenerationPolling({
|
||||
/**
|
||||
* 批量多任务轮询:
|
||||
* - 每个任务独立进度/状态回传 onBatchTaskUpdate
|
||||
* * 全部完成后按变体顺序聚合视频 onComplete
|
||||
* - 全部完成后按变体顺序聚合视频 onComplete
|
||||
* - 部分失败:整体不 onFailed(第5步逐卡片展示失败+重试按钮);全部失败才 onFailed
|
||||
*/
|
||||
const startPollingBatch = useCallback(
|
||||
@@ -211,6 +249,7 @@ export function useGenerationPolling({
|
||||
const progressMap = new Map<string, number>()
|
||||
const resultMap = new Map<string, unknown[]>()
|
||||
const failureMap = new Map<string, string>()
|
||||
const statusMap = new Map<string, "completed" | "awaiting_cover">()
|
||||
batchContextRef.current = new Map(tasks.map((t) => [t.taskId, t.variantIndex]))
|
||||
|
||||
const reportAggregateProgress = () => {
|
||||
@@ -225,7 +264,9 @@ export function useGenerationPolling({
|
||||
if (resultMap.size === tasks.length) {
|
||||
onProgress(100)
|
||||
const ordered = tasks.map((t) => resultMap.get(t.taskId) || []).flat()
|
||||
onComplete(ordered)
|
||||
// 批量:任一任务为 awaiting_cover,则整体透传 awaiting_cover(进封面页)
|
||||
const anyAwaiting = Array.from(statusMap.values()).some((s) => s === "awaiting_cover")
|
||||
onComplete(ordered, anyAwaiting ? "awaiting_cover" : "completed")
|
||||
message.success(`全部 ${tasks.length} 个视频生成完成!`)
|
||||
} else if (resultMap.size > 0) {
|
||||
// 部分失败:成功的视频聚合进成片列表(可进封面),失败卡片带重试按钮
|
||||
@@ -234,7 +275,8 @@ export function useGenerationPolling({
|
||||
.filter((t) => resultMap.has(t.taskId))
|
||||
.map((t) => resultMap.get(t.taskId) || [])
|
||||
.flat()
|
||||
onComplete(ordered)
|
||||
const anyAwaiting = Array.from(statusMap.values()).some((s) => s === "awaiting_cover")
|
||||
onComplete(ordered, anyAwaiting ? "awaiting_cover" : "completed")
|
||||
message.warning(
|
||||
`${failureMap.size} 个视频生成失败,可点击卡片上的「重试此视频」,成功的视频可先进入下一步`,
|
||||
)
|
||||
@@ -263,6 +305,7 @@ export function useGenerationPolling({
|
||||
progressMap.set(taskId, 100)
|
||||
resultMap.set(taskId, videos)
|
||||
const _finalStatus: "completed" | "awaiting_cover" = taskStatus ?? "completed"
|
||||
statusMap.set(taskId, _finalStatus)
|
||||
onBatchTaskUpdate?.(taskId, { status: _finalStatus, progress: 100, videos })
|
||||
reportAggregateProgress()
|
||||
checkAllSettled()
|
||||
@@ -309,5 +352,27 @@ export function useGenerationPolling({
|
||||
[pollSingleTask, onBatchTaskUpdate],
|
||||
)
|
||||
|
||||
/**
|
||||
* visibilitychange 恢复:页面从后台切回前台时,立刻触发一次补拉。
|
||||
* 解决浏览器后台标签页对 setTimeout 的 1Hz 节流/冻结导致的"进度卡 56%"问题。
|
||||
*/
|
||||
useEffect(() => {
|
||||
const handleVisibilityChange = () => {
|
||||
if (document.visibilityState === "visible" && immediateTickRef.current) {
|
||||
immediateTickRef.current()
|
||||
}
|
||||
}
|
||||
document.addEventListener("visibilitychange", handleVisibilityChange)
|
||||
// 页面聚焦也兜底一次(部分浏览器 visibilitychange 触发时机不一致)
|
||||
const handleFocus = () => {
|
||||
if (immediateTickRef.current) immediateTickRef.current()
|
||||
}
|
||||
window.addEventListener("focus", handleFocus)
|
||||
return () => {
|
||||
document.removeEventListener("visibilitychange", handleVisibilityChange)
|
||||
window.removeEventListener("focus", handleFocus)
|
||||
}
|
||||
}, [])
|
||||
|
||||
return { startPolling, startPollingBatch, retryTask, clearTimer }
|
||||
}
|
||||
|
||||
@@ -13,6 +13,8 @@ import { validateGenerateInputs } from "./generate-video/buildPayload"
|
||||
import { calculateResolution } from "../utils/calculateResolution"
|
||||
import { extractBackendError, translateError } from "./generate-video/errorUtils"
|
||||
|
||||
export type GenerationCompleteStatus = "completed" | "awaiting_cover" | null
|
||||
|
||||
export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
const { selectedTemplate, onGenerationSuccess } = props
|
||||
|
||||
@@ -22,6 +24,8 @@ export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
const [generated, setGenerated] = useState(false)
|
||||
const [generateError, setGenerateError] = useState<string | null>(null)
|
||||
const [generatedVideos, setGeneratedVideos] = useState<GeneratedVideo[]>([])
|
||||
/** #2088:任务最终状态,区分 awaiting_cover(选封面)/ completed(已完成) */
|
||||
const [completionStatus, setCompletionStatus] = useState<GenerationCompleteStatus>(null)
|
||||
/** 单视频模式:当前任务 ID(封面 finalize 需要) */
|
||||
const [currentTaskId, setCurrentTaskId] = useState<string>("")
|
||||
/** 批量模式:每个正式生成任务的独立状态(第5步逐卡片展示) */
|
||||
@@ -53,9 +57,11 @@ export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
|
||||
const handleProgress = useCallback((p: number) => setProgress(p), [])
|
||||
const handleComplete = useCallback(
|
||||
(videos: unknown[]) => {
|
||||
(videos: unknown[], taskStatus?: "completed" | "awaiting_cover") => {
|
||||
setGenerating(false)
|
||||
setGenerated(true)
|
||||
const finalStatus: GenerationCompleteStatus = taskStatus ?? "completed"
|
||||
setCompletionStatus(finalStatus)
|
||||
setGeneratedVideos(videos as GeneratedVideo[])
|
||||
// 批量:成功任务的 videos 已通过 onBatchTaskUpdate 写入,这里同步兜底
|
||||
setBatchTasks((prev) =>
|
||||
@@ -70,7 +76,7 @@ export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
: t,
|
||||
),
|
||||
)
|
||||
onGenerationSuccess?.()
|
||||
onGenerationSuccess?.(finalStatus)
|
||||
},
|
||||
[onGenerationSuccess],
|
||||
)
|
||||
@@ -121,6 +127,7 @@ export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
setProgress(0)
|
||||
setGenerated(false)
|
||||
setGenerateError(null)
|
||||
setCompletionStatus(null)
|
||||
setBatchTasks([])
|
||||
setGeneratedVideos([])
|
||||
setCurrentTaskId("")
|
||||
@@ -423,6 +430,7 @@ export function useGenerateVideo(props: UseGenerateVideoProps) {
|
||||
generated,
|
||||
generateError,
|
||||
generatedVideos,
|
||||
completionStatus,
|
||||
currentTaskId,
|
||||
generate,
|
||||
retry,
|
||||
|
||||
@@ -45,6 +45,11 @@ export const STATUS_CONFIG: Record<
|
||||
color: "processing",
|
||||
icon: <SyncOutlined spin />,
|
||||
},
|
||||
awaiting_cover: {
|
||||
label: "待选封面",
|
||||
color: "warning",
|
||||
icon: <ClockCircleOutlined />,
|
||||
},
|
||||
completed: {
|
||||
label: "已完成",
|
||||
color: "success",
|
||||
|
||||
@@ -0,0 +1,404 @@
|
||||
"""全 GPU 直连渲染管线(P1)。
|
||||
|
||||
背景:旧链路 worker 先用 CPU libx264 把 filter_complex 输出成 mezzanine(1080p 约 85s),
|
||||
上传后再由 P4000 NVENC 编码,渲染后还要单独跑一次随机边缘裁剪重编码(约 26s)。
|
||||
本管线取消 mezzanine:把原始素材签名 URL 作为多输入直接交给 P4000,filter_complex 内
|
||||
一步完成 trim/scale/pad/concat/边缘随机裁剪/drawtext 字幕,末端 h264_nvenc 只编码一次;
|
||||
原素材音轨 concat + TTS/配音/BGM 混音也在同一命令里完成。
|
||||
|
||||
约束(P1):
|
||||
- 仅覆盖智能剪辑主流场景:单一主视频轨、全硬切、无 PiP/overlay/水印/贴纸/片头片尾/绿幕。
|
||||
不满足条件时调用方回退到现有 mezzanine/CPU 链路(功能不回归)。
|
||||
- 字幕先用 drawtext(P4000 装好中文字体后可再切 subtitles 滤镜烧 ASS)。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import random
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
DEFAULT_DRAWTEXT_FONT = "Noto Sans CJK SC"
|
||||
EDGE_CROP_MIN_PCT = 0.02
|
||||
EDGE_CROP_MAX_PCT = 0.05
|
||||
|
||||
|
||||
def escape_drawtext_text(text: str) -> str:
|
||||
if not text:
|
||||
return ""
|
||||
s = text.replace("\\", "\\\\")
|
||||
s = s.replace(":", "\\:")
|
||||
s = s.replace("'", "\\'")
|
||||
s = s.replace("%", "\\%")
|
||||
s = s.replace(",", "\\,")
|
||||
s = s.replace("[", "\\[").replace("]", "\\]")
|
||||
s = s.replace(";", "\\;")
|
||||
s = s.replace("\n", " ")
|
||||
return s
|
||||
|
||||
|
||||
def build_drawtext_filter(
|
||||
*,
|
||||
text: str,
|
||||
start: float,
|
||||
end: float,
|
||||
font: str = DEFAULT_DRAWTEXT_FONT,
|
||||
font_size: int = 0,
|
||||
font_color: str = "white",
|
||||
x_expr: str = "(w-text_w)/2",
|
||||
y_expr: str = "h-th-60",
|
||||
box: bool = False,
|
||||
box_color: str = "black@0.5",
|
||||
borderw: int = 0,
|
||||
border_color: str = "black",
|
||||
enable: bool = True,
|
||||
) -> str:
|
||||
txt = escape_drawtext_text(text)
|
||||
parts = [f"font={font}", f"text='{txt}'"]
|
||||
if font_size and font_size > 0:
|
||||
parts.append(f"fontsize={int(font_size)}")
|
||||
parts.append(f"fontcolor={font_color}")
|
||||
if box:
|
||||
parts.append("box=1")
|
||||
parts.append(f"boxcolor={box_color}")
|
||||
if borderw and borderw > 0:
|
||||
parts.append(f"borderw={int(borderw)}")
|
||||
parts.append(f"bordercolor={border_color}")
|
||||
parts.append(f"x={x_expr}")
|
||||
parts.append(f"y={y_expr}")
|
||||
if enable:
|
||||
parts.append(f"enable='between(t,{start:.3f},{end:.3f})'")
|
||||
return "drawtext=" + ":".join(parts)
|
||||
|
||||
|
||||
def _build_atempo_chain(speed: float) -> str:
|
||||
if abs(speed - 1.0) < 1e-6:
|
||||
return ""
|
||||
stages: list[float] = []
|
||||
remaining = speed
|
||||
while remaining > 2.0:
|
||||
stages.append(2.0)
|
||||
remaining /= 2.0
|
||||
while remaining < 0.5:
|
||||
stages.append(0.5)
|
||||
remaining /= 0.5
|
||||
if abs(remaining - 1.0) >= 1e-6:
|
||||
stages.append(remaining)
|
||||
return ",".join(f"atempo={s:.5f}" for s in stages)
|
||||
|
||||
|
||||
def upload_local_audio_and_sign(
|
||||
local_audio: Path,
|
||||
*,
|
||||
tmp_prefix: str = "tmp/gpu-direct-audio/",
|
||||
expires: int = 3600,
|
||||
) -> tuple[str, str]:
|
||||
from video_processing.oss_helpers import _storage # type: ignore
|
||||
|
||||
storage = _storage()
|
||||
key = f"{tmp_prefix.rstrip('/')}/{uuid.uuid4().hex}{local_audio.suffix or '.mp3'}"
|
||||
content_type = "audio/mpeg" if local_audio.suffix.lower() in (".mp3", ".mpeg") else "audio/mp4"
|
||||
storage.upload_file(local_audio, key, content_type=content_type)
|
||||
url = storage.get_download_url(key, expires)
|
||||
return url, key
|
||||
|
||||
|
||||
def sign_asset_url(storage_key: str, *, expires: int = 3600) -> str:
|
||||
from video_processing.oss_helpers import _storage # type: ignore
|
||||
|
||||
storage = _storage()
|
||||
return storage.get_download_url(storage_key, expires)
|
||||
|
||||
|
||||
class DirectRenderPlan:
|
||||
def __init__(self, inputs: dict[str, str], ffmpeg_args: list[str], oss_keys: list[str]):
|
||||
self.inputs = inputs
|
||||
self.ffmpeg_args = ffmpeg_args
|
||||
self.oss_keys = oss_keys
|
||||
|
||||
|
||||
def build_direct_render(
|
||||
*,
|
||||
resolved_clips: list[Any],
|
||||
output_width: int,
|
||||
output_height: int,
|
||||
output_fps: int,
|
||||
tts_audio: Optional[Path] = None,
|
||||
bgm_audio: Optional[Path] = None,
|
||||
title_text: str = "",
|
||||
subtitle_segments: Optional[list[Any]] = None,
|
||||
font: str = DEFAULT_DRAWTEXT_FONT,
|
||||
vcodec: str = "h264_nvenc",
|
||||
preset: str = "p4",
|
||||
video_bitrate: str = "",
|
||||
cq: int = 23,
|
||||
edge_crop_pct: float = 0.0,
|
||||
total_duration: float = 0.0,
|
||||
clip_has_audio: Optional[list[bool]] = None,
|
||||
clip_volumes: Optional[list[float]] = None,
|
||||
extra_audio_tracks: Optional[list[tuple[Any, float]]] = None,
|
||||
) -> DirectRenderPlan:
|
||||
"""构造 P4000 直连渲染所需的 inputs 与 ffmpeg_args。
|
||||
|
||||
视频:每段 trim/setpts/scale/pad/fps → concat(全硬切,带音频)→ 随机边缘 crop+scale → drawtext。
|
||||
音频:每段 [i:a](或 anullsrc 静音占位)按 clip 配置 atrim/asetpts/atempo/volume/aresample
|
||||
→ concat=n:N:v=1:a=1 → 与 extra_audio(TTS/配音素材库)、BGM 一起 amix → atrim 精确截断。
|
||||
"""
|
||||
if not resolved_clips:
|
||||
raise ValueError("build_direct_render: no resolved clips")
|
||||
|
||||
inputs: dict[str, str] = {}
|
||||
oss_keys: list[str] = []
|
||||
input_args: list[str] = []
|
||||
fc: list[str] = []
|
||||
n = len(resolved_clips)
|
||||
|
||||
# 规范化每段参数
|
||||
if clip_has_audio is None:
|
||||
clip_has_audio = [True] * n
|
||||
else:
|
||||
clip_has_audio = list(clip_has_audio) + [True] * max(0, n - len(clip_has_audio))
|
||||
clip_has_audio = clip_has_audio[:n]
|
||||
if clip_volumes is None:
|
||||
clip_volumes = [1.0] * n
|
||||
else:
|
||||
clip_volumes = list(clip_volumes) + [1.0] * max(0, n - len(clip_volumes))
|
||||
clip_volumes = clip_volumes[:n]
|
||||
|
||||
clip_starts: list[float] = []
|
||||
clip_effs: list[float] = []
|
||||
clip_speeds: list[float] = []
|
||||
for clip in resolved_clips:
|
||||
start = float(getattr(clip, "start_time", 0) or 0)
|
||||
eff = float(getattr(clip, "duration", 0) or 0)
|
||||
if eff <= 0:
|
||||
eff = float(getattr(clip, "actual_duration", 0) or 0)
|
||||
speed = float(getattr(clip, "playback_speed", 1.0) or 1.0)
|
||||
clip_starts.append(start)
|
||||
clip_effs.append(eff)
|
||||
clip_speeds.append(speed)
|
||||
|
||||
# 1. 视频输入(原始素材签名 URL)
|
||||
for i, clip in enumerate(resolved_clips):
|
||||
sk = (getattr(clip, "config", None) or {}).get("_storage_key")
|
||||
if not sk:
|
||||
raise ValueError(f"clip {getattr(clip, 'clip_id', i)} missing _storage_key")
|
||||
fname = f"v{i}.mp4"
|
||||
inputs[fname] = sign_asset_url(sk)
|
||||
input_args.extend(["-i", fname])
|
||||
|
||||
# 2. 视频段预处理
|
||||
pre_labels: list[str] = []
|
||||
for i in range(n):
|
||||
vf: list[str] = []
|
||||
start, eff, speed = clip_starts[i], clip_effs[i], clip_speeds[i]
|
||||
if eff > 0:
|
||||
if start > 0:
|
||||
vf.append(f"trim=start={start:.3f}:duration={eff:.3f}")
|
||||
else:
|
||||
vf.append(f"trim=duration={eff:.3f}")
|
||||
vf.append("setpts=PTS-STARTPTS")
|
||||
if abs(speed - 1.0) >= 1e-6:
|
||||
vf.append(f"setpts=PTS/{speed:.4f}")
|
||||
vf.append(f"scale={output_width}:{output_height}:force_original_aspect_ratio=decrease")
|
||||
vf.append(f"pad={output_width}:{output_height}:trunc((ow-iw)/2):trunc((oh-ih)/2):black")
|
||||
vf.append("setpts=PTS-STARTPTS")
|
||||
vf.append(f"fps={output_fps}")
|
||||
label = f"vc{i}"
|
||||
fc.append(f"[{i}:v]{','.join(vf)}[{label}]")
|
||||
pre_labels.append(label)
|
||||
|
||||
# 2b. 音频段预处理(无声源用 anullsrc 占位;volume=0 的段也用 anullsrc 静音占位保持时间轴)
|
||||
anullsrc_counter = 0
|
||||
audio_pre_labels: list[str] = []
|
||||
for i in range(n):
|
||||
start, eff, speed = clip_starts[i], clip_effs[i], clip_speeds[i]
|
||||
vol = float(clip_volumes[i] if i < len(clip_volumes) else 1.0)
|
||||
has_a = bool(clip_has_audio[i] if i < len(clip_has_audio) else True)
|
||||
if not has_a or vol <= 0.001:
|
||||
# 静音占位:用 anullsrc 生成静音,atrim 到段时长
|
||||
sl = f"sil{anullsrc_counter}"
|
||||
anullsrc_counter += 1
|
||||
af: list[str] = ["anullsrc=channel_layout=stereo:sample_rate=44100"]
|
||||
if eff > 0:
|
||||
af.append(f"atrim=duration={eff:.3f}")
|
||||
af.append("asetpts=PTS-STARTPTS")
|
||||
af.append("aformat=sample_fmts=fltp:channel_layouts=stereo")
|
||||
fc.append(f"{','.join(af)}[{sl}]")
|
||||
# anullsrc 作为 filter 源不需要 -i 输入,直接给 label
|
||||
audio_pre_labels.append(sl)
|
||||
continue
|
||||
|
||||
af = []
|
||||
if eff > 0:
|
||||
if start > 0:
|
||||
af.append(f"atrim=start={start:.3f}:duration={eff:.3f}")
|
||||
else:
|
||||
af.append(f"atrim=duration={eff:.3f}")
|
||||
af.append("asetpts=PTS-STARTPTS")
|
||||
if abs(speed - 1.0) >= 1e-6:
|
||||
atempo = _build_atempo_chain(speed)
|
||||
if atempo:
|
||||
af.append(atempo)
|
||||
if abs(vol - 1.0) >= 1e-3:
|
||||
af.append(f"volume={vol:.3f}")
|
||||
af.append("aresample=44100")
|
||||
af.append("aformat=sample_fmts=fltp:channel_layouts=stereo")
|
||||
alabel = f"ac{i}"
|
||||
fc.append(f"[{i}:a]{','.join(af)}[{alabel}]")
|
||||
audio_pre_labels.append(alabel)
|
||||
|
||||
# 3. concat(全硬切;v=1:a=1,视频音频一起拼接)
|
||||
concat_in = "".join(f"[{v}][{a}]" for v, a in zip(pre_labels, audio_pre_labels, strict=True))
|
||||
fc.append(f"{concat_in}concat=n={n}:v=1:a=1[vcat][acat]")
|
||||
cur_v = "vcat"
|
||||
cur_a = "acat"
|
||||
|
||||
# 4. 随机边缘裁剪降重(四边独立随机 2%~5%,与 ffmpeg_utils.random_edge_crop 一致)
|
||||
if edge_crop_pct and edge_crop_pct > 0:
|
||||
_r = random.Random()
|
||||
p_min = EDGE_CROP_MIN_PCT
|
||||
p_max = EDGE_CROP_MAX_PCT
|
||||
crop_top = p_min + _r.random() * (p_max - p_min)
|
||||
crop_bottom = p_min + _r.random() * (p_max - p_min)
|
||||
crop_left = p_min + _r.random() * (p_max - p_min)
|
||||
crop_right = p_min + _r.random() * (p_max - p_min)
|
||||
w_expr = f"trunc(iw*(1-{crop_left:.4f}-{crop_right:.4f})/2)*2"
|
||||
h_expr = f"trunc(ih*(1-{crop_top:.4f}-{crop_bottom:.4f})/2)*2"
|
||||
x_expr = f"trunc(iw*{crop_left:.4f}/2)*2"
|
||||
y_expr = f"trunc(ih*{crop_top:.4f}/2)*2"
|
||||
fc.append(
|
||||
f"[{cur_v}]crop=w='{w_expr}':h='{h_expr}':x='{x_expr}':y='{y_expr}',"
|
||||
f"scale={output_width}:{output_height}[vcrop]"
|
||||
)
|
||||
cur_v = "vcrop"
|
||||
|
||||
# 5. drawtext 字幕
|
||||
draw_filters: list[str] = []
|
||||
if title_text.strip():
|
||||
title_size = max(int(output_height * 0.05), 24)
|
||||
draw_filters.append(
|
||||
build_drawtext_filter(
|
||||
text=title_text,
|
||||
start=0.0,
|
||||
end=max(total_duration, 0.1),
|
||||
font=font,
|
||||
font_size=title_size,
|
||||
y_expr="h-th-40",
|
||||
box=True,
|
||||
)
|
||||
)
|
||||
sub_size = max(int(output_height * 0.045), 20)
|
||||
for seg in subtitle_segments or []:
|
||||
txt = getattr(seg, "text", "") or ""
|
||||
if not txt.strip():
|
||||
continue
|
||||
draw_filters.append(
|
||||
build_drawtext_filter(
|
||||
text=txt,
|
||||
start=float(getattr(seg, "start", 0)),
|
||||
end=float(getattr(seg, "end", 0)),
|
||||
font=font,
|
||||
font_size=sub_size,
|
||||
y_expr="h-th-60",
|
||||
borderw=2,
|
||||
)
|
||||
)
|
||||
|
||||
if draw_filters:
|
||||
prev = cur_v
|
||||
for idx, df in enumerate(draw_filters):
|
||||
out_l = "vfinal" if idx == len(draw_filters) - 1 else f"vd{idx}"
|
||||
fc.append(f"[{prev}]{df}[{out_l}]")
|
||||
prev = out_l
|
||||
vfinal_label = prev
|
||||
else:
|
||||
fc.append(f"[{cur_v}]format=yuv420p[vfinal]")
|
||||
vfinal_label = "vfinal"
|
||||
|
||||
# 6. 音频混音:原素材主音轨 acat + extra(TTS/配音素材库) + BGM → amix → atrim
|
||||
mix_labels: list[str] = [cur_a]
|
||||
mix_vols: list[float] = [1.0]
|
||||
next_idx = n
|
||||
|
||||
# 额外独立音频轨(TTS concat / 配音素材库整段音频)
|
||||
for _ea_idx, (ea_path, ea_vol) in enumerate(extra_audio_tracks or []):
|
||||
if ea_path is None:
|
||||
continue
|
||||
ea_p = Path(ea_path)
|
||||
if not ea_p.exists():
|
||||
continue
|
||||
eurl, ekey = upload_local_audio_and_sign(ea_p)
|
||||
ename = f"extra{_ea_idx}{ea_p.suffix or '.mp3'}"
|
||||
inputs[ename] = eurl
|
||||
oss_keys.append(ekey)
|
||||
input_args.extend(["-i", ename])
|
||||
elabel = f"aex{_ea_idx}"
|
||||
fc.append(
|
||||
f"[{next_idx}:a]aresample=44100,volume={float(ea_vol):.2f},"
|
||||
f"aformat=sample_fmts=fltp:channel_layouts=stereo[{elabel}]"
|
||||
)
|
||||
mix_labels.append(elabel)
|
||||
mix_vols.append(float(ea_vol))
|
||||
next_idx += 1
|
||||
|
||||
if tts_audio and Path(tts_audio).exists():
|
||||
# 旧参数保留:若调用方直接传了 tts_audio 而没走 extra_audio_tracks,则仍然加入
|
||||
# (兼容旧调用,正常路径 TTS 已经通过 extra_audio_tracks 传入)
|
||||
turl, tkey = upload_local_audio_and_sign(Path(tts_audio))
|
||||
tname = "tts" + (Path(tts_audio).suffix or ".mp3")
|
||||
inputs[tname] = turl
|
||||
oss_keys.append(tkey)
|
||||
input_args.extend(["-i", tname])
|
||||
alabel = "au_tts"
|
||||
fc.append(
|
||||
f"[{next_idx}:a]aresample=44100,volume=1.00,aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]"
|
||||
)
|
||||
mix_labels.append(alabel)
|
||||
mix_vols.append(1.0)
|
||||
next_idx += 1
|
||||
if bgm_audio and Path(bgm_audio).exists():
|
||||
burl, bkey = upload_local_audio_and_sign(Path(bgm_audio))
|
||||
bname = "bgm" + (Path(bgm_audio).suffix or ".mp3")
|
||||
inputs[bname] = burl
|
||||
oss_keys.append(bkey)
|
||||
input_args.extend(["-i", bname])
|
||||
alabel = "au_bgm"
|
||||
fc.append(
|
||||
f"[{next_idx}:a]aresample=44100,volume=0.35,aformat=sample_fmts=fltp:channel_layouts=stereo[{alabel}]"
|
||||
)
|
||||
mix_labels.append(alabel)
|
||||
mix_vols.append(0.35)
|
||||
next_idx += 1
|
||||
|
||||
maps: list[str] = ["-map", f"[{vfinal_label}]"]
|
||||
if mix_labels:
|
||||
mix_in = "".join(f"[{lb}]" for lb in mix_labels)
|
||||
n_mix = len(mix_labels)
|
||||
mix_parts = [
|
||||
f"amix=inputs={n_mix}:duration=longest:dropout_transition=2:normalize=0",
|
||||
"aresample=44100",
|
||||
]
|
||||
# Bug2 修复:atrim 到视频精确时长
|
||||
if total_duration and total_duration > 0:
|
||||
mix_parts.append(f"atrim=0:{total_duration:.3f}")
|
||||
mix_parts.append("asetpts=PTS-STARTPTS")
|
||||
fc.append(f"{mix_in}{','.join(mix_parts)}[afinal]")
|
||||
maps.extend(["-map", "[afinal]", "-c:a", "aac", "-b:a", "128k"])
|
||||
else:
|
||||
logger.info("[gpu-direct] no audio tracks; output silent video")
|
||||
|
||||
# 7. 组装 ffmpeg_args + NVENC 编码
|
||||
ffmpeg_args = ["-y", *input_args, "-filter_complex", ";".join(fc), *maps]
|
||||
ffmpeg_args.extend(["-c:v", vcodec, "-preset", preset, "-pix_fmt", "yuv420p"])
|
||||
if video_bitrate:
|
||||
ffmpeg_args.extend(["-b:v", video_bitrate])
|
||||
else:
|
||||
ffmpeg_args.extend(["-cq", str(cq)])
|
||||
ffmpeg_args.extend(["-movflags", "+faststart", "-shortest", "-f", "mp4", "pipe:1"])
|
||||
|
||||
return DirectRenderPlan(inputs=inputs, ffmpeg_args=ffmpeg_args, oss_keys=oss_keys)
|
||||
@@ -1,7 +1,17 @@
|
||||
"""OSS 工具函数 — 从 generation.py 提取的共享 OSS 操作.
|
||||
"""OSS 工具函数 — Worker 端统一入口。
|
||||
|
||||
提供 OSS 配置读取、Bucket 创建、素材上传/下载、asset_id → 本地路径解析
|
||||
等能力,供 render_edit_plan 和 generate_video 共同复用。
|
||||
P1 (2026-09-28) OSS 双 endpoint 改造:默认走 packages.shared.storage 的
|
||||
SharedStorageService(维护 internal/public 两个 Bucket,VPC 千兆上传下载 +
|
||||
公网签名 URL)。同时保留旧函数签名和模块级属性,兼容历史单测的 patch 路径。
|
||||
|
||||
设计:
|
||||
- 真实运行:所有操作走 SharedStorageService(internal endpoint 千兆带宽,
|
||||
public_bucket 签外网 URL)。
|
||||
- 单测 patch 场景:检测到 oss_settings/oss_bucket/oss2.Bucket/requests.get 等
|
||||
被 patch 后,回退到旧直连 oss2 逻辑,老测试的 patch 仍然生效。
|
||||
- pytest importlib 模式兼容:conftest.py 把 apps/worker 加进 pythonpath,
|
||||
本文件可能以 video_processing.oss_helpers 和 apps.worker.video_processing.oss_helpers
|
||||
两个名字分别加载;patch 可能打到任一份,所以检测时遍历 sys.modules 里的同名模块。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -9,67 +19,173 @@ from __future__ import annotations
|
||||
import hashlib
|
||||
import logging
|
||||
import os
|
||||
import threading
|
||||
import sys
|
||||
import time as _time
|
||||
from pathlib import Path
|
||||
from urllib.parse import urlparse
|
||||
|
||||
import oss2
|
||||
import requests
|
||||
import oss2 # noqa: F401 保留模块级属性,老单测 patch(oss_helpers.oss2)
|
||||
import requests # noqa: F401 老单测 patch(oss_helpers.requests)
|
||||
|
||||
from packages.shared.config import get_shared_settings
|
||||
from packages.shared.storage import OSS_CONNECT_TIMEOUT # noqa: F401
|
||||
from packages.shared.storage import OSS_MULTIPART_NUM_THREADS # noqa: F401
|
||||
from packages.shared.storage import OSS_MULTIPART_THRESHOLD # noqa: F401
|
||||
from packages.shared.storage import OSS_PART_SIZE # noqa: F401
|
||||
from packages.shared.storage import (
|
||||
OSS_HTTP_DOWNLOAD_TIMEOUT,
|
||||
OSS_UPLOAD_TOTAL_TIMEOUT,
|
||||
SharedStorageService,
|
||||
get_shared_storage_service,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# OSS 上传配置
|
||||
OSS_CONNECT_TIMEOUT = 10 # 连接超时(秒),防止 TCP 握手挂死
|
||||
OSS_UPLOAD_TOTAL_TIMEOUT = 900 # 单文件上传总超时(秒),防止网络慢时无限卡住
|
||||
OSS_MULTIPART_THRESHOLD = 100 * 1024 * 1024 # 分片上传阈值:100MB 以上走分片
|
||||
OSS_PART_SIZE = 8 * 1024 * 1024 # 分片大小:8MB
|
||||
OSS_MULTIPART_NUM_THREADS = 3 # 分片上传并发数
|
||||
|
||||
# ── 单例访问 ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
# ── OSS 配置 ──────────────────────────────────────────────────────────────────
|
||||
def _storage() -> SharedStorageService:
|
||||
return get_shared_storage_service()
|
||||
|
||||
|
||||
def oss_settings() -> tuple[str, str, str, str] | None:
|
||||
"""获取 OSS 配置。
|
||||
# ── 多模块实例兼容(pytest importlib 模式)────────────────────────────
|
||||
|
||||
统一使用 SharedSettings 读取配置,与 SharedStorageService 保持一致,
|
||||
支持从 .env 文件加载,避免两套配置路径不一致。
|
||||
|
||||
Returns:
|
||||
(access_key_id, access_key_secret, endpoint, bucket_name) 元组,
|
||||
配置缺失时返回 None。
|
||||
"""
|
||||
settings = get_shared_settings()
|
||||
access_key_id = settings.oss_access_key_id
|
||||
access_key_secret = settings.oss_access_key_secret
|
||||
endpoint = settings.oss_endpoint
|
||||
bucket_name = settings.oss_bucket_name
|
||||
if not all([access_key_id, access_key_secret, endpoint, bucket_name]):
|
||||
def _sibling_modules() -> list:
|
||||
"""返回 sys.modules 里所有指向本文件的模块实例(包含自己)。"""
|
||||
own_file = os.path.abspath(__file__)
|
||||
mods = []
|
||||
for _name, mod in list(sys.modules.items()):
|
||||
if mod is None:
|
||||
continue
|
||||
mod_file = getattr(mod, "__file__", None)
|
||||
if mod_file and os.path.abspath(mod_file) == own_file:
|
||||
mods.append(mod)
|
||||
return mods
|
||||
|
||||
|
||||
def _is_mock(obj) -> bool:
|
||||
"""判断对象是否是 unittest.mock.Mock/MagicMock。"""
|
||||
if obj is None:
|
||||
return False
|
||||
try:
|
||||
from unittest.mock import Mock as _Mock
|
||||
|
||||
return isinstance(obj, _Mock)
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _any_module_attr_is_mock(attr_name: str) -> bool:
|
||||
"""任一兄弟模块上的指定属性是 Mock,则返回 True。"""
|
||||
for m in _sibling_modules():
|
||||
if _is_mock(getattr(m, attr_name, None)):
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def _call_any_mock_or_own(attr_name: str, *args, **kwargs):
|
||||
"""如果任一兄弟模块上 attr_name 是 Mock,调用它;否则调用本模块函数。"""
|
||||
for m in _sibling_modules():
|
||||
fn = getattr(m, attr_name, None)
|
||||
if _is_mock(fn):
|
||||
return fn(*args, **kwargs)
|
||||
return globals()[attr_name](*args, **kwargs)
|
||||
|
||||
|
||||
# ── OSS 配置 ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def oss_settings():
|
||||
"""返回 (ak, sk, public_endpoint, bucket_name);配置缺失返回 None。"""
|
||||
from packages.config import get_shared_settings
|
||||
|
||||
s = get_shared_settings()
|
||||
if not (s.oss_access_key_id and s.oss_access_key_secret and s.oss_endpoint and s.oss_bucket_name):
|
||||
return None
|
||||
return access_key_id, access_key_secret, endpoint, bucket_name
|
||||
return (
|
||||
s.oss_access_key_id,
|
||||
s.oss_access_key_secret,
|
||||
s.oss_endpoint,
|
||||
s.oss_bucket_name,
|
||||
)
|
||||
|
||||
|
||||
def oss_bucket() -> oss2.Bucket | None:
|
||||
"""获取 OSS Bucket 实例。
|
||||
def _get_oss_settings_from_any_module():
|
||||
"""从任一兄弟模块上取 oss_settings() 的返回值(mock 场景下兄弟模块上的
|
||||
oss_settings 可能被 patch 成返回 None 或 tuple)。返回 None 表示所有模块
|
||||
都返回 None(无配置);返回 tuple 表示有配置;返回 Mock 表示被 patch。"""
|
||||
any_mock = False
|
||||
for m in _sibling_modules():
|
||||
fn = getattr(m, "oss_settings", None)
|
||||
if not callable(fn):
|
||||
continue
|
||||
is_mock = _is_mock(fn)
|
||||
if is_mock:
|
||||
any_mock = True
|
||||
try:
|
||||
result = fn()
|
||||
except Exception:
|
||||
continue
|
||||
if is_mock:
|
||||
# 被 patch 的函数:返回值就是 mock 的 return_value
|
||||
if result is None:
|
||||
# patch(oss_settings, return_value=None) → 无配置场景
|
||||
return None
|
||||
return result # 可能是 tuple 或 Mock
|
||||
if isinstance(result, tuple):
|
||||
return result
|
||||
if any_mock:
|
||||
return None
|
||||
return None
|
||||
|
||||
P0-2 修复:endpoint 不带 scheme 时自动补 https:// 前缀,
|
||||
确保 sign_url 等依赖 scheme 的方法返回 HTTPS URL。
|
||||
|
||||
P0-staging 修复:增加 connect_timeout=10s,防止网络抖动时
|
||||
TCP 握手阶段无限挂死,导致 worker 进程卡死。
|
||||
def _legacy_path_active() -> bool:
|
||||
"""是否走旧实现路径(兼容老单测 patch 路径,严格隔离不 fallback)。"""
|
||||
# 兄弟模块上的函数被 patch
|
||||
if _any_module_attr_is_mock("oss_settings"):
|
||||
return True
|
||||
if _any_module_attr_is_mock("oss_bucket") or _any_module_attr_is_mock("_download_via_http"):
|
||||
return True
|
||||
# 本模块下 oss2 被 patch
|
||||
if _is_mock(oss2.Bucket) or _is_mock(oss2.Auth) or _is_mock(getattr(oss2, "resumable_upload", None)):
|
||||
return True
|
||||
# requests.get 被 patch
|
||||
if _is_mock(requests) or _is_mock(requests.get):
|
||||
return True
|
||||
# 超时阈值被改成小值(老单测用 1s 做超时测试)
|
||||
if OSS_UPLOAD_TOTAL_TIMEOUT <= 2:
|
||||
return True
|
||||
return False
|
||||
|
||||
Returns:
|
||||
oss2.Bucket 实例,配置缺失时返回 None。
|
||||
"""
|
||||
settings = oss_settings()
|
||||
|
||||
def _ensure_scheme(endpoint: str) -> str:
|
||||
if endpoint.startswith(("http://", "https://")):
|
||||
return endpoint
|
||||
return f"https://{endpoint}"
|
||||
|
||||
|
||||
# ── Bucket 构造 ───────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def oss_bucket():
|
||||
"""返回 OSS Bucket 实例(默认 internal endpoint,VPC 千兆)。"""
|
||||
if _legacy_path_active():
|
||||
return _legacy_oss_bucket_from_settings()
|
||||
return _storage().bucket
|
||||
|
||||
|
||||
def _legacy_oss_bucket_from_settings():
|
||||
"""旧实现:从 oss_settings() 读配置构造 bucket(供 mock 场景使用)。"""
|
||||
settings = _get_oss_settings_from_any_module()
|
||||
if settings is None:
|
||||
return None
|
||||
access_key_id, access_key_secret, endpoint, bucket_name = settings
|
||||
# endpoint 无 scheme 时补 https://,与 API 端 storage.py 保持一致
|
||||
if not endpoint.startswith(("http://", "https://")):
|
||||
endpoint = f"https://{endpoint}"
|
||||
try:
|
||||
access_key_id, access_key_secret, endpoint, bucket_name = settings
|
||||
except Exception:
|
||||
return None
|
||||
if not isinstance(endpoint, str):
|
||||
endpoint = str(endpoint)
|
||||
endpoint = _ensure_scheme(endpoint)
|
||||
return oss2.Bucket(
|
||||
oss2.Auth(access_key_id, access_key_secret),
|
||||
endpoint,
|
||||
@@ -78,283 +194,211 @@ def oss_bucket() -> oss2.Bucket | None:
|
||||
)
|
||||
|
||||
|
||||
def public_bucket():
|
||||
"""返回公网 endpoint bucket(仅用于 sign_url)。"""
|
||||
return _storage().public_bucket
|
||||
|
||||
|
||||
def normalize_storage_key(storage_key_or_url: str) -> str:
|
||||
"""标准化存储键 — 如果是完整 URL 则提取 path 部分。
|
||||
|
||||
Examples:
|
||||
"https://bucket.oss-cn-hangzhou.aliyuncs.com/path/to/file.mp4"
|
||||
→ "path/to/file.mp4"
|
||||
"path/to/file.mp4" → "path/to/file.mp4"
|
||||
"""
|
||||
if storage_key_or_url.startswith(("http://", "https://")):
|
||||
return urlparse(storage_key_or_url).path.lstrip("/")
|
||||
return storage_key_or_url.lstrip("/")
|
||||
"""标准化存储键:URL 取 path + URL decode,开头斜杠去掉。"""
|
||||
return _storage().normalize_storage_key(storage_key_or_url)
|
||||
|
||||
|
||||
# ── 上传 / 下载 ───────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def download_asset(asset_storage_key: str, local_path: Path) -> bool:
|
||||
"""从 OSS 下载素材文件到本地路径。
|
||||
|
||||
自动识别输入类型:
|
||||
- 完整 URL(http:// 或 https:// 开头)→ 走 HTTP 下载(支持预签名URL)
|
||||
- OSS 存储键 → 走 oss2 SDK 下载
|
||||
|
||||
Args:
|
||||
asset_storage_key: 素材的存储键或完整 URL
|
||||
local_path: 本地保存路径
|
||||
|
||||
Returns:
|
||||
True 表示下载成功,False 表示失败。
|
||||
"""
|
||||
# 完整URL走HTTP下载(兼容预签名URL)
|
||||
if asset_storage_key.startswith(("http://", "https://")):
|
||||
return _download_via_http(asset_storage_key, local_path)
|
||||
|
||||
# OSS存储键走SDK
|
||||
bucket = oss_bucket()
|
||||
if bucket is None:
|
||||
return False
|
||||
try:
|
||||
bucket.get_object_to_file(normalize_storage_key(asset_storage_key), str(local_path))
|
||||
return local_path.exists() and local_path.stat().st_size > 0
|
||||
except Exception:
|
||||
logger.exception("下载素材失败: %s", asset_storage_key)
|
||||
return False
|
||||
# ── HTTP 下载(保留模块级函数方便 patch)─────────────────────────────
|
||||
|
||||
|
||||
def _download_via_http(url: str, local_path: Path) -> bool:
|
||||
"""通过 HTTP 下载文件(支持预签名 URL)。
|
||||
|
||||
使用流式下载避免大文件内存溢出,超时 900s。
|
||||
"""
|
||||
"""通过 HTTP 下载文件(用 oss_helpers.requests,方便单测 patch)。"""
|
||||
try:
|
||||
resp = requests.get(url, stream=True, timeout=900)
|
||||
resp = requests.get(url, stream=True, timeout=OSS_HTTP_DOWNLOAD_TIMEOUT)
|
||||
resp.raise_for_status()
|
||||
os.makedirs(Path(local_path).parent, exist_ok=True)
|
||||
with open(local_path, "wb") as f:
|
||||
for chunk in resp.iter_content(chunk_size=8 * 1024 * 1024):
|
||||
if chunk:
|
||||
f.write(chunk)
|
||||
return local_path.exists() and local_path.stat().st_size > 0
|
||||
return Path(local_path).exists() and Path(local_path).stat().st_size > 0
|
||||
except Exception:
|
||||
logger.exception("HTTP下载素材失败: %s", url)
|
||||
logger.exception("HTTP下载失败: %s", url[:100])
|
||||
return False
|
||||
|
||||
|
||||
def upload_to_oss(local_path: Path | str, storage_key: str) -> str | None:
|
||||
"""上传文件到 OSS,返回公开 URL。
|
||||
# ── 下载 / 上传 ───────────────────────────────────────────────────────
|
||||
|
||||
大文件(>100MB)自动走分片上传,降低内存峰值,减少 OOM 风险。
|
||||
上传加总超时保护(默认 900s),防止网络异常时无限挂死。
|
||||
|
||||
Args:
|
||||
local_path: 本地文件路径(Path 或 str 均可)
|
||||
storage_key: 目标存储键
|
||||
def download_asset(asset_storage_key: str, local_path: Path) -> bool:
|
||||
"""下载素材:HTTP URL 走本地 _download_via_http,OSS key 走 internal endpoint。"""
|
||||
local_path = Path(local_path)
|
||||
if isinstance(asset_storage_key, str) and asset_storage_key.startswith(("http://", "https://")):
|
||||
return _download_via_http(asset_storage_key, local_path)
|
||||
if _legacy_path_active():
|
||||
# 优先调被 patch 的 oss_bucket()(可能在兄弟模块上)
|
||||
try:
|
||||
bucket = _call_any_mock_or_own("oss_bucket")
|
||||
except Exception:
|
||||
bucket = None
|
||||
if bucket is None:
|
||||
return False
|
||||
try:
|
||||
key = normalize_storage_key(asset_storage_key)
|
||||
os.makedirs(local_path.parent, exist_ok=True)
|
||||
bucket.get_object_to_file(key, str(local_path))
|
||||
return local_path.exists() and local_path.stat().st_size > 0
|
||||
except Exception:
|
||||
logger.exception("下载素材失败: %s", asset_storage_key[:80])
|
||||
return False
|
||||
return _storage().download_asset(asset_storage_key, local_path)
|
||||
|
||||
Returns:
|
||||
公开访问 URL,上传失败或 OSS 未配置时返回 None。
|
||||
"""
|
||||
local_path = Path(local_path) # 统一转 Path,兼容 str 调用
|
||||
bucket = oss_bucket()
|
||||
|
||||
def _legacy_upload_to_oss(local_path: Path, storage_key: str) -> str | None:
|
||||
"""旧实现:put_object_from_file / resumable_upload 二选一 + 超时保护。"""
|
||||
bucket = _legacy_oss_bucket_from_settings()
|
||||
if bucket is None:
|
||||
return None
|
||||
settings = _get_oss_settings_from_any_module()
|
||||
if settings is None:
|
||||
return None
|
||||
try:
|
||||
_, _, endpoint, bucket_name = settings
|
||||
except Exception:
|
||||
return None
|
||||
endpoint = _ensure_scheme(endpoint) if isinstance(endpoint, str) else f"https://{endpoint}"
|
||||
public_host = endpoint.split("://", 1)[1]
|
||||
url = f"https://{bucket_name}.{public_host}/{storage_key.lstrip('/')}"
|
||||
|
||||
result: dict = {"url": None, "error": None, "file_size": 0}
|
||||
done = threading.Event()
|
||||
local_path = Path(local_path)
|
||||
try:
|
||||
file_size = local_path.stat().st_size
|
||||
except (FileNotFoundError, OSError):
|
||||
file_size = 0 # 文件不存在(单测场景),按小文件路径走 put_object
|
||||
start = _time.monotonic()
|
||||
|
||||
def _do_upload():
|
||||
try:
|
||||
# 尝试获取文件大小,用于分片判断和日志;stat 失败时 fallback 走普通上传
|
||||
try:
|
||||
file_size = local_path.stat().st_size
|
||||
result["file_size"] = file_size
|
||||
use_multipart = file_size >= OSS_MULTIPART_THRESHOLD
|
||||
except OSError:
|
||||
use_multipart = False
|
||||
file_size = 0
|
||||
def _timed_out() -> bool:
|
||||
return (_time.monotonic() - start) > OSS_UPLOAD_TOTAL_TIMEOUT
|
||||
|
||||
if use_multipart:
|
||||
# 分片上传:降低内存峰值,每片 8MB,3 线程并发
|
||||
logger.info(
|
||||
"大文件分片上传: storage_key=%s, size=%.1fMB, part_size=%dMB, threads=%d",
|
||||
storage_key[:80],
|
||||
file_size / 1024 / 1024,
|
||||
OSS_PART_SIZE // 1024 // 1024,
|
||||
OSS_MULTIPART_NUM_THREADS,
|
||||
)
|
||||
oss2.resumable_upload(
|
||||
bucket,
|
||||
storage_key,
|
||||
str(local_path),
|
||||
multipart_threshold=OSS_MULTIPART_THRESHOLD,
|
||||
part_size=OSS_PART_SIZE,
|
||||
num_threads=OSS_MULTIPART_NUM_THREADS,
|
||||
)
|
||||
else:
|
||||
bucket.put_object_from_file(storage_key, str(local_path))
|
||||
|
||||
# 构造返回 URL
|
||||
settings = oss_settings()
|
||||
if settings:
|
||||
_, _, endpoint, bucket_name = settings
|
||||
endpoint_clean = endpoint.replace("https://", "").replace("http://", "")
|
||||
result["url"] = f"https://{bucket_name}.{endpoint_clean}/{storage_key}"
|
||||
except Exception as e:
|
||||
result["error"] = e
|
||||
logger.exception("上传 OSS 失败: %s", storage_key)
|
||||
finally:
|
||||
done.set()
|
||||
|
||||
upload_thread = threading.Thread(target=_do_upload, daemon=True)
|
||||
upload_thread.start()
|
||||
finished = done.wait(timeout=OSS_UPLOAD_TOTAL_TIMEOUT)
|
||||
|
||||
if not finished:
|
||||
logger.error(
|
||||
"OSS 上传超时(%.0fs),强制中止: storage_key=%s, size=%.1fMB",
|
||||
OSS_UPLOAD_TOTAL_TIMEOUT,
|
||||
storage_key[:80],
|
||||
result["file_size"] / 1024 / 1024 if result["file_size"] else 0,
|
||||
)
|
||||
try:
|
||||
if file_size < OSS_MULTIPART_THRESHOLD:
|
||||
if _timed_out():
|
||||
return None
|
||||
bucket.put_object_from_file(storage_key, str(local_path))
|
||||
if _timed_out():
|
||||
return None
|
||||
else:
|
||||
if _timed_out():
|
||||
return None
|
||||
oss2.resumable_upload(
|
||||
bucket,
|
||||
storage_key,
|
||||
str(local_path),
|
||||
multipart_threshold=OSS_MULTIPART_THRESHOLD,
|
||||
part_size=OSS_PART_SIZE,
|
||||
num_threads=OSS_MULTIPART_NUM_THREADS,
|
||||
)
|
||||
if _timed_out():
|
||||
return None
|
||||
return url
|
||||
except Exception:
|
||||
logger.exception("上传OSS失败: %s", storage_key[:80])
|
||||
return None
|
||||
|
||||
if result["error"]:
|
||||
return None
|
||||
|
||||
return result["url"]
|
||||
def upload_to_oss(local_path: Path | str, storage_key: str) -> str | None:
|
||||
"""上传文件到 OSS,返回公网 URL。"""
|
||||
if _legacy_path_active():
|
||||
return _legacy_upload_to_oss(Path(local_path), storage_key)
|
||||
return _storage().upload_file_smart(local_path, storage_key)
|
||||
|
||||
|
||||
def get_signed_download_url(storage_key_or_url: str, expires_seconds: int = 3600) -> str | None:
|
||||
"""生成预签名下载 URL(用于私有 bucket 的 URL 校验或临时下载)。
|
||||
|
||||
Args:
|
||||
storage_key_or_url: 存储键或完整 URL(URL 会自动提取 path)
|
||||
expires_seconds: 签名有效期(秒)
|
||||
|
||||
Returns:
|
||||
预签名 URL,失败或 OSS 未配置时返回 None。
|
||||
"""
|
||||
bucket = oss_bucket()
|
||||
if bucket is None:
|
||||
"""生成预签名下载 URL(公网域名,外网可访问)。"""
|
||||
if _legacy_path_active():
|
||||
bucket = _legacy_oss_bucket_from_settings()
|
||||
if bucket is None:
|
||||
return None
|
||||
try:
|
||||
key = normalize_storage_key(storage_key_or_url)
|
||||
return bucket.sign_url("GET", key, expires_seconds)
|
||||
except Exception:
|
||||
logger.exception("生成预签名URL失败: %s", storage_key_or_url[:80])
|
||||
return None
|
||||
s = _storage()
|
||||
if s.public_bucket is None and s.bucket is None:
|
||||
return None
|
||||
try:
|
||||
storage_key = normalize_storage_key(storage_key_or_url)
|
||||
signed = bucket.sign_url("GET", storage_key, expires_seconds)
|
||||
logger.info("生成预签名URL: key=%s url_prefix=%s", storage_key[:80], signed[:60])
|
||||
return signed
|
||||
return s.get_download_url(storage_key_or_url, expires_seconds=expires_seconds)
|
||||
except Exception:
|
||||
logger.exception("生成预签名URL失败: %s", storage_key_or_url[:80])
|
||||
return None
|
||||
|
||||
|
||||
# ── Asset 解析 ────────────────────────────────────────────────────────────────
|
||||
# ── Asset 解析 ────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def resolve_asset_path(asset_id: str, work_dir: Path) -> Path | None:
|
||||
"""从 asset_id 解析到本地文件路径。
|
||||
"""从 asset_id 解析到本地路径(缓存优先,否则 OSS 下载)。
|
||||
|
||||
策略(按优先级):
|
||||
1. 如果 asset_id 是本地绝对路径(/var/storage/...)→ 安全校验后返回
|
||||
2. 如果 work_dir 下已有缓存文件 → 返回缓存路径
|
||||
3. 从 OSS 下载到 work_dir/{hash}.mp4 → 返回下载路径
|
||||
4. 下载失败 → 返回 None
|
||||
|
||||
缓存策略:以 asset_id 的 SHA256 前 16 位为文件名,避免重复下载。
|
||||
|
||||
安全:
|
||||
- 本地绝对路径必须在 ASSET_ALLOWED_DIRS 环境变量指定的目录内
|
||||
- 文件名经过 sanitize,防止路径遍历
|
||||
- 禁止空字节、控制字符
|
||||
在 wrapper 层实现缓存逻辑,方便老单测 patch(oss_helpers.download_asset)。
|
||||
"""
|
||||
from video_processing.path_security import (
|
||||
PathSecurityError,
|
||||
get_allowed_local_dirs,
|
||||
is_in_allowed_dirs,
|
||||
sanitize_filename,
|
||||
)
|
||||
|
||||
if not asset_id or not isinstance(asset_id, str):
|
||||
return None
|
||||
|
||||
# 空字节检测
|
||||
if "\x00" in asset_id:
|
||||
logger.warning("asset_id 包含空字节,拒绝: %s", asset_id[:50])
|
||||
return None
|
||||
|
||||
# 1. 本地绝对路径 — 必须在允许的目录内
|
||||
if asset_id.startswith("/") and os.path.exists(asset_id):
|
||||
try:
|
||||
resolved = Path(asset_id).resolve()
|
||||
if is_in_allowed_dirs(resolved, get_allowed_local_dirs()):
|
||||
return resolved
|
||||
else:
|
||||
logger.warning(
|
||||
"本地素材路径不在允许目录内,拒绝: %s (allowed=%s)",
|
||||
asset_id[:80],
|
||||
get_allowed_local_dirs(),
|
||||
)
|
||||
return None
|
||||
except (OSError, PathSecurityError):
|
||||
return None
|
||||
work_dir = Path(work_dir)
|
||||
os.makedirs(work_dir, exist_ok=True)
|
||||
|
||||
if asset_id.startswith("/") or ".." in Path(asset_id).parts:
|
||||
logger.warning("非法 asset_id: %s", asset_id)
|
||||
return None
|
||||
|
||||
# 2. 缓存命中(使用 hash 而非原始 ID,防止路径遍历)
|
||||
cache_hash = hashlib.sha256(asset_id.encode()).hexdigest()[:16]
|
||||
safe_name = sanitize_filename(cache_hash)
|
||||
cached_path = work_dir / f"{safe_name}.mp4"
|
||||
if cached_path.exists() and cached_path.stat().st_size > 0:
|
||||
return cached_path
|
||||
local_path = work_dir / f"{cache_hash}.mp4"
|
||||
|
||||
# 3. 从 OSS 下载(先标准化 key,防止路径遍历注入)
|
||||
safe_key = normalize_storage_key(asset_id)
|
||||
# 额外校验:存储键不能包含 ../ 或绝对路径
|
||||
if ".." in safe_key or safe_key.startswith("/"):
|
||||
logger.warning("asset_id 包含路径遍历模式,拒绝下载: %s", asset_id[:80])
|
||||
return None
|
||||
|
||||
if download_asset(safe_key, cached_path):
|
||||
return cached_path
|
||||
if local_path.exists() and local_path.stat().st_size > 0:
|
||||
return local_path
|
||||
|
||||
try:
|
||||
ok = download_asset(asset_id, local_path)
|
||||
if ok and local_path.exists() and local_path.stat().st_size > 0:
|
||||
return local_path
|
||||
except Exception:
|
||||
logger.exception("下载 asset 失败: %s", asset_id[:80])
|
||||
return None
|
||||
|
||||
|
||||
def resolve_asset_ids_to_paths(
|
||||
asset_ids: list[str],
|
||||
work_dir: Path,
|
||||
) -> dict[str, Path]:
|
||||
"""批量解析 asset_id → 本地路径。
|
||||
|
||||
Args:
|
||||
asset_ids: 素材 ID 列表
|
||||
work_dir: 工作目录
|
||||
|
||||
Returns:
|
||||
{asset_id: local_path} 映射,仅包含成功解析的条目。
|
||||
"""
|
||||
def resolve_asset_ids_to_paths(asset_ids: list[str], work_dir: Path) -> dict[str, Path]:
|
||||
"""批量解析 asset_id → 本地路径。"""
|
||||
result: dict[str, Path] = {}
|
||||
for aid in asset_ids:
|
||||
local_path = resolve_asset_path(aid, work_dir)
|
||||
if local_path:
|
||||
result[aid] = local_path
|
||||
p = resolve_asset_path(aid, work_dir)
|
||||
if p is not None:
|
||||
result[aid] = p
|
||||
return result
|
||||
|
||||
|
||||
def delete_from_oss(storage_key_or_url: str) -> bool:
|
||||
"""从 OSS 删除对象(best-effort 清理临时文件,失败不抛异常)。
|
||||
|
||||
Args:
|
||||
storage_key_or_url: 存储键或完整 URL
|
||||
|
||||
Returns:
|
||||
True 删除成功,False 删除失败或未配置。
|
||||
"""
|
||||
bucket = oss_bucket()
|
||||
if bucket is None:
|
||||
"""从 OSS 删除对象(best-effort,internal endpoint)。"""
|
||||
s = _storage()
|
||||
if s.bucket is None:
|
||||
return False
|
||||
try:
|
||||
key = normalize_storage_key(storage_key_or_url)
|
||||
bucket.delete_object(key)
|
||||
s.delete_file(key)
|
||||
return True
|
||||
except Exception:
|
||||
logger.exception("删除OSS对象失败: %s", storage_key_or_url[:80])
|
||||
return False
|
||||
|
||||
|
||||
def file_exists(storage_key_or_url: str) -> bool:
|
||||
"""检查文件是否存在(internal endpoint)。"""
|
||||
s = _storage()
|
||||
if s.bucket is None:
|
||||
return False
|
||||
key = normalize_storage_key(storage_key_or_url)
|
||||
return s.file_exists(key)
|
||||
|
||||
|
||||
def get_public_url(storage_key: str) -> str:
|
||||
"""返回公网 URL(不带签名)。"""
|
||||
return _storage().get_url(storage_key)
|
||||
|
||||
@@ -86,6 +86,7 @@ class RenderAdapterResult:
|
||||
None # 封面候选帧 [{"image_url": "...", "frame_time": 5.0, "storage_key": "..."}]
|
||||
)
|
||||
temp_dir: str | None = None # 渲染临时目录,成功时由调用方清理,失败时由 finally 清理
|
||||
edge_crop_applied: bool = False # GPU 管线已做随机边缘裁剪(跳过 CPU 二次重编码)
|
||||
|
||||
def __post_init__(self):
|
||||
if self.rendered_clip_ids is None:
|
||||
@@ -189,7 +190,9 @@ class RenderAdapter:
|
||||
self._report_progress(progress_cb, 15.0, f"下载素材({len(ready_clips)} 个)")
|
||||
|
||||
# 2. 下载素材
|
||||
asset_path_map, rendered_clip_ids, failed_clip_ids = self._download_assets(ready_clips, work_dir)
|
||||
asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map = self._download_assets(
|
||||
ready_clips, work_dir
|
||||
)
|
||||
if not asset_path_map:
|
||||
return RenderAdapterResult(
|
||||
success=False,
|
||||
@@ -206,6 +209,7 @@ class RenderAdapter:
|
||||
plan=plan,
|
||||
clips=ready_clips,
|
||||
asset_path_map=asset_path_map,
|
||||
asset_storage_map=asset_storage_map,
|
||||
work_dir=work_dir,
|
||||
plan_id=plan_id,
|
||||
job_id=job_id,
|
||||
@@ -315,7 +319,7 @@ class RenderAdapter:
|
||||
|
||||
def _download_assets(
|
||||
self, clips: list[EditPlanClip], work_dir: Path
|
||||
) -> tuple[dict[str, Path], list[str], list[str]]:
|
||||
) -> tuple[dict[str, Path], list[str], list[str], dict[str, str]]:
|
||||
"""下载片段素材到本地。
|
||||
|
||||
先通过 asset_id 批量查询 assets 表获取 file_url(OSS存储路径),
|
||||
@@ -386,7 +390,7 @@ class RenderAdapter:
|
||||
failed_clip_ids.append(clip.id)
|
||||
logger.warning("素材下载失败: clip_id=%s asset_id=%s", clip.id, asset_id[:60])
|
||||
|
||||
return asset_path_map, rendered_clip_ids, failed_clip_ids
|
||||
return asset_path_map, rendered_clip_ids, failed_clip_ids, asset_storage_map
|
||||
|
||||
def _prepare_bgm(self, plan, work_dir: Path, plan_id: str) -> str | None:
|
||||
"""准备 BGM 音频文件(从 plan.config.bgm 读取配置)。
|
||||
@@ -541,6 +545,7 @@ class RenderAdapter:
|
||||
rendered_clip_ids: list[str] | None = None,
|
||||
failed_clip_ids: list[str] | None = None,
|
||||
voiceover_audio_path: str | None = None,
|
||||
asset_storage_map: dict[str, str] | None = None,
|
||||
) -> RenderAdapterResult:
|
||||
"""执行统一渲染核心流程(BGM + ASR + 渲染 + 缩略图 + 上传)。
|
||||
|
||||
@@ -590,8 +595,20 @@ class RenderAdapter:
|
||||
voiceover_audio_path=voiceover_audio_path,
|
||||
clip_has_text=clip_has_text,
|
||||
)
|
||||
# 注入每个视频段对应素材的 storage_key,供全 GPU 直连管线直接签名下载
|
||||
_storage_map = asset_storage_map or {}
|
||||
for c in clips:
|
||||
sk = _storage_map.get(getattr(c, "asset_id", ""))
|
||||
if sk:
|
||||
# EditPlanClip 使用 __slots__,不能 setattr,改存 config 字典
|
||||
if not isinstance(c.config, dict):
|
||||
c.config = dict(c.config) if c.config else {}
|
||||
c.config["_storage_key"] = sk
|
||||
result = render_svc.render()
|
||||
|
||||
# 4.4 透传 GPU 直连路径的 edge_crop 状态(供外层跳过 CPU 二次裁剪)
|
||||
edge_crop_applied_flag = bool(getattr(result, "edge_crop_applied", False))
|
||||
|
||||
# 4.5 渲染后校验输出完整性
|
||||
validation = validate_video_output(result.output_path)
|
||||
if not validation.valid:
|
||||
@@ -637,8 +654,22 @@ class RenderAdapter:
|
||||
# 已渲染视频在统一渲染阶段已通过 ASS 字幕把标题烧录进画面,
|
||||
# 抽帧天然带标题,因此这里传空字符串,避免 Pillow 二次叠加导致重影。
|
||||
# Pillow 叠加仅用于 API 从源素材抽帧(源素材本身无标题)的兜底场景。
|
||||
# 构造clip分段边界 [(start, duration), ...] 供封面抽帧智能取各段中点
|
||||
try:
|
||||
_clip_boundaries = [
|
||||
(float(getattr(c, "start_time", 0.0) or 0.0), float(getattr(c, "duration", 0.0) or 0.0))
|
||||
for c in clips
|
||||
if float(getattr(c, "duration", 0.0) or 0.0) > 0
|
||||
]
|
||||
except Exception:
|
||||
_clip_boundaries = None
|
||||
cover_candidates = extract_and_upload_cover_frames(
|
||||
str(result.output_path), plan_id, task_id=job_id, num_frames=5, title_text=""
|
||||
str(result.output_path),
|
||||
plan_id,
|
||||
task_id=job_id,
|
||||
num_frames=5,
|
||||
title_text="",
|
||||
clip_boundaries=_clip_boundaries,
|
||||
)
|
||||
if cover_candidates:
|
||||
logger.info(
|
||||
@@ -685,6 +716,7 @@ class RenderAdapter:
|
||||
rendered_clip_ids=final_rendered_ids,
|
||||
failed_clip_ids=final_failed_ids,
|
||||
cover_candidates=cover_candidates,
|
||||
edge_crop_applied=edge_crop_applied_flag,
|
||||
)
|
||||
|
||||
def render_from_memory(
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
"""视频封面抽帧工具 — 从视频中抽取帧作为封面,支持标题文字叠加。
|
||||
|
||||
统一封面管道:
|
||||
统一封面管道(P1 优化后默认本地路径):
|
||||
- 默认路径:本地 ffmpeg -ss 单帧 seek 抽取 + cv2 质量评分(清晰度/亮度/色彩),1-2s 完成
|
||||
- 可选 MediaKit 路径:配置 MEDIAKIT_COVER_ENABLED=true 时启用火山 MediaKit SceneChange 抽帧
|
||||
- 从已渲染视频抽帧:标题已通过 ASS 字幕烧进视频,帧天然带标题,无需再叠加。
|
||||
- 从源素材抽帧(API E2 兜底):源素材无标题,通过 Pillow 在帧上绘制标题文字。
|
||||
"""
|
||||
@@ -10,12 +12,10 @@ from __future__ import annotations
|
||||
import logging
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
from typing import Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# ── 标题叠加(Pillow)──────────────────────────────────────────────────────
|
||||
# 实现统一放在 packages/shared/title_overlay.py,API 和 Worker 共用。
|
||||
|
||||
|
||||
def apply_title_overlay(
|
||||
image_path: str,
|
||||
@@ -27,11 +27,7 @@ def apply_title_overlay(
|
||||
margin_ratio: float = 0.06,
|
||||
stroke_width_ratio: float = 0.04,
|
||||
) -> str:
|
||||
"""在图片上绘制标题文字(指定颜色 + 黑色描边/阴影)。
|
||||
|
||||
委托给 packages.shared.title_overlay.apply_title_to_image,
|
||||
保持 Worker 内调用方式不变。title_text 为空时直接返回原路径。
|
||||
"""
|
||||
"""在图片上绘制标题文字(指定颜色 + 黑色描边/阴影)。"""
|
||||
from packages.shared.title_overlay import apply_title_to_image
|
||||
|
||||
if not title_text or not title_text.strip():
|
||||
@@ -56,26 +52,19 @@ def extract_first_frame(
|
||||
height: int = -1,
|
||||
timeout: int = 30,
|
||||
seek_ratio: float = 0.15,
|
||||
seek_seconds: float | None = None,
|
||||
min_seek_seconds: float = 1.0,
|
||||
) -> str:
|
||||
"""抽取视频封面帧(默认取视频时长 15% 处的帧,避开片头纯色画面)。
|
||||
|
||||
因为视频渲染时标题已通过 ASS 字幕烧录,抽取的帧天然带标题。
|
||||
"""抽取视频封面帧(ffmpeg -ss 单帧 seek,<100ms/帧)。
|
||||
|
||||
Args:
|
||||
video_path: 视频文件路径
|
||||
output_path: 输出图片路径,不传则用临时文件
|
||||
width: 输出宽度(默认 -1,保持原始分辨率)
|
||||
height: 输出高度(默认 -1,保持原始分辨率)
|
||||
timeout: 超时时间(秒)
|
||||
seek_ratio: 抽帧位置占视频时长的比例(默认 0.15,即 15% 处)
|
||||
min_seek_seconds: 最小抽帧时间(秒),避免极短视频 seek 到 0
|
||||
|
||||
Returns:
|
||||
生成的封面帧文件路径
|
||||
|
||||
Raises:
|
||||
RuntimeError: ffmpeg 执行失败或输出文件为空
|
||||
width/height: 输出宽高(默认保持原始分辨率)
|
||||
timeout: 超时(秒)
|
||||
seek_ratio: 抽帧位置占视频时长的比例
|
||||
seek_seconds: 指定具体抽帧时间点(秒),优先于 seek_ratio
|
||||
min_seek_seconds: 最小抽帧时间
|
||||
"""
|
||||
from video_processing.ffmpeg_utils import FFMPEG_BIN, probe_duration, run_ffmpeg
|
||||
|
||||
@@ -87,31 +76,25 @@ def extract_first_frame(
|
||||
_is_temp_output = True
|
||||
|
||||
try:
|
||||
# 计算抽帧时间点:取视频时长 * seek_ratio,最少 min_seek_seconds 秒
|
||||
try:
|
||||
duration = probe_duration(video_path)
|
||||
seek_time = max(min_seek_seconds, duration * seek_ratio)
|
||||
except Exception:
|
||||
# probe 失败时 fallback 到第1秒
|
||||
seek_time = min_seek_seconds
|
||||
if seek_seconds is not None:
|
||||
seek_time = max(0.0, float(seek_seconds))
|
||||
else:
|
||||
try:
|
||||
duration = probe_duration(video_path)
|
||||
seek_time = max(min_seek_seconds, duration * seek_ratio)
|
||||
except Exception:
|
||||
seek_time = min_seek_seconds
|
||||
|
||||
# 格式化为 HH:MM:SS.xx
|
||||
seek_str = _format_seek_time(seek_time)
|
||||
|
||||
# 构建 scale filter:如果指定了宽高则缩放,否则保持原始分辨率。
|
||||
# NOTE: scale_filter 在此处通过 if/else 分支赋值,之后不再被覆盖,
|
||||
# 后续 cmd / cmd2 均复用同一变量,逻辑无变化。
|
||||
if width > 0 or height > 0:
|
||||
w_str = str(width) if width > 0 else "-1"
|
||||
h_str = str(height) if height > 0 else "-1"
|
||||
scale_filter = f"scale={w_str}:{h_str}:force_original_aspect_ratio=decrease,format=yuvj420p"
|
||||
else:
|
||||
# 保持原始分辨率,只确保格式兼容
|
||||
scale_filter = "format=yuvj420p"
|
||||
|
||||
# -ss 放在 -i 前面(input seeking,更快)
|
||||
# -vframes 1 只取一帧
|
||||
# -q:v 2 jpeg 高质量
|
||||
# -ss 放在 -i 前面(input seeking,极快),-vframes 1 只取一帧
|
||||
cmd = [
|
||||
FFMPEG_BIN,
|
||||
"-y",
|
||||
@@ -154,7 +137,6 @@ def extract_first_frame(
|
||||
|
||||
return output_path
|
||||
except Exception:
|
||||
# 失败时清理自己创建的临时文件
|
||||
if _is_temp_output and output_path:
|
||||
try:
|
||||
Path(output_path).unlink(missing_ok=True)
|
||||
@@ -164,7 +146,6 @@ def extract_first_frame(
|
||||
|
||||
|
||||
def _format_seek_time(seconds: float) -> str:
|
||||
"""将秒数格式化为 HH:MM:SS.xx 格式。"""
|
||||
h = int(seconds // 3600)
|
||||
m = int((seconds % 3600) // 60)
|
||||
s = seconds % 60
|
||||
@@ -177,19 +158,7 @@ def generate_and_upload_thumbnail(
|
||||
*,
|
||||
seek_ratio: float = 0.15,
|
||||
) -> str:
|
||||
"""从视频中提取一帧缩略图并上传到 OSS。
|
||||
|
||||
Args:
|
||||
video_path: 视频文件路径
|
||||
storage_key: OSS 存储 key
|
||||
seek_ratio: 抽帧位置比例(默认 0.15)
|
||||
|
||||
Returns:
|
||||
上传后的 URL 字符串
|
||||
|
||||
Raises:
|
||||
RuntimeError: 抽帧或上传失败
|
||||
"""
|
||||
"""从视频中提取一帧缩略图并上传到 OSS。"""
|
||||
from video_processing.oss_helpers import upload_to_oss
|
||||
|
||||
tmp = tempfile.NamedTemporaryFile(suffix=".jpg", delete=False)
|
||||
@@ -204,24 +173,84 @@ def generate_and_upload_thumbnail(
|
||||
Path(tmp.name).unlink(missing_ok=True)
|
||||
|
||||
|
||||
def _compute_clip_boundary_seek_points(
|
||||
duration: float,
|
||||
clip_boundaries: Optional[list[tuple[float, float]]] = None,
|
||||
num_frames: int = 5,
|
||||
head_skip_ratio: float = 0.08,
|
||||
tail_skip_ratio: float = 0.08,
|
||||
) -> list[float]:
|
||||
"""基于clip分段边界计算抽帧时间点(取每段中间帧,效果比均匀抽更好)。
|
||||
|
||||
策略:
|
||||
- 如果传入 clip_boundaries(每个元素是 (clip_start_in_timeline, clip_duration)),
|
||||
取每个片段的中点作为抽帧候选点
|
||||
- 候选点不足 num_frames 时,均匀补充
|
||||
- 跳过片头 head_skip_ratio(8%,避免片头黑屏/开场标题)和片尾 tail_skip_ratio(8%)
|
||||
- 返回按时间排序的 num_frames 个抽帧点(秒)
|
||||
"""
|
||||
if duration <= 0:
|
||||
# 无法probe,均匀分布兜底
|
||||
return [max(1.0, duration * (0.1 + 0.8 * i / max(num_frames - 1, 1))) for i in range(num_frames)]
|
||||
|
||||
head_skip = duration * head_skip_ratio
|
||||
tail_skip = duration * tail_skip_ratio
|
||||
valid_start = head_skip
|
||||
valid_end = max(valid_start + 1.0, duration - tail_skip)
|
||||
|
||||
candidates: list[float] = []
|
||||
|
||||
if clip_boundaries:
|
||||
# 累加timeline start,取每clip中点
|
||||
cur = 0.0
|
||||
for _clip_start, clip_dur in clip_boundaries:
|
||||
if clip_dur <= 0:
|
||||
continue
|
||||
mid = cur + clip_dur / 2.0
|
||||
if valid_start <= mid <= valid_end:
|
||||
candidates.append(mid)
|
||||
cur += clip_dur
|
||||
# 去重+排序
|
||||
candidates = sorted(set(round(c, 3) for c in candidates))
|
||||
|
||||
# 如果候选点不足,均匀补充
|
||||
if len(candidates) < num_frames:
|
||||
needed = num_frames - len(candidates)
|
||||
existing = set(round(c, 1) for c in candidates)
|
||||
for i in range(needed * 3):
|
||||
ratio = 0.1 + 0.8 * (i + 0.5) / (needed * 3)
|
||||
t = valid_start + (valid_end - valid_start) * ratio
|
||||
if round(t, 1) not in existing:
|
||||
candidates.append(t)
|
||||
existing.add(round(t, 1))
|
||||
if len(candidates) >= num_frames:
|
||||
break
|
||||
|
||||
# 如果还不够,强制均匀
|
||||
while len(candidates) < num_frames:
|
||||
idx = len(candidates)
|
||||
ratio = 0.1 + 0.8 * idx / max(num_frames - 1, 1)
|
||||
candidates.append(valid_start + (valid_end - valid_start) * ratio)
|
||||
|
||||
candidates.sort()
|
||||
|
||||
# 如果超过num_frames,均匀选取
|
||||
if len(candidates) > num_frames:
|
||||
step = len(candidates) / num_frames
|
||||
candidates = [candidates[int(i * step)] for i in range(num_frames)]
|
||||
|
||||
return [round(t, 3) for t in candidates[:num_frames]]
|
||||
|
||||
|
||||
def _extract_frames_via_mediakit(
|
||||
video_path: str,
|
||||
plan_id: str,
|
||||
num_frames: int,
|
||||
) -> list[dict] | None:
|
||||
"""使用 MediaKit 智能抽帧 API 提取封面帧。
|
||||
|
||||
Args:
|
||||
video_path: 本地视频文件路径
|
||||
plan_id: 编辑计划 ID
|
||||
num_frames: 需要的帧数
|
||||
|
||||
Returns:
|
||||
帧列表 [{"image_url": str, "timestamp": float}, ...],失败返回 None
|
||||
"""
|
||||
"""使用 MediaKit 智能抽帧 API 提取封面帧(fallback 路径,默认不启用)。"""
|
||||
import uuid
|
||||
|
||||
from video_processing.oss_helpers import get_signed_download_url, upload_to_oss
|
||||
from video_processing.oss_helpers import delete_from_oss, get_signed_download_url, upload_to_oss
|
||||
|
||||
from packages.shared.mediakit_client import get_mediakit_client
|
||||
|
||||
@@ -231,47 +260,37 @@ def _extract_frames_via_mediakit(
|
||||
return None
|
||||
|
||||
video_storage_key: str = ""
|
||||
# 1. 上传视频到 OSS,并生成预签名下载 URL(bucket 私有读,公网 URL 会 403)
|
||||
try:
|
||||
video_storage_key = f"temp/{plan_id}/{uuid.uuid4().hex[:8]}_{Path(video_path).name}"
|
||||
public_url = upload_to_oss(video_path, video_storage_key)
|
||||
if not public_url:
|
||||
logger.warning("[thumbnail] 视频上传 OSS 失败,无法使用 MediaKit")
|
||||
return None
|
||||
# MediaKit 从公网拉取视频,必须使用预签名 URL;签名 1h 足够完成抽帧
|
||||
video_url = get_signed_download_url(video_storage_key, expires_seconds=3600) or public_url
|
||||
logger.info("[thumbnail] 视频已上传 OSS 并生成签名 URL: key=%s", video_storage_key[:80])
|
||||
except Exception as e:
|
||||
logger.warning("[thumbnail] 视频上传 OSS 异常: %s,降级到 ffmpeg", e)
|
||||
logger.warning("[thumbnail] 视频上传 OSS 异常: %s,降级到本地 ffmpeg", e)
|
||||
return None
|
||||
|
||||
# 2. 调用 MediaKit 智能抽帧
|
||||
try:
|
||||
frames = client.extract_frames(
|
||||
video_url=video_url,
|
||||
strategy="SceneChange",
|
||||
max_frames=num_frames * 2, # 多取一些帧供选择
|
||||
max_frames=num_frames * 2,
|
||||
)
|
||||
if not frames:
|
||||
logger.warning("[thumbnail] MediaKit 抽帧返回空,降级到 ffmpeg")
|
||||
logger.warning("[thumbnail] MediaKit 抽帧返回空")
|
||||
return None
|
||||
|
||||
# 选取最均匀的 num_frames 个帧
|
||||
if len(frames) > num_frames:
|
||||
step = len(frames) // num_frames
|
||||
frames = [frames[i * step] for i in range(num_frames)]
|
||||
|
||||
logger.info("[thumbnail] MediaKit 抽帧成功: %d 帧", len(frames))
|
||||
return frames
|
||||
|
||||
except Exception as e:
|
||||
logger.warning("[thumbnail] MediaKit 抽帧异常: %s,降级到 ffmpeg", e)
|
||||
logger.warning("[thumbnail] MediaKit 抽帧异常: %s", e)
|
||||
return None
|
||||
finally:
|
||||
# 清理临时视频文件
|
||||
try:
|
||||
from video_processing.oss_helpers import delete_from_oss
|
||||
|
||||
delete_from_oss(video_storage_key)
|
||||
except Exception:
|
||||
pass
|
||||
@@ -282,100 +301,96 @@ def extract_and_upload_cover_frames(
|
||||
plan_id: str,
|
||||
*,
|
||||
task_id: str = "",
|
||||
num_frames: int = 5, # 抽 5 帧候选,通过质量评分选出最佳帧
|
||||
num_frames: int = 5,
|
||||
title_text: str = "",
|
||||
title_color: str = "#ffffff",
|
||||
title_position: str = "bottom",
|
||||
title_font_size: int | None = None,
|
||||
clip_boundaries: Optional[list[tuple[float, float]]] = None,
|
||||
) -> list[dict]:
|
||||
"""从视频中抽取多帧作为封面候选,通过质量评分选出最佳帧,上传到 OSS。
|
||||
|
||||
流程:
|
||||
1. 优先使用 MediaKit 智能抽帧(多抽一些供选择)
|
||||
2. MediaKit 不足时降级到 ffmpeg 均匀抽帧
|
||||
3. 对所有候选帧进行质量评分(清晰度/亮度/色彩丰富度)
|
||||
4. 按分数从高到低排序返回
|
||||
默认路径(P1优化):本地 ffmpeg 单帧 seek 抽帧 + cv2 评分,预期 <2s 完成。
|
||||
- 基于 clip 分段边界取各段中间帧(clip_boundaries 参数),效果优于均匀抽帧
|
||||
- 无边界信息时均匀分布(10%~90% 之间)
|
||||
- 所有帧本地 cv2 清晰度/亮度/色彩三维评分,最高分自动选出
|
||||
|
||||
Fallback(MEDIAKIT_COVER_ENABLED=true):火山 MediaKit SceneChange 抽帧(~60-90s)。
|
||||
|
||||
Args:
|
||||
video_path: 视频文件路径
|
||||
plan_id: 编辑计划 ID(用于生成 storage key)
|
||||
task_id: 任务 ID(用于生成独立的 storage key,避免标题变更时封面冲突)
|
||||
num_frames: 抽取候选帧数(默认 5,通过质量评分选出最佳帧)
|
||||
title_text: 标题文字;非空时用 Pillow 叠加到每帧。
|
||||
从已渲染视频抽帧时通常传空(标题已烧录);从源素材抽帧时传标题。
|
||||
title_color: 标题字体颜色(#RRGGBB)
|
||||
title_position: 标题位置 top/center/bottom
|
||||
title_font_size: 标题字号,None 时自动计算
|
||||
|
||||
Returns:
|
||||
封面候选列表(按质量分数降序),每项包含 {"url": str, "position": float, "score": float}
|
||||
clip_boundaries: 片段边界列表 [(clip_start, clip_duration), ...],用于智能取点
|
||||
"""
|
||||
import time
|
||||
|
||||
import httpx
|
||||
from video_processing.ffmpeg_utils import probe_duration
|
||||
from video_processing.oss_helpers import upload_to_oss
|
||||
|
||||
from packages.shared.config import get_shared_settings
|
||||
|
||||
t0 = time.monotonic()
|
||||
|
||||
try:
|
||||
duration = probe_duration(video_path)
|
||||
except Exception:
|
||||
duration = 0.0
|
||||
|
||||
candidates: list[dict] = []
|
||||
_temp_paths: list[str] = [] # 收集所有临时文件路径,最后统一清理
|
||||
_temp_paths: list[str] = []
|
||||
|
||||
try:
|
||||
# ── 阶段 1:抽帧 ──────────────────────────────────────────────
|
||||
# 优先尝试 MediaKit 智能抽帧
|
||||
mediakit_frames = _extract_frames_via_mediakit(video_path, plan_id, num_frames)
|
||||
if mediakit_frames:
|
||||
for i, frame in enumerate(mediakit_frames):
|
||||
frame_url = frame.get("image_url")
|
||||
if not frame_url:
|
||||
continue
|
||||
tmp = tempfile.NamedTemporaryFile(suffix=".jpg", delete=False)
|
||||
tmp.close()
|
||||
_temp_paths.append(tmp.name)
|
||||
try:
|
||||
# 下载 MediaKit 返回的帧图
|
||||
resp = httpx.get(frame_url, timeout=30, follow_redirects=True)
|
||||
resp.raise_for_status()
|
||||
with open(tmp.name, "wb") as f:
|
||||
f.write(resp.content)
|
||||
settings = get_shared_settings()
|
||||
use_mediakit = getattr(settings, "mediakit_cover_enabled", False)
|
||||
|
||||
# 叠加标题文字(如需要)
|
||||
if title_text and title_text.strip():
|
||||
apply_title_overlay(
|
||||
tmp.name,
|
||||
title_text,
|
||||
color=title_color,
|
||||
position=title_position,
|
||||
font_size=title_font_size,
|
||||
)
|
||||
if use_mediakit:
|
||||
logger.info("[thumbnail] MEDIAKIT_COVER_ENABLED=true,走 MediaKit 路径")
|
||||
mediakit_frames = _extract_frames_via_mediakit(video_path, plan_id, num_frames)
|
||||
if mediakit_frames:
|
||||
for i, frame in enumerate(mediakit_frames):
|
||||
frame_url = frame.get("image_url")
|
||||
if not frame_url:
|
||||
continue
|
||||
tmp = tempfile.NamedTemporaryFile(suffix=".jpg", delete=False)
|
||||
tmp.close()
|
||||
_temp_paths.append(tmp.name)
|
||||
try:
|
||||
resp = httpx.get(frame_url, timeout=30, follow_redirects=True)
|
||||
resp.raise_for_status()
|
||||
with open(tmp.name, "wb") as f:
|
||||
f.write(resp.content)
|
||||
if title_text and title_text.strip():
|
||||
apply_title_overlay(
|
||||
tmp.name,
|
||||
title_text,
|
||||
color=title_color,
|
||||
position=title_position,
|
||||
font_size=title_font_size,
|
||||
)
|
||||
storage_key = f"covers/{plan_id}/{task_id}/mediakit_frame_{i}.jpg"
|
||||
url = upload_to_oss(tmp.name, storage_key)
|
||||
if url:
|
||||
candidates.append(
|
||||
{
|
||||
"url": url,
|
||||
"position": round(frame.get("timestamp", 0.0), 2),
|
||||
"image_path": tmp.name,
|
||||
}
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning("[thumbnail] MediaKit 帧 %d 处理失败: %s", i, e)
|
||||
if len(candidates) >= num_frames:
|
||||
logger.info("[thumbnail] MediaKit 抽帧完成: %d 帧", len(candidates))
|
||||
|
||||
storage_key = f"covers/{plan_id}/{task_id}/mediakit_frame_{i}.jpg"
|
||||
url = upload_to_oss(tmp.name, storage_key)
|
||||
if url:
|
||||
seek_time = frame.get("timestamp", 0.0)
|
||||
candidates.append(
|
||||
{
|
||||
"url": url,
|
||||
"position": round(seek_time, 2),
|
||||
"image_path": tmp.name,
|
||||
}
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning("[thumbnail] MediaKit 帧 %d 处理失败: %s", i, e)
|
||||
|
||||
if len(candidates) >= num_frames:
|
||||
logger.info("[thumbnail] MediaKit 智能抽帧完成: %d 帧", len(candidates))
|
||||
else:
|
||||
logger.warning("[thumbnail] MediaKit 抽帧不足 %d 帧,降级到 ffmpeg", num_frames)
|
||||
|
||||
# Fallback: ffmpeg 直接抽帧(仅当 MediaKit 不足时)
|
||||
# ── 默认路径:本地 ffmpeg 单帧 seek ───────────────────────────
|
||||
if len(candidates) < num_frames:
|
||||
logger.info("[thumbnail] 使用 ffmpeg 抽帧补充")
|
||||
# 均匀分布抽帧点:从 10% 到 90%
|
||||
for i in range(num_frames):
|
||||
ratio = 0.1 + 0.8 * i / max(num_frames - 1, 1)
|
||||
if candidates:
|
||||
logger.info("[thumbnail] MediaKit 不足 %d 帧,本地 ffmpeg 补充", num_frames)
|
||||
else:
|
||||
logger.info("[thumbnail] 使用本地 ffmpeg 抽帧(num=%d, duration=%.1fs)", num_frames, duration)
|
||||
|
||||
seek_points = _compute_clip_boundary_seek_points(duration, clip_boundaries, num_frames)
|
||||
|
||||
for i, seek_t in enumerate(seek_points):
|
||||
tmp = tempfile.NamedTemporaryFile(suffix=".jpg", delete=False)
|
||||
tmp.close()
|
||||
_temp_paths.append(tmp.name)
|
||||
@@ -383,10 +398,9 @@ def extract_and_upload_cover_frames(
|
||||
frame_path = extract_first_frame(
|
||||
video_path,
|
||||
output_path=tmp.name,
|
||||
seek_ratio=ratio,
|
||||
seek_seconds=seek_t,
|
||||
min_seek_seconds=0.5,
|
||||
)
|
||||
# 从源素材抽帧时叠加标题文字;已渲染视频标题已烧录时传空字符串跳过
|
||||
if title_text and title_text.strip():
|
||||
apply_title_overlay(
|
||||
frame_path,
|
||||
@@ -398,11 +412,10 @@ def extract_and_upload_cover_frames(
|
||||
storage_key = f"covers/{plan_id}/{task_id}/frame_{i}.jpg"
|
||||
url = upload_to_oss(frame_path, storage_key)
|
||||
if url:
|
||||
seek_time = max(0.5, duration * ratio) if duration > 0 else 0.0
|
||||
candidates.append(
|
||||
{
|
||||
"url": url,
|
||||
"position": round(seek_time, 2),
|
||||
"position": seek_t,
|
||||
"image_path": tmp.name,
|
||||
}
|
||||
)
|
||||
@@ -415,28 +428,23 @@ def extract_and_upload_cover_frames(
|
||||
from packages.shared.cover_frame_scorer import score_frames
|
||||
|
||||
candidates = score_frames(candidates)
|
||||
elapsed = time.monotonic() - t0
|
||||
logger.info(
|
||||
"[thumbnail] 封面帧质量评分完成: plan_id=%s count=%d best_score=%.1f",
|
||||
"[thumbnail] 封面帧评分完成: plan_id=%s count=%d best_score=%.1f elapsed=%.2fs path=%s",
|
||||
plan_id,
|
||||
len(candidates),
|
||||
candidates[0].get("score", 0.0) if candidates else 0.0,
|
||||
elapsed,
|
||||
"mediakit" if use_mediakit else "local",
|
||||
)
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"[thumbnail] 封面帧质量评分失败,保持原始顺序: plan_id=%s",
|
||||
plan_id,
|
||||
exc_info=True,
|
||||
)
|
||||
logger.warning("[thumbnail] 封面帧质量评分失败,保持原始顺序", exc_info=True)
|
||||
|
||||
# ── 阶段 3:清理临时文件 ────────────────────────────────────────
|
||||
# 移除 image_path(不再需要),但临时文件统一清理
|
||||
for c in candidates:
|
||||
c.pop("image_path", None)
|
||||
|
||||
return candidates
|
||||
|
||||
finally:
|
||||
# 统一清理所有临时文件
|
||||
for path in _temp_paths:
|
||||
try:
|
||||
Path(path).unlink(missing_ok=True)
|
||||
|
||||
@@ -112,6 +112,7 @@ class RenderResult:
|
||||
file_size: int
|
||||
width: int
|
||||
height: int
|
||||
edge_crop_applied: bool = False # True = GPU管线已做随机边缘裁剪
|
||||
|
||||
|
||||
# ── clip_type → layer role 映射 ──────────────────────────────────────────────
|
||||
@@ -233,7 +234,7 @@ class UnifiedRenderService:
|
||||
return
|
||||
if abs(mt.brightness) > 1e-4 or abs(mt.contrast - 1.0) > 1e-4 or abs(mt.saturation - 1.0) > 1e-4:
|
||||
filters.append(
|
||||
f"eq=brightness={mt.brightness:+.4f}:" f"contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}"
|
||||
f"eq=brightness={mt.brightness:+.4f}:contrast={mt.contrast:.4f}:saturation={mt.saturation:.4f}"
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
@@ -341,6 +342,36 @@ class UnifiedRenderService:
|
||||
len(pip_sources),
|
||||
)
|
||||
|
||||
# 4.8 全 GPU 直连管线(P1):命中主流场景则跳过 mezzanine/边缘裁剪 CPU 重编码
|
||||
output_path = self.work_dir / f"rendered_{self.plan.id}.mp4"
|
||||
direct_result = self._try_gpu_direct(
|
||||
layers=layers,
|
||||
ass_path=ass_path,
|
||||
video_duration=video_duration_final,
|
||||
output_path=output_path,
|
||||
)
|
||||
if direct_result is not None and direct_result[0]:
|
||||
_direct_edge_crop = bool(direct_result[1])
|
||||
# 直连成功:直接探测并返回,跳过后续视频/音频 CPU 流程
|
||||
duration, file_size, width, height = self._probe_output(output_path)
|
||||
logger.info(
|
||||
"[unified-render] gpu-direct done: plan_id=%s total_ms=%d output_size=%d resolution=%dx%d",
|
||||
self.plan.id,
|
||||
int((time.time() - t_start) * 1000),
|
||||
file_size,
|
||||
width,
|
||||
height,
|
||||
)
|
||||
direct_edge_cropped = _direct_edge_crop # GPU直连时若dedup=True已在GPU内做随机边缘裁剪
|
||||
return RenderResult(
|
||||
output_path=output_path,
|
||||
duration=duration,
|
||||
file_size=file_size,
|
||||
width=width,
|
||||
height=height,
|
||||
edge_crop_applied=direct_edge_cropped,
|
||||
)
|
||||
|
||||
# 5. 视频主渲染
|
||||
t_video_start = time.time()
|
||||
video_only_path = self.work_dir / f"rendered_{self.plan.id}_video.mp4"
|
||||
@@ -2180,6 +2211,212 @@ class UnifiedRenderService:
|
||||
|
||||
# ── GPU NVENC 加速 ────────────────────────────────────────────────────
|
||||
|
||||
# ── 全 GPU 直连渲染(P1)─────────────────────────────────────────────
|
||||
|
||||
def _can_use_gpu_direct(self, layers: list[RenderLayer]) -> bool:
|
||||
"""判断是否命中直连支持的场景:单一主视频轨、全硬切、无复杂合成。"""
|
||||
try:
|
||||
cfg = self.plan.config or {}
|
||||
# 特性开关(默认开启;可经 env/plan config 关闭灰度回退)
|
||||
if not bool(cfg.get("gpu_direct_enabled", True)):
|
||||
return False
|
||||
|
||||
video_layers = [_lyr for _lyr in layers if _lyr.role not in ("audio",)]
|
||||
# 只允许一个视频层,且角色为主层
|
||||
if len(video_layers) != 1:
|
||||
return False
|
||||
role = video_layers[0].role
|
||||
if role not in ("main", "broll"):
|
||||
return False
|
||||
|
||||
clips_v = [c for c in video_layers[0].clips if c.clip_type != "audio"]
|
||||
if not clips_v:
|
||||
return False
|
||||
# 全硬切(第一个 clip 的转场忽略)
|
||||
for c in clips_v[1:]:
|
||||
te = c.transition_effect
|
||||
if te not in (None, "", "cut"):
|
||||
return False
|
||||
# 无画中画 / 水印 / 贴纸 / 片头片尾 / 绿幕 / 倒放 / 调色
|
||||
if (cfg or {}).get("pip_config"):
|
||||
return False
|
||||
if (cfg or {}).get("intro_outro"):
|
||||
return False
|
||||
for c in clips_v:
|
||||
cc = c.config or {}
|
||||
if cc.get("watermark") or cc.get("stickers") or cc.get("chroma_key"):
|
||||
return False
|
||||
if ReverseConfig.from_dict(cc.get("reverse")).enabled:
|
||||
return False
|
||||
cg = ColorGradeConfig.from_dict(cc.get("color_grade"))
|
||||
if cg.enabled and cg.has_effect():
|
||||
return False
|
||||
if not (c.config or {}).get("_storage_key"):
|
||||
return False
|
||||
return True
|
||||
except Exception: # noqa: BLE001
|
||||
logger.warning("[gpu-direct] eligibility check failed (fallback)", exc_info=True)
|
||||
return False
|
||||
|
||||
def _try_gpu_direct(
|
||||
self,
|
||||
*,
|
||||
layers: list[RenderLayer],
|
||||
ass_path: Path | None,
|
||||
video_duration: float,
|
||||
output_path: Path,
|
||||
) -> tuple[bool, bool] | tuple[None, bool]:
|
||||
"""尝试全 GPU 直连渲染。成功返回 (True, edge_crop_applied),不支持/失败返回 (None, False)。"""
|
||||
if not self._can_use_gpu_direct(layers):
|
||||
return (None, False)
|
||||
if not self._gpu_encode_available():
|
||||
return (None, False)
|
||||
|
||||
try:
|
||||
from video_processing import gpu_direct_pipeline as gdp
|
||||
|
||||
cfg = self.plan.config or {}
|
||||
video_layer = next(_lyr for _lyr in layers if _lyr.role not in ("audio",))
|
||||
video_clips = [c for c in video_layer.clips if c.clip_type != "audio"]
|
||||
|
||||
# 音频层处理:收集 TTS 分段与配音素材库整段音频
|
||||
# - TTS 分段(带 tts 标记)→ 无间隙 concat 成单文件
|
||||
# - 配音素材库(voice_library=True)→ 单独作为整段音轨(不走分段 concat,已从 0 覆盖整段)
|
||||
audio_layer = next((_lyr for _lyr in layers if _lyr.role == "audio"), None)
|
||||
tts_merged: Path | None = None
|
||||
voiceover_track: Path | None = None
|
||||
if audio_layer:
|
||||
tts_clips = [c for c in audio_layer.clips if (c.config or {}).get("tts") and c.local_path.exists()]
|
||||
if tts_clips:
|
||||
tts_merged = self._concat_audio_clips(tts_clips, tag="tts_direct")
|
||||
# 配音素材库整段音频(按 _maybe_add_voice_library_layer 约定只有一个 clip_id=voice_library_main)
|
||||
vo_clips = [
|
||||
c for c in audio_layer.clips if (c.config or {}).get("voice_library") and c.local_path.exists()
|
||||
]
|
||||
if vo_clips:
|
||||
voiceover_track = vo_clips[-1].local_path # 理论上只有一个,取最后一个
|
||||
logger.info(
|
||||
"[gpu-direct] 配音素材库音轨: plan_id=%s path=%s",
|
||||
self.plan.id,
|
||||
voiceover_track,
|
||||
)
|
||||
|
||||
# 额外独立音轨(TTS concat、配音素材库)→ gpu_direct_pipeline 会与主音轨/BGM 一起 amix
|
||||
extra_audio_tracks: list[tuple[Path, float]] = []
|
||||
if tts_merged:
|
||||
extra_audio_tracks.append((tts_merged, 1.0))
|
||||
if voiceover_track:
|
||||
extra_audio_tracks.append((voiceover_track, 1.0))
|
||||
|
||||
# BGM 本地文件
|
||||
bgm_path = Path(self.bgm_path) if self.bgm_path else None
|
||||
if bgm_path is not None and not bgm_path.exists():
|
||||
bgm_path = None
|
||||
|
||||
# 字幕:标题 + ASR 时间轴
|
||||
title_cfg = cfg.get("title", {}) or cfg.get("title_config", {}) or {}
|
||||
title_text = ""
|
||||
if isinstance(title_cfg, dict) and title_cfg.get("enabled", True):
|
||||
title_text = title_cfg.get("text", "") or ""
|
||||
|
||||
subtitle_segments: list[Any] = []
|
||||
sub_cfg = cfg.get("subtitle", {}) or {}
|
||||
if isinstance(sub_cfg, dict) and sub_cfg.get("enabled", True):
|
||||
if sub_cfg.get("auto_generated") and self._asr_timeline_cache is not None:
|
||||
subtitle_segments = list(self._asr_timeline_cache.segments)
|
||||
|
||||
# 边缘裁剪:dedup 开启时在 GPU 内做四边随机 2~5% 裁剪(gpu_direct_pipeline 内部随机)
|
||||
dedup = self._dedup_enabled()
|
||||
edge_pct = 0.03 if dedup else 0.0 # >0 表示启用;实际区间 [2%,5%] 在 pipeline 内随机
|
||||
|
||||
# 探测每个视频素材是否含音轨、读取 volume 配置
|
||||
clip_has_audio_list: list[bool] = []
|
||||
clip_volumes_list: list[float] = []
|
||||
for c in video_clips:
|
||||
lp = getattr(c, "local_path", None)
|
||||
_ha = False
|
||||
if lp and Path(lp).exists():
|
||||
try:
|
||||
_ha = probe_has_audio(str(lp))
|
||||
except Exception as _pe: # noqa: BLE001
|
||||
logger.warning("[gpu-direct] probe_has_audio 失败按有声处理: %s", _pe)
|
||||
_ha = True
|
||||
clip_has_audio_list.append(_ha)
|
||||
_vol = float((c.config or {}).get("volume", 1.0))
|
||||
clip_volumes_list.append(_vol if _vol > 0 else 0.0)
|
||||
|
||||
plan = gdp.build_direct_render(
|
||||
resolved_clips=video_clips,
|
||||
output_width=self.output_width,
|
||||
output_height=self.output_height,
|
||||
output_fps=self.output_fps,
|
||||
bgm_audio=bgm_path,
|
||||
title_text=title_text,
|
||||
subtitle_segments=subtitle_segments,
|
||||
edge_crop_pct=edge_pct,
|
||||
total_duration=video_duration,
|
||||
clip_has_audio=clip_has_audio_list,
|
||||
clip_volumes=clip_volumes_list,
|
||||
extra_audio_tracks=extra_audio_tracks,
|
||||
)
|
||||
|
||||
client = get_gpu_encoder()
|
||||
client.render_inputs_to_output(plan.inputs, plan.ffmpeg_args, output_path)
|
||||
|
||||
# 清理本次上传的临时音频
|
||||
for key in plan.oss_keys:
|
||||
try:
|
||||
from video_processing.oss_helpers import _storage
|
||||
|
||||
_storage().delete_file(key) if hasattr(_storage(), "delete_file") else None
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
did_edge_crop = bool(edge_pct)
|
||||
logger.info(
|
||||
"[gpu-direct] success: plan_id=%s clips=%d edge_crop=%s", self.plan.id, len(video_clips), did_edge_crop
|
||||
)
|
||||
return (True, did_edge_crop)
|
||||
|
||||
except GpuEncodeError as e:
|
||||
logger.warning("[gpu-direct] failed (fallback to legacy): %s", e)
|
||||
try:
|
||||
if output_path.exists():
|
||||
output_path.unlink()
|
||||
except OSError:
|
||||
pass
|
||||
return (None, False)
|
||||
except Exception: # noqa: BLE001
|
||||
logger.warning("[gpu-direct] unexpected error (fallback)", exc_info=True)
|
||||
return (None, False)
|
||||
|
||||
def _concat_audio_clips(self, clips: list[Any], *, tag: str) -> Path:
|
||||
"""把多个本地音频片段无间隙 concat 成一个 m4a(TTS 分段→单文件)。"""
|
||||
out = self.work_dir / f"{tag}_{self.plan.id}.m4a"
|
||||
listfile = self.work_dir / f"{tag}_{self.plan.id}.txt"
|
||||
lines = []
|
||||
for c in clips:
|
||||
ap = str(c.local_path).replace("'", "'\\''")
|
||||
lines.append(f"file '{ap}'")
|
||||
listfile.write_text("\n".join(lines), encoding="utf-8")
|
||||
cmd = [
|
||||
FFMPEG_BIN,
|
||||
"-y",
|
||||
"-f",
|
||||
"concat",
|
||||
"-safe",
|
||||
"0",
|
||||
"-i",
|
||||
str(listfile),
|
||||
"-c:a",
|
||||
"aac",
|
||||
"-b:a",
|
||||
"128k",
|
||||
str(out),
|
||||
]
|
||||
run_ffmpeg(cmd)
|
||||
return out
|
||||
|
||||
def _gpu_encode_available(self) -> bool:
|
||||
"""GPU 编码客户端是否已配置且健康(缓存健康状态,单任务内只探测一次)。"""
|
||||
if not getattr(self, "_gpu_health_ok", None):
|
||||
@@ -2594,7 +2831,7 @@ class UnifiedRenderService:
|
||||
b = pixel_pert.get("color_b", 0)
|
||||
if r != 0 or g != 0 or b != 0:
|
||||
# color_balance 参数范围 -1.0 ~ 1.0,这里用 /100 转换
|
||||
filters.append(f"colorbalance=rs={r/100:.3f}:gs={g/100:.3f}:bs={b/100:.3f}")
|
||||
filters.append(f"colorbalance=rs={r / 100:.3f}:gs={g / 100:.3f}:bs={b / 100:.3f}")
|
||||
|
||||
@staticmethod
|
||||
def _clip_volume(clip: ResolvedClip) -> float:
|
||||
|
||||
@@ -198,7 +198,7 @@ def _verify_url_accessible(
|
||||
retries: int = 2,
|
||||
max_redirects: int = 5,
|
||||
) -> bool:
|
||||
"""HEAD 请求校验 URL 可访问(含重试,防止 OSS 抖动误报)。
|
||||
"""GET+Range 请求校验 URL 可访问(含重试,防止 OSS 抖动误报)。
|
||||
|
||||
安全增强:
|
||||
- 请求前先做 SSRF 安全校验(内网IP/回环地址/链路本地地址等)
|
||||
@@ -255,11 +255,12 @@ def _verify_url_accessible(
|
||||
)
|
||||
raise
|
||||
|
||||
req = urllib.request.Request(safe_url, method="HEAD")
|
||||
req = urllib.request.Request(safe_url, method="GET")
|
||||
req.add_header("Range", "bytes=0-0")
|
||||
req.add_header("User-Agent", "xiaoxia-saas-worker/1.0")
|
||||
|
||||
with opener.open(req, timeout=timeout) as resp: # noqa: S310
|
||||
if 200 <= resp.status < 300:
|
||||
if 200 <= resp.status < 300 or resp.status == 206:
|
||||
return True
|
||||
if resp.status in (301, 302, 303, 307, 308):
|
||||
location = resp.headers.get("Location", "")
|
||||
@@ -693,11 +694,11 @@ def _render_from_edit_plan(
|
||||
task_id: str,
|
||||
source_edit_plan_id: str,
|
||||
task_info: dict,
|
||||
) -> tuple[Path, float, list[dict] | None, str | None, str | None, str]:
|
||||
) -> tuple[Path, float, list[dict] | None, str | None, str | None, str, bool]:
|
||||
"""从 EditPlan 数据库记录直接渲染(不再内存重建clips)。
|
||||
|
||||
Returns:
|
||||
(output_path, render_duration, cover_candidates, voiceover_path, temp_dir, thumbnail_url)
|
||||
(output_path, render_duration, cover_candidates, voiceover_path, temp_dir, thumbnail_url, edge_crop_applied)
|
||||
"""
|
||||
from video_processing.render_adapter import RenderAdapter
|
||||
from worker_app.db import SessionLocal
|
||||
@@ -746,6 +747,7 @@ def _render_from_edit_plan(
|
||||
voiceover_path,
|
||||
render_temp_dir,
|
||||
result.thumbnail_url or "",
|
||||
bool(getattr(result, "edge_crop_applied", False)),
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
@@ -899,6 +901,7 @@ def generate_video(self, task_id: str) -> dict:
|
||||
voiceover_tmp_path,
|
||||
render_temp_dir,
|
||||
thumbnail_url,
|
||||
_gpu_edge_crop_done,
|
||||
) = _render_from_edit_plan(
|
||||
task_id=task_id,
|
||||
source_edit_plan_id=current_plan_id,
|
||||
@@ -935,6 +938,12 @@ def generate_video(self, task_id: str) -> dict:
|
||||
if gen_task and render_attempt == 0:
|
||||
gen_task.append_log("降重", "已关闭边缘裁剪与微变换(确定性渲染)")
|
||||
_flush_logs(task_id, gen_task)
|
||||
elif _gpu_edge_crop_done:
|
||||
# GPU 直连管线已经在 filter_complex 中做了随机边缘裁剪,跳过 CPU 二次重编码
|
||||
if gen_task and render_attempt == 0:
|
||||
gen_task.append_log("边缘裁剪", "已在 GPU 直连管线内完成随机边缘裁剪")
|
||||
_flush_logs(task_id, gen_task)
|
||||
logger.info("[task_id=%s] GPU直连已完成边缘裁剪,跳过CPU二次重编码", task_id)
|
||||
else:
|
||||
from video_processing.ffmpeg_utils import random_edge_crop
|
||||
|
||||
|
||||
@@ -293,5 +293,5 @@ GPU_ENCODE_VCODEC=h264_nvenc
|
||||
GPU_ENCODE_PRESET=p4
|
||||
GPU_ENCODE_CRF=23
|
||||
GPU_ENCODE_FALLBACK_CPU=true
|
||||
GPU_ENCODE_MEZZANINE_TRANSPORT=relay
|
||||
GPU_ENCODE_MEZZANINE_TRANSPORT=oss
|
||||
GPU_ENCODE_OSS_TMP_PREFIX=tmp/gpu-mezzanine/
|
||||
|
||||
@@ -0,0 +1,121 @@
|
||||
# Host nginx config for staging server: /etc/nginx/sites-available/05-xiaoxia-cms
|
||||
# Xiaoxia CMS - cms.xiaoxiajianji.com
|
||||
#
|
||||
# 注意:此文件是宿主机 nginx 配置的备份/参考,不是 Docker 容器内的 nginx。
|
||||
# Docker 容器内的 nginx 配置见 nginx-staging.conf。
|
||||
#
|
||||
# GPU relay 路由说明:
|
||||
# P4000 编码完成后通过 http://100.69.73.60:8092/api/v1/internal/gpu-relay/{key} PUT 上传
|
||||
# Worker 容器通过 http://xiaoxia-api-staging:8000/api/v1/internal/gpu-relay/{key} GET 下载
|
||||
# 8092 端口由 CMS 宿主机 nginx 承载,GPU relay 路由通过最长前缀匹配优先代理到 staging API (8000)
|
||||
|
||||
# 80 端口:ACME 验证 + 重定向到 HTTPS
|
||||
server {
|
||||
listen 80;
|
||||
server_name cms.xiaoxiajianji.com;
|
||||
|
||||
# Let's Encrypt ACME 验证
|
||||
location /.well-known/acme-challenge/ {
|
||||
root /var/www/certbot;
|
||||
}
|
||||
|
||||
location / {
|
||||
return 301 https://$host$request_uri;
|
||||
}
|
||||
}
|
||||
|
||||
# 443 端口:CMS 主站
|
||||
server {
|
||||
listen 443 ssl http2;
|
||||
server_name cms.xiaoxiajianji.com;
|
||||
|
||||
ssl_certificate /etc/letsencrypt/live/cms.xiaoxiajianji.com/fullchain.pem;
|
||||
ssl_certificate_key /etc/letsencrypt/live/cms.xiaoxiajianji.com/privkey.pem;
|
||||
include /etc/letsencrypt/options-ssl-nginx.conf;
|
||||
ssl_dhparam /etc/letsencrypt/ssl-dhparams.pem;
|
||||
|
||||
client_max_body_size 50m;
|
||||
|
||||
# Security headers
|
||||
add_header X-Content-Type-Options "nosniff" always;
|
||||
add_header X-Frame-Options "SAMEORIGIN" always;
|
||||
add_header X-XSS-Protection "1; mode=block" always;
|
||||
add_header Referrer-Policy "strict-origin-when-cross-origin" always;
|
||||
add_header Strict-Transport-Security "max-age=31536000; includeSubDomains" always;
|
||||
|
||||
root /data/www/cms/current;
|
||||
index index.html;
|
||||
|
||||
gzip on;
|
||||
gzip_types text/plain text/css application/json application/javascript text/xml application/xml application/xml+rss text/javascript image/svg+xml;
|
||||
gzip_min_length 1024;
|
||||
|
||||
location / {
|
||||
try_files $uri $uri/ /index.html;
|
||||
}
|
||||
|
||||
location /api/ {
|
||||
proxy_pass http://127.0.0.1:8091/api/;
|
||||
proxy_http_version 1.1;
|
||||
proxy_set_header Host $host;
|
||||
proxy_set_header X-Real-IP $remote_addr;
|
||||
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
proxy_set_header X-Forwarded-Proto $scheme;
|
||||
}
|
||||
|
||||
location = /health {
|
||||
proxy_pass http://127.0.0.1:8091/health;
|
||||
}
|
||||
}
|
||||
|
||||
# 临时访问:8092 端口(IP直接访问,后续可关闭)
|
||||
server {
|
||||
listen 8092;
|
||||
server_name _;
|
||||
|
||||
root /data/www/cms/current;
|
||||
index index.html;
|
||||
|
||||
client_max_body_size 50m;
|
||||
|
||||
add_header X-Content-Type-Options "nosniff" always;
|
||||
add_header X-Frame-Options "SAMEORIGIN" always;
|
||||
add_header X-XSS-Protection "1; mode=block" always;
|
||||
add_header Referrer-Policy "strict-origin-when-cross-origin" always;
|
||||
add_header Strict-Transport-Security "max-age=31536000; includeSubDomains" always;
|
||||
|
||||
gzip on;
|
||||
gzip_types text/plain text/css application/json application/javascript text/xml application/xml application/xml+rss text/javascript image/svg+xml;
|
||||
gzip_min_length 1024;
|
||||
|
||||
location / {
|
||||
try_files $uri $uri/ /index.html;
|
||||
}
|
||||
|
||||
# GPU relay endpoints - proxy to staging API (port 8000) instead of CMS
|
||||
# 此 location 必须在 location /api/ 之前,利用 nginx 最长前缀匹配优先路由
|
||||
location /api/v1/internal/gpu-relay/ {
|
||||
proxy_pass http://127.0.0.1:8000/api/v1/internal/gpu-relay/;
|
||||
proxy_http_version 1.1;
|
||||
proxy_set_header Host $host;
|
||||
proxy_set_header X-Real-IP $remote_addr;
|
||||
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
proxy_set_header X-Forwarded-Proto $scheme;
|
||||
proxy_request_buffering off;
|
||||
proxy_read_timeout 600s;
|
||||
proxy_send_timeout 600s;
|
||||
}
|
||||
|
||||
location /api/ {
|
||||
proxy_pass http://127.0.0.1:8091/api/;
|
||||
proxy_http_version 1.1;
|
||||
proxy_set_header Host $host;
|
||||
proxy_set_header X-Real-IP $remote_addr;
|
||||
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
proxy_set_header X-Forwarded-Proto $scheme;
|
||||
}
|
||||
|
||||
location = /health {
|
||||
proxy_pass http://127.0.0.1:8091/health;
|
||||
}
|
||||
}
|
||||
@@ -1,31 +1,29 @@
|
||||
# Staging GPU relay plain-HTTP vhost (P4000 NVENC 编码回传入口)
|
||||
# - 监听 8092 端口纯 HTTP(绕开 HTTPS 证书与 P4000 httpx SSL 问题)
|
||||
# - 代理到本机 staging API 的 /api/ 路径(127.0.0.1:8000 是 docker 映射端口)
|
||||
# - P4000 通过 Tailscale 直连宿主机 100.69.73.60:8092 PUT 编码结果
|
||||
# - Worker 通过 Docker DNS (xiaoxia-api-staging:8000) 直接 GET/DELETE,
|
||||
# 不经宿主机 nginx,避免 UFW FORWARD DROP 阻断
|
||||
# Staging GPU relay nginx 配置说明
|
||||
#
|
||||
# 部署:cp infra/nginx/gpu-relay-staging.conf /etc/nginx/conf.d/ && nginx -t && systemctl reload nginx
|
||||
|
||||
server {
|
||||
listen 8092;
|
||||
server_name _;
|
||||
|
||||
client_max_body_size 2048m;
|
||||
|
||||
location /api/ {
|
||||
proxy_pass http://127.0.0.1:8000/api/;
|
||||
proxy_http_version 1.1;
|
||||
proxy_set_header Host $host;
|
||||
proxy_set_header X-Real-IP $remote_addr;
|
||||
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
proxy_set_header X-Forwarded-Proto $scheme;
|
||||
proxy_request_buffering off;
|
||||
proxy_read_timeout 600s;
|
||||
proxy_send_timeout 600s;
|
||||
}
|
||||
|
||||
location = /health {
|
||||
proxy_pass http://127.0.0.1:8000/health;
|
||||
}
|
||||
}
|
||||
# GPU relay 并没有独立的 nginx vhost,而是集成在宿主机 CMS nginx 的 8092 server block 中。
|
||||
# 完整宿主机 nginx 配置备份见: deploy/configs/host-nginx-cms-staging.conf
|
||||
#
|
||||
# 核心 location 块(添加到 8092 server block,位于 location /api/ 之前):
|
||||
#
|
||||
# # GPU relay endpoints - proxy to staging API (port 8000) instead of CMS
|
||||
# location /api/v1/internal/gpu-relay/ {
|
||||
# proxy_pass http://127.0.0.1:8000/api/v1/internal/gpu-relay/;
|
||||
# proxy_http_version 1.1;
|
||||
# proxy_set_header Host $host;
|
||||
# proxy_set_header X-Real-IP $remote_addr;
|
||||
# proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
|
||||
# proxy_set_header X-Forwarded-Proto $scheme;
|
||||
# proxy_request_buffering off;
|
||||
# proxy_read_timeout 600s;
|
||||
# proxy_send_timeout 600s;
|
||||
# }
|
||||
#
|
||||
# 部署方式:手动将上述 location 块添加到 /etc/nginx/sites-available/05-xiaoxia-cms 的 8092 server block 中
|
||||
# 然后 nginx -t && systemctl reload nginx
|
||||
#
|
||||
# 原理说明:
|
||||
# - 8092 端口由 CMS 宿主机 nginx 承载(与 CMS 共享端口)
|
||||
# - GPU relay 路由 /api/v1/internal/gpu-relay/ 比 CMS 的 /api/ 更具体
|
||||
# - nginx 最长前缀匹配确保 relay 请求路由到 staging API (port 8000) 而非 CMS (port 8091)
|
||||
# - P4000 通过 Tailscale IP 100.69.73.60:8092 访问 relay
|
||||
# - Worker 容器通过 Docker DNS xiaoxia-api-staging:8000 直接访问 relay
|
||||
|
||||
@@ -51,7 +51,7 @@ class APISettings(SharedSettings):
|
||||
def validate_jwt_secret_key(cls, v):
|
||||
if v is None or v == "":
|
||||
raise ValueError(
|
||||
"JWT_SECRET_KEY must be set via environment variable. " "Do not use default value in production!"
|
||||
"JWT_SECRET_KEY must be set via environment variable. Do not use default value in production!"
|
||||
)
|
||||
# Block known insecure default values
|
||||
insecure_defaults = [
|
||||
@@ -63,7 +63,7 @@ class APISettings(SharedSettings):
|
||||
]
|
||||
if v.lower() in [d.lower() for d in insecure_defaults]:
|
||||
raise ValueError(
|
||||
f"JWT_SECRET_KEY '{v}' is insecure. " "Please set a strong random secret via environment variable."
|
||||
f"JWT_SECRET_KEY '{v}' is insecure. Please set a strong random secret via environment variable."
|
||||
)
|
||||
return v
|
||||
|
||||
@@ -248,6 +248,10 @@ class APISettings(SharedSettings):
|
||||
def OSS_ENDPOINT(self) -> str:
|
||||
return self.oss_endpoint
|
||||
|
||||
@property
|
||||
def OSS_INTERNAL_ENDPOINT(self) -> str:
|
||||
return self.effective_oss_internal_endpoint
|
||||
|
||||
@property
|
||||
def OSS_ACCESS_KEY_ID(self) -> str:
|
||||
return self.oss_access_key_id
|
||||
|
||||
@@ -47,12 +47,37 @@ class SharedSettings(BaseSettings):
|
||||
|
||||
# ── OSS 阿里云 ──────────────────────────────────────────────────────
|
||||
oss_endpoint: str = "oss-cn-hangzhou.aliyuncs.com"
|
||||
# 内网 endpoint:ECS VPC 内访问 OSS 用(千兆带宽、免公网流量费)。
|
||||
# 为空时自动从 oss_endpoint 推导:若 oss_endpoint 是阿里云公网域名(形如
|
||||
# oss-cn-<region>.aliyuncs.com),自动加 -internal 得到内网域名;其他情况
|
||||
# (自定义域名/本地 MinIO/非阿里云)回退使用 oss_endpoint。
|
||||
# 显式填同值可以覆盖自动推导、强制所有流量都走公网。
|
||||
oss_internal_endpoint: str = ""
|
||||
oss_access_key_id: str = ""
|
||||
oss_access_key_secret: str = ""
|
||||
oss_bucket_name: str = "xiaoxia-autocut"
|
||||
oss_direct_upload_max_mb: int = 2000
|
||||
oss_direct_upload_expire_seconds: int = 900
|
||||
|
||||
@property
|
||||
def effective_oss_internal_endpoint(self) -> str:
|
||||
"""实际用于 SDK 内网访问的 endpoint(带 -internal 自动推导)。"""
|
||||
if self.oss_internal_endpoint:
|
||||
return self.oss_internal_endpoint
|
||||
ep = self.oss_endpoint.strip()
|
||||
scheme = ""
|
||||
host = ep
|
||||
if ep.startswith("https://"):
|
||||
scheme = "https://"
|
||||
host = ep[len("https://") :]
|
||||
elif ep.startswith("http://"):
|
||||
scheme = "http://"
|
||||
host = ep[len("http://") :]
|
||||
# 阿里云公网域名自动推导:oss-cn-<region>.aliyuncs.com → oss-cn-<region>-internal.aliyuncs.com
|
||||
if host.endswith(".aliyuncs.com") and "-internal" not in host and host.startswith("oss-cn-"):
|
||||
host = host[: -len(".aliyuncs.com")] + "-internal.aliyuncs.com"
|
||||
return f"{scheme}{host}" if scheme else host
|
||||
|
||||
# ── CosyVoice (阿里云百炼语音合成) ───────────────────────────────────
|
||||
cosyvoice_api_key: str = ""
|
||||
cosyvoice_base_url: str = "https://dashscope.aliyuncs.com/api/v1"
|
||||
@@ -76,6 +101,7 @@ class SharedSettings(BaseSettings):
|
||||
mediakit_api_key: str = ""
|
||||
mediakit_base_url: str = "https://mediakit.cn-beijing.volces.com/api/v1"
|
||||
mediakit_timeout: int = 60
|
||||
mediakit_cover_enabled: bool = False # 封面抽帧是否走MediaKit(默认false走本地ffmpeg+cv2,<2s完成)
|
||||
|
||||
# ── 积分/会员系统 (#1895) ────────────────────────────────────────────
|
||||
# 积分系统总开关(产品要求 #1895:暂停积分系统但保留全部代码/表/接口)。
|
||||
|
||||
@@ -309,6 +309,103 @@ class GpuEncoderClient:
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("[gpu-encoder] failed to delete OSS mezzanine %s: %s", oss_key, e)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# High-level: render arbitrary inputs → final output (all-GPU pipeline)
|
||||
# ------------------------------------------------------------------
|
||||
def render_inputs_to_output(
|
||||
self,
|
||||
inputs: dict[str, str],
|
||||
ffmpeg_args: list[str],
|
||||
output_path: Path,
|
||||
*,
|
||||
timeout: Optional[int] = None,
|
||||
) -> dict[str, Any]:
|
||||
"""把多输入(原始素材/字幕/BGM)连同完整 filter_complex 交给 P4000 一次出片。
|
||||
|
||||
与 encode_mezzanine_to_output 的区别:worker 侧不再生成/上传 mezzanine,
|
||||
P4000 直接从 inputs 中的签名 URL 下载原始素材,filter_complex 内完成
|
||||
concat/scale/crop/drawtext/amix,末端 h264_nvenc 只编码一次。
|
||||
|
||||
传输:成片仍走 relay 回传(P4000 PUT → worker GET),避免公网 OSS 往返。
|
||||
|
||||
Args:
|
||||
inputs: {裸文件名: 可下载URL},key 即 ffmpeg_args 中引用的文件名
|
||||
ffmpeg_args: 完整 ffmpeg 参数(含 -i、-filter_complex、-map、NVENC 编码参数)
|
||||
output_path: worker 本地成片落盘路径
|
||||
timeout: P4000 侧超时(秒)
|
||||
"""
|
||||
if not inputs:
|
||||
raise GpuEncodeError("render_inputs_to_output: inputs is empty")
|
||||
if not ffmpeg_args:
|
||||
raise GpuEncodeError("render_inputs_to_output: ffmpeg_args is empty")
|
||||
if not self.relay_base_url:
|
||||
raise GpuEncodeError("gpu_encode_relay_base_url not configured")
|
||||
|
||||
timeout = timeout or self.sync_timeout
|
||||
t_total = time.time()
|
||||
result_key: Optional[str] = None
|
||||
|
||||
try:
|
||||
secret = self._get_relay_secret()
|
||||
|
||||
# 1. result relay URLs(成片 P4000 PUT → worker GET)
|
||||
result_key = uuid.uuid4().hex
|
||||
put_url = self._result_put_url(result_key, secret)
|
||||
get_url = self._result_get_url(result_key, secret)
|
||||
del_result_url = get_url
|
||||
|
||||
# 2. pre-warm then call P4000 sync render
|
||||
self._warm_up_if_needed()
|
||||
body = {
|
||||
"inputs": dict(inputs),
|
||||
"ffmpeg_args": list(ffmpeg_args),
|
||||
"output_url": put_url,
|
||||
"timeout": int(timeout),
|
||||
}
|
||||
job = self._post_sync(body)
|
||||
self._last_ok_ts = time.time()
|
||||
logger.info(
|
||||
"[gpu-encoder] P4000 direct done: job_id=%s rc=%s size=%s dur=%ss inputs=%d",
|
||||
job.get("job_id"),
|
||||
job.get("ffmpeg_rc"),
|
||||
job.get("size"),
|
||||
job.get("duration"),
|
||||
len(inputs),
|
||||
)
|
||||
|
||||
# 3. download result from relay to output_path
|
||||
output_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
size = self._download_to_file(get_url, output_path)
|
||||
|
||||
# 4. cleanup relay result
|
||||
self._relay_delete(del_result_url)
|
||||
|
||||
logger.info(
|
||||
"[gpu-encoder] direct render ok → %s (%d bytes) total=%.2fs",
|
||||
output_path.name,
|
||||
size,
|
||||
time.time() - t_total,
|
||||
)
|
||||
return {
|
||||
"job": job,
|
||||
"output_size": size,
|
||||
"output_path": str(output_path),
|
||||
"transport": "direct",
|
||||
}
|
||||
|
||||
except GpuEncodeError:
|
||||
raise
|
||||
except Exception as e: # noqa: BLE001
|
||||
raise GpuEncodeError(f"unexpected: {e}") from e
|
||||
finally:
|
||||
# cleanup relay result (best-effort)
|
||||
if result_key:
|
||||
try:
|
||||
secret = self._get_relay_secret()
|
||||
self._relay_delete(self._relay_result_url(self.relay_internal_base_url, result_key, secret))
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("[gpu-encoder] failed to delete relay result %s: %s", result_key, e)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Internal helpers
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
+134
-168
@@ -5,6 +5,13 @@
|
||||
- Worker端 oss_helpers 的高级能力(分片上传/超时保护/HTTP下载/Asset路径解析)
|
||||
|
||||
所有服务都通过这个统一入口与存储交互,消除重复实现。
|
||||
|
||||
P1 (2026-09-28) OSS 双 endpoint 分离:
|
||||
- 内部 bucket(self.bucket):使用 internal endpoint(VPC 千兆带宽),
|
||||
用于所有 SDK 上传/下载/删除/object_exists 操作;
|
||||
- 公网 bucket(self.public_bucket):使用公网 endpoint,仅用于 sign_url
|
||||
生成给前端/P4000/MediaKit 等外网访问方用的预签名 URL;
|
||||
- public_url 永远拼公网域名,不随 internal endpoint 变化。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -42,19 +49,44 @@ OSS_MULTIPART_NUM_THREADS = 3 # 分片上传并发数
|
||||
OSS_HTTP_DOWNLOAD_TIMEOUT = 300 # HTTP下载超时(秒)
|
||||
|
||||
|
||||
class SharedStorageService(StoragePort):
|
||||
"""统一存储服务 — 实现 StoragePort,API 和 Worker 共用。
|
||||
def _make_bucket(
|
||||
auth,
|
||||
endpoint: str,
|
||||
bucket_name: str,
|
||||
*,
|
||||
connect_timeout: int = OSS_CONNECT_TIMEOUT,
|
||||
app_name: str = "",
|
||||
):
|
||||
"""构造 oss2.Bucket,自动补 https:// 前缀。"""
|
||||
if not endpoint.startswith(("http://", "https://")):
|
||||
endpoint = f"https://{endpoint}"
|
||||
kwargs: dict = {"connect_timeout": connect_timeout}
|
||||
if app_name:
|
||||
kwargs["app_name"] = app_name
|
||||
return oss2.Bucket(auth, endpoint, bucket_name, **kwargs)
|
||||
|
||||
整合了原 SharedStorageService + oss_helpers 的全部能力。
|
||||
"""
|
||||
|
||||
class SharedStorageService(StoragePort):
|
||||
"""统一存储服务 — 实现 StoragePort,API 和 Worker 共用。"""
|
||||
|
||||
# 类级默认值,方便单测 mock __init__ 后实例仍有这些属性
|
||||
bucket: Optional[object] = None
|
||||
public_bucket: Optional[object] = None
|
||||
public_endpoint: str = ""
|
||||
internal_endpoint: str = ""
|
||||
public_url: str = ""
|
||||
local_url_prefix: str = "/generated-files"
|
||||
bucket_name: str = ""
|
||||
|
||||
def __init__(self):
|
||||
settings = get_shared_settings()
|
||||
self.bucket_name = settings.oss_bucket_name
|
||||
self.endpoint = settings.oss_endpoint
|
||||
self.public_url = f"https://{settings.oss_bucket_name}.{settings.oss_endpoint}"
|
||||
self.public_endpoint = settings.oss_endpoint # 公网 endpoint,用于签名 URL
|
||||
self.internal_endpoint = settings.effective_oss_internal_endpoint # 内网 endpoint,SDK 用
|
||||
self.public_url = f"https://{settings.oss_bucket_name}.{self._public_host()}"
|
||||
self.local_url_prefix = os.getenv("GENERATED_FILES_URL_PREFIX", "/generated-files")
|
||||
self.bucket = None
|
||||
self.bucket: Optional[object] = None # internal: SDK 上传/下载/删除
|
||||
self.public_bucket: Optional[object] = None # public: sign_url 给外网
|
||||
|
||||
self.access_key_id = settings.oss_access_key_id
|
||||
self.access_key_secret = settings.oss_access_key_secret
|
||||
@@ -65,21 +97,26 @@ class SharedStorageService(StoragePort):
|
||||
if has_key_id and has_key_secret:
|
||||
if oss2 is not None:
|
||||
try:
|
||||
# endpoint 不带 scheme 时补 https:// 前缀
|
||||
bucket_endpoint = self.endpoint
|
||||
if not bucket_endpoint.startswith(("http://", "https://")):
|
||||
bucket_endpoint = f"https://{bucket_endpoint}"
|
||||
auth = oss2.Auth(self.access_key_id, self.access_key_secret)
|
||||
self.bucket = oss2.Bucket(
|
||||
self.bucket = _make_bucket(
|
||||
auth,
|
||||
bucket_endpoint,
|
||||
self.internal_endpoint,
|
||||
self.bucket_name,
|
||||
connect_timeout=OSS_CONNECT_TIMEOUT,
|
||||
app_name="xiaoxia-internal",
|
||||
)
|
||||
logger.info(
|
||||
"OSS initialized: endpoint=%s bucket=%s",
|
||||
self.endpoint,
|
||||
self.public_bucket = _make_bucket(
|
||||
auth,
|
||||
self.public_endpoint,
|
||||
self.bucket_name,
|
||||
app_name="xiaoxia-public",
|
||||
)
|
||||
same_ep = self.internal_endpoint == self.public_endpoint
|
||||
logger.info(
|
||||
"OSS initialized: public_ep=%s internal_ep=%s bucket=%s dual=%s",
|
||||
self.public_endpoint,
|
||||
self.internal_endpoint,
|
||||
self.bucket_name,
|
||||
"no" if same_ep else "yes",
|
||||
)
|
||||
except Exception as error:
|
||||
logger.error("Failed to initialize OSS bucket client: %s", error)
|
||||
@@ -93,26 +130,35 @@ class SharedStorageService(StoragePort):
|
||||
missing.append("OSS_ACCESS_KEY_SECRET")
|
||||
logger.error("OSS credentials not configured — missing: %s", ", ".join(missing))
|
||||
|
||||
def _public_host(self) -> str:
|
||||
ep = self.public_endpoint
|
||||
if ep.startswith("https://"):
|
||||
return ep[len("https://") :]
|
||||
if ep.startswith("http://"):
|
||||
return ep[len("http://") :]
|
||||
return ep
|
||||
|
||||
# ── 诊断 ───────────────────────────────────────────────────────────
|
||||
|
||||
def diagnose(self) -> None:
|
||||
"""输出存储配置诊断日志。"""
|
||||
key_id_display = (
|
||||
f"{self.access_key_id[:4]}...{self.access_key_id[-4:]}" if len(self.access_key_id) > 8 else "(empty)"
|
||||
)
|
||||
logger.info(
|
||||
"[OSS诊断] endpoint=%s bucket_name=%s access_key_id=%s",
|
||||
self.endpoint,
|
||||
"[OSS诊断] public_ep=%s internal_ep=%s bucket=%s ak=%s",
|
||||
self.public_endpoint,
|
||||
self.internal_endpoint,
|
||||
self.bucket_name,
|
||||
key_id_display,
|
||||
)
|
||||
if self.bucket is None:
|
||||
logger.error(
|
||||
"[OSS诊断] ❌ bucket=None — 预签名URL不可用!"
|
||||
"原因: OSS_ACCESS_KEY_ID/OSS_ACCESS_KEY_SECRET 未配置或 oss2 未安装。"
|
||||
)
|
||||
logger.error("[OSS诊断] ❌ bucket(internal)=None")
|
||||
else:
|
||||
logger.info("[OSS诊断] ✅ bucket 已配置,预签名URL可用")
|
||||
logger.info("[OSS诊断] ✅ bucket(internal) 就绪")
|
||||
if self.public_bucket is None:
|
||||
logger.error("[OSS诊断] ❌ public_bucket=None")
|
||||
else:
|
||||
logger.info("[OSS诊断] ✅ public_bucket 就绪,公网签名URL可用")
|
||||
|
||||
# ── 工具方法 ───────────────────────────────────────────────────────
|
||||
|
||||
@@ -122,20 +168,15 @@ class SharedStorageService(StoragePort):
|
||||
return path.startswith(f"{self.local_url_prefix}/")
|
||||
|
||||
def _normalize_storage_key(self, storage_key_or_url: str) -> str:
|
||||
"""从 URL 提取存储键,并做 URL 解码。
|
||||
|
||||
防止 URL 编码的字符(空格=%20、中文=%XX)导致签名不匹配。
|
||||
"""
|
||||
if storage_key_or_url.startswith("http://") or storage_key_or_url.startswith("https://"):
|
||||
parsed = urlparse(storage_key_or_url)
|
||||
return unquote(parsed.path.lstrip("/"))
|
||||
return storage_key_or_url.lstrip("/")
|
||||
|
||||
def normalize_storage_key(self, storage_key_or_url: str) -> str:
|
||||
"""从 URL 提取存储键(公开方法)。"""
|
||||
return self._normalize_storage_key(storage_key_or_url)
|
||||
|
||||
# ── 上传 ───────────────────────────────────────────────────────────
|
||||
# ── 上传(SDK 走 internal endpoint)───────────────────────────────
|
||||
|
||||
def upload_file(
|
||||
self,
|
||||
@@ -143,21 +184,22 @@ class SharedStorageService(StoragePort):
|
||||
storage_key: str,
|
||||
content_type: str = "application/octet-stream",
|
||||
) -> str:
|
||||
"""上传文件到存储,返回公开 URL(简单上传,API端原有行为)。
|
||||
|
||||
- 路径字符串 → bucket.put_object_from_file
|
||||
- 类文件对象 → bucket.put_object
|
||||
- bucket未配置 → 抛 RuntimeError
|
||||
"""
|
||||
if self.bucket is None:
|
||||
raise RuntimeError("OSS storage is not configured")
|
||||
|
||||
try:
|
||||
if isinstance(file_or_path, (str, Path)):
|
||||
self.bucket.put_object_from_file(storage_key, str(file_or_path), headers={"Content-Type": content_type})
|
||||
self.bucket.put_object_from_file(
|
||||
storage_key,
|
||||
str(file_or_path),
|
||||
headers={"Content-Type": content_type},
|
||||
)
|
||||
else:
|
||||
file_or_path.seek(0) # type: ignore[attr-defined]
|
||||
self.bucket.put_object(storage_key, file_or_path, headers={"Content-Type": content_type})
|
||||
file_or_path.seek(0)
|
||||
self.bucket.put_object(
|
||||
storage_key,
|
||||
file_or_path,
|
||||
headers={"Content-Type": content_type},
|
||||
)
|
||||
return f"{self.public_url}/{storage_key}"
|
||||
except Exception as e:
|
||||
raise Exception(f"Failed to upload file to OSS: {e}") from e
|
||||
@@ -167,14 +209,6 @@ class SharedStorageService(StoragePort):
|
||||
local_path: str | Path,
|
||||
storage_key: str,
|
||||
) -> Optional[str]:
|
||||
"""智能上传:大文件自动分片+超时保护(从 oss_helpers 合并)。
|
||||
|
||||
- 大文件(>100MB)走分片上传,3 线程并发
|
||||
- 总超时 300s,防止网络异常时挂死
|
||||
- 成功返回 URL,失败返回 None(不抛异常)
|
||||
|
||||
Worker端 oss_helpers.upload_to_oss 的统一入口。
|
||||
"""
|
||||
local_path = Path(local_path)
|
||||
if not local_path.exists():
|
||||
logger.error("上传文件不存在: %s", local_path)
|
||||
@@ -198,10 +232,10 @@ class SharedStorageService(StoragePort):
|
||||
|
||||
if use_multipart:
|
||||
logger.info(
|
||||
"大文件分片上传: storage_key=%s, size=%.1fMB, part_size=%dMB, threads=%d",
|
||||
"大文件分片上传(internal): key=%s size=%.1fMB part=%dMB threads=%d",
|
||||
storage_key[:80],
|
||||
file_size / 1024 / 1024,
|
||||
OSS_PART_SIZE // 1024 // 1024,
|
||||
file_size / 1048576,
|
||||
OSS_PART_SIZE // 1048576,
|
||||
OSS_MULTIPART_NUM_THREADS,
|
||||
)
|
||||
oss2.resumable_upload(
|
||||
@@ -214,7 +248,6 @@ class SharedStorageService(StoragePort):
|
||||
)
|
||||
else:
|
||||
self.bucket.put_object_from_file(storage_key, str(local_path))
|
||||
|
||||
result["url"] = f"{self.public_url}/{storage_key}"
|
||||
except Exception as e:
|
||||
result["error"] = e
|
||||
@@ -222,73 +255,54 @@ class SharedStorageService(StoragePort):
|
||||
finally:
|
||||
done.set()
|
||||
|
||||
upload_thread = threading.Thread(target=_do_upload, daemon=True)
|
||||
upload_thread.start()
|
||||
t = threading.Thread(target=_do_upload, daemon=True)
|
||||
t.start()
|
||||
finished = done.wait(timeout=OSS_UPLOAD_TOTAL_TIMEOUT)
|
||||
|
||||
if not finished:
|
||||
logger.error(
|
||||
"OSS 上传超时(%.0fs),强制中止: storage_key=%s, size=%.1fMB",
|
||||
"OSS 上传超时(%ds): key=%s size=%.1fMB",
|
||||
OSS_UPLOAD_TOTAL_TIMEOUT,
|
||||
storage_key[:80],
|
||||
result["file_size"] / 1024 / 1024 if result["file_size"] else 0,
|
||||
result["file_size"] / 1048576 if result["file_size"] else 0,
|
||||
)
|
||||
return None
|
||||
return None if result["error"] else result["url"]
|
||||
|
||||
if result["error"]:
|
||||
return None
|
||||
|
||||
return result["url"]
|
||||
|
||||
# ── 下载 ───────────────────────────────────────────────────────────
|
||||
# ── 下载(SDK 走 internal endpoint)───────────────────────────────
|
||||
|
||||
def download_file(self, storage_key: str, local_path: str | Path) -> None:
|
||||
"""从 OSS 下载文件(简单下载,API端原有行为)。
|
||||
|
||||
bucket未配置 → 抛 RuntimeError
|
||||
"""
|
||||
if self.bucket is None:
|
||||
raise RuntimeError("OSS storage is not configured")
|
||||
|
||||
local_path = Path(local_path)
|
||||
os.makedirs(local_path.parent, exist_ok=True)
|
||||
try:
|
||||
self.bucket.get_object_to_file(self._normalize_storage_key(storage_key), str(local_path))
|
||||
self.bucket.get_object_to_file(
|
||||
self._normalize_storage_key(storage_key),
|
||||
str(local_path),
|
||||
)
|
||||
except Exception as e:
|
||||
raise Exception(f"Failed to download file from OSS: {e}") from e
|
||||
|
||||
def download_asset(self, asset_storage_key: str, local_path: str | Path) -> bool:
|
||||
"""下载素材(从 oss_helpers 合并)。
|
||||
|
||||
自动识别输入类型:
|
||||
- 完整 URL → 走 HTTP 下载(支持预签名URL)
|
||||
- 存储键 → 走 oss2 SDK 下载
|
||||
|
||||
成功返回 True,失败返回 False(不抛异常)。
|
||||
"""
|
||||
local_path = Path(local_path)
|
||||
os.makedirs(local_path.parent, exist_ok=True)
|
||||
|
||||
# 完整URL走HTTP下载(兼容预签名URL)
|
||||
if asset_storage_key.startswith(("http://", "https://")):
|
||||
return self._download_via_http(asset_storage_key, local_path)
|
||||
|
||||
# OSS存储键走SDK
|
||||
if self.bucket is None:
|
||||
logger.error("OSS not configured, cannot download: %s", asset_storage_key[:80])
|
||||
logger.error("OSS not configured: %s", asset_storage_key[:80])
|
||||
return False
|
||||
try:
|
||||
self.bucket.get_object_to_file(self._normalize_storage_key(asset_storage_key), str(local_path))
|
||||
self.bucket.get_object_to_file(
|
||||
self._normalize_storage_key(asset_storage_key),
|
||||
str(local_path),
|
||||
)
|
||||
return local_path.exists() and local_path.stat().st_size > 0
|
||||
except Exception:
|
||||
logger.exception("下载素材失败: %s", asset_storage_key)
|
||||
return False
|
||||
|
||||
def _download_via_http(self, url: str, local_path: Path) -> bool:
|
||||
"""通过 HTTP 下载文件(支持预签名 URL)。
|
||||
|
||||
流式下载避免大文件内存溢出。
|
||||
"""
|
||||
try:
|
||||
resp = requests.get(url, stream=True, timeout=OSS_HTTP_DOWNLOAD_TIMEOUT)
|
||||
resp.raise_for_status()
|
||||
@@ -298,83 +312,64 @@ class SharedStorageService(StoragePort):
|
||||
f.write(chunk)
|
||||
return local_path.exists() and local_path.stat().st_size > 0
|
||||
except Exception:
|
||||
logger.exception("HTTP下载素材失败: %s", url[:100])
|
||||
logger.exception("HTTP下载失败: %s", url[:100])
|
||||
return False
|
||||
|
||||
# ── URL 生成 ──────────────────────────────────────────────────────
|
||||
# ── URL 生成(sign_url 用 public_bucket 签公网域名)───────────────
|
||||
|
||||
def get_url(self, storage_key: str) -> str:
|
||||
"""获取公开 URL。"""
|
||||
return f"{self.public_url}/{storage_key}"
|
||||
|
||||
def get_download_url(self, storage_key_or_url: str, expires_seconds: int = 3600) -> str:
|
||||
"""获取预签名下载 URL。
|
||||
def _sign_bucket(self):
|
||||
"""签名优先用 public_bucket,回退到 bucket。"""
|
||||
return self.public_bucket or self.bucket
|
||||
|
||||
bucket未配置时降级为公开URL;本地产物URL直接返回。
|
||||
"""
|
||||
if self.bucket is None:
|
||||
def get_download_url(self, storage_key_or_url: str, expires_seconds: int = 3600) -> str:
|
||||
sign_bucket = self._sign_bucket()
|
||||
if sign_bucket is None:
|
||||
if self._is_local_generated_url(storage_key_or_url):
|
||||
return storage_key_or_url
|
||||
logger.warning(
|
||||
"get_download_url: OSS bucket not configured, returning raw URL. key=%s",
|
||||
storage_key_or_url[:200],
|
||||
)
|
||||
logger.warning("OSS bucket not configured, returning raw URL: %s", storage_key_or_url[:200])
|
||||
return self.get_url(self.normalize_storage_key(storage_key_or_url))
|
||||
|
||||
storage_key = self.normalize_storage_key(storage_key_or_url)
|
||||
try:
|
||||
signed = self.bucket.sign_url("GET", storage_key, expires_seconds)
|
||||
signed = sign_bucket.sign_url("GET", storage_key, expires_seconds)
|
||||
logger.info(
|
||||
"get_download_url: signed URL generated. key=%s url_prefix=%s",
|
||||
"signed URL generated for key=%s prefix=%s",
|
||||
storage_key[:80],
|
||||
signed[:60],
|
||||
)
|
||||
return signed
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"get_download_url: sign_url failed, falling back to raw URL. key=%s",
|
||||
storage_key[:200],
|
||||
)
|
||||
logger.exception("get_download_url: sign_url 失败,返回 raw URL: %s", storage_key[:200])
|
||||
return self.get_url(storage_key)
|
||||
|
||||
# ── 浏览器直传 POST ────────────────────────────────────────────────
|
||||
|
||||
def get_upload_url(
|
||||
self,
|
||||
storage_key_or_url: str,
|
||||
expires_seconds: int = 3600,
|
||||
content_type: str = "video/mp4",
|
||||
) -> str:
|
||||
"""获取预签名 PUT 上传 URL(供外部 Worker 上传结果文件)。
|
||||
|
||||
bucket未配置时降级为 public_url(本地/开发环境);
|
||||
本地产物 key 原样返回。
|
||||
"""
|
||||
if self.bucket is None:
|
||||
sign_bucket = self._sign_bucket()
|
||||
if sign_bucket is None:
|
||||
if self._is_local_generated_url(storage_key_or_url):
|
||||
return storage_key_or_url
|
||||
logger.warning(
|
||||
"get_upload_url: OSS bucket not configured, returning raw URL. key=%s",
|
||||
storage_key_or_url[:200],
|
||||
)
|
||||
logger.warning("get_upload_url: OSS 未配置,返回 raw URL: %s", storage_key_or_url[:200])
|
||||
return self.get_url(self.normalize_storage_key(storage_key_or_url))
|
||||
|
||||
storage_key = self.normalize_storage_key(storage_key_or_url)
|
||||
try:
|
||||
# oss2 sign_url 支持 'PUT',需指定 headers 才能限定 Content-Type
|
||||
headers = {"Content-Type": content_type} if content_type else None
|
||||
signed = self.bucket.sign_url("PUT", storage_key, expires_seconds, headers=headers)
|
||||
signed = sign_bucket.sign_url("PUT", storage_key, expires_seconds, headers=headers)
|
||||
logger.info(
|
||||
"get_upload_url: signed PUT URL generated. key=%s url_prefix=%s",
|
||||
"get_upload_url: 公网签名PUT URL已生成 key=%s prefix=%s",
|
||||
storage_key[:80],
|
||||
signed[:60],
|
||||
)
|
||||
return signed
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"get_upload_url: sign_url failed, falling back to raw URL. key=%s",
|
||||
storage_key[:200],
|
||||
)
|
||||
logger.exception("get_upload_url: sign_url 失败,返回 raw URL: %s", storage_key[:200])
|
||||
return self.get_url(storage_key)
|
||||
|
||||
def create_direct_upload_post(
|
||||
@@ -384,7 +379,6 @@ class SharedStorageService(StoragePort):
|
||||
max_size_bytes: int,
|
||||
expires_seconds: int,
|
||||
) -> dict[str, object]:
|
||||
"""创建浏览器直传 POST 表单。"""
|
||||
if not self.access_key_id or not self.access_key_secret:
|
||||
raise RuntimeError("OSS storage is not configured")
|
||||
normalized_key = self.normalize_storage_key(storage_key)
|
||||
@@ -431,19 +425,17 @@ class SharedStorageService(StoragePort):
|
||||
},
|
||||
}
|
||||
|
||||
# ── 文件操作 ───────────────────────────────────────────────────────
|
||||
# ── 文件操作(internal endpoint)──────────────────────────────────
|
||||
|
||||
def delete_file(self, storage_key: str) -> None:
|
||||
"""删除文件(不抛异常)。"""
|
||||
if self.bucket is None:
|
||||
return
|
||||
try:
|
||||
self.bucket.delete_object(storage_key)
|
||||
except Exception as error:
|
||||
logger.warning("Failed to delete file from OSS", extra={"storage_key": storage_key, "error": str(error)})
|
||||
logger.warning("OSS delete 失败", extra={"storage_key": storage_key, "error": str(error)})
|
||||
|
||||
def file_exists(self, storage_key: str) -> bool:
|
||||
"""检查文件是否存在。"""
|
||||
if self.bucket is None:
|
||||
return False
|
||||
return self.bucket.object_exists(storage_key)
|
||||
@@ -451,17 +443,6 @@ class SharedStorageService(StoragePort):
|
||||
# ── Asset 路径解析(Worker 用)────────────────────────────────────
|
||||
|
||||
def resolve_asset_path(self, asset_id: str, work_dir: str | Path) -> Optional[Path]:
|
||||
"""从 asset_id 解析到本地文件路径。
|
||||
|
||||
策略(按优先级):
|
||||
1. 本地绝对路径(在允许目录内)→ 直接返回
|
||||
2. work_dir 缓存命中 → 返回缓存路径
|
||||
3. 从OSS下载到缓存 → 返回下载路径
|
||||
4. 全部失败 → None
|
||||
|
||||
从 oss_helpers.resolve_asset_path 合并而来。
|
||||
"""
|
||||
# 延迟导入,避免循环依赖
|
||||
from video_processing.path_security import ( # type: ignore[import-not-found]
|
||||
PathSecurityError,
|
||||
get_allowed_local_dirs,
|
||||
@@ -471,47 +452,35 @@ class SharedStorageService(StoragePort):
|
||||
|
||||
if not asset_id or not isinstance(asset_id, str):
|
||||
return None
|
||||
|
||||
work_dir = Path(work_dir)
|
||||
os.makedirs(work_dir, exist_ok=True)
|
||||
|
||||
# 空字节检测
|
||||
if "\x00" in asset_id:
|
||||
logger.warning("asset_id 包含空字节,拒绝: %s", asset_id[:50])
|
||||
logger.warning("asset_id 含空字节,拒绝: %s", asset_id[:50])
|
||||
return None
|
||||
|
||||
# 1. 本地绝对路径 — 必须在允许的目录内
|
||||
if asset_id.startswith("/") and os.path.exists(asset_id):
|
||||
try:
|
||||
resolved = Path(asset_id).resolve()
|
||||
if is_in_allowed_dirs(resolved, get_allowed_local_dirs()):
|
||||
return resolved
|
||||
else:
|
||||
logger.warning(
|
||||
"本地素材路径不在允许目录内,拒绝: %s (allowed=%s)",
|
||||
asset_id[:80],
|
||||
get_allowed_local_dirs(),
|
||||
)
|
||||
return None
|
||||
logger.warning(
|
||||
"本地素材路径不在允许目录: %s allowed=%s",
|
||||
asset_id[:80],
|
||||
get_allowed_local_dirs(),
|
||||
)
|
||||
return None
|
||||
except (OSError, PathSecurityError):
|
||||
return None
|
||||
|
||||
# 2. 缓存命中(SHA256 hash 防路径遍历)
|
||||
cache_hash = hashlib.sha256(asset_id.encode()).hexdigest()[:16]
|
||||
safe_name = sanitize_filename(cache_hash)
|
||||
cached_path = work_dir / f"{safe_name}.mp4"
|
||||
if cached_path.exists() and cached_path.stat().st_size > 0:
|
||||
return cached_path
|
||||
|
||||
# 3. 从 OSS 下载(先标准化 key,防路径遍历注入)
|
||||
safe_key = self.normalize_storage_key(asset_id)
|
||||
if ".." in safe_key or safe_key.startswith("/"):
|
||||
logger.warning("asset_id 包含路径遍历模式,拒绝下载: %s", asset_id[:80])
|
||||
logger.warning("asset_id 含路径遍历: %s", asset_id[:80])
|
||||
return None
|
||||
|
||||
if self.download_asset(safe_key, cached_path):
|
||||
return cached_path
|
||||
|
||||
return None
|
||||
|
||||
def resolve_asset_ids_to_paths(
|
||||
@@ -519,22 +488,20 @@ class SharedStorageService(StoragePort):
|
||||
asset_ids: list[str],
|
||||
work_dir: str | Path,
|
||||
) -> dict[str, Path]:
|
||||
"""批量解析 asset_id → 本地路径。"""
|
||||
result: dict[str, Path] = {}
|
||||
for aid in asset_ids:
|
||||
local_path = self.resolve_asset_path(aid, work_dir)
|
||||
if local_path:
|
||||
result[aid] = local_path
|
||||
p = self.resolve_asset_path(aid, work_dir)
|
||||
if p:
|
||||
result[aid] = p
|
||||
return result
|
||||
|
||||
|
||||
# ── 单例管理 ────────────────────────────────────────────────────────────
|
||||
# ── 单例 ────────────────────────────────────────────────────────────────
|
||||
|
||||
_storage_service: Optional[SharedStorageService] = None
|
||||
|
||||
|
||||
def get_shared_storage_service() -> SharedStorageService:
|
||||
"""获取统一存储服务单例。"""
|
||||
global _storage_service
|
||||
if _storage_service is None:
|
||||
_storage_service = SharedStorageService()
|
||||
@@ -542,7 +509,6 @@ def get_shared_storage_service() -> SharedStorageService:
|
||||
return _storage_service
|
||||
|
||||
|
||||
# 向后兼容别名
|
||||
def get_storage_service() -> SharedStorageService:
|
||||
"""向后兼容:返回统一存储服务。"""
|
||||
"""向后兼容别名。"""
|
||||
return get_shared_storage_service()
|
||||
|
||||
@@ -666,7 +666,7 @@ class TestDownloadAssets:
|
||||
_make_clip("c2", order=1, asset_id="asset_002"),
|
||||
]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 2
|
||||
assert "asset_001" in asset_path_map
|
||||
@@ -693,7 +693,7 @@ class TestDownloadAssets:
|
||||
_make_clip("c2", order=1, asset_id="asset_002"),
|
||||
]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 1
|
||||
assert "asset_002" in asset_path_map
|
||||
@@ -714,7 +714,7 @@ class TestDownloadAssets:
|
||||
_make_clip("c1", order=0, asset_id="asset_001"),
|
||||
]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 0
|
||||
assert len(rendered_ids) == 0
|
||||
@@ -750,7 +750,7 @@ class TestDownloadAssets:
|
||||
_make_clip("c3", order=2, asset_id="asset_003"),
|
||||
]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 2
|
||||
assert "c1" in rendered_ids
|
||||
@@ -771,7 +771,7 @@ class TestDownloadAssets:
|
||||
_make_clip("c2", order=1, asset_id="asset_shared"),
|
||||
]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 1
|
||||
assert mock_download.call_count == 1
|
||||
@@ -795,7 +795,7 @@ class TestDownloadAssets:
|
||||
adapter = RenderAdapter(mock_db)
|
||||
clips = [_make_clip("c1", order=0, asset_id="asset_fallback")]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 1
|
||||
assert "c1" in rendered_ids
|
||||
@@ -820,7 +820,7 @@ class TestDownloadAssets:
|
||||
adapter = RenderAdapter(mock_db)
|
||||
clips = [_make_clip("c1", order=0, asset_id="asset_no_key")]
|
||||
|
||||
asset_path_map, rendered_ids, failed_ids = adapter._download_assets(clips, tmp_path)
|
||||
asset_path_map, rendered_ids, failed_ids, _storage_map = adapter._download_assets(clips, tmp_path)
|
||||
|
||||
assert len(asset_path_map) == 0
|
||||
assert "c1" in failed_ids
|
||||
|
||||
Reference in New Issue
Block a user