fix(#1714): 素材上传去重防重 + complete幂等传递 + 轮询收敛(前端P0) #1717

Merged
auto-approve-bot merged 2 commits from feat/1714-upload-dedup into develop 2026-09-05 16:40:06 +08:00
12 changed files with 720 additions and 66 deletions
+35 -3
View File
@@ -4,6 +4,7 @@
import apiClient from "../client"
import { getOrCreateDefaultProject } from "../projects"
import type { DirectUploadPrepareResult, DirectUploadCompleteResult } from "./types"
import { computeFileHash, makeClientUploadId } from "./uploadDedup"
/** 预签名直传准备 */
export const prepareDirectUpload = async (data: {
@@ -12,8 +13,13 @@ export const prepareDirectUpload = async (data: {
filename: string
content_type: string
file_size: number
/** 前端算好的文件内容哈希(SHA-256 hex),打开后端 file_hash 去重闸门 */
file_hash?: string
/** 前端生成的上传幂等 token,同一次逻辑上传(含重试)保持不变 */
client_upload_id?: string
}): Promise<DirectUploadPrepareResult> => {
const response = await apiClient.post("/upload/direct/prepare", data)
// prepare 单独放宽到 30s(全局 axios 实例只有 10sstaging 抖动时易超时)
const response = await apiClient.post("/upload/direct/prepare", data, { timeout: 30_000 })
return response.data
}
@@ -22,8 +28,14 @@ export const completeDirectUpload = async (data: {
project_id: string
library_id: string
storage_key: string
/** 前端算好的文件内容哈希(与 prepare 一致),后端按 hash 幂等去重 */
file_hash?: string
/** 前端上传幂等 token(与 prepare 一致),同一次上传重发 complete 不重复建记录 */
client_upload_id?: string
}): Promise<DirectUploadCompleteResult> => {
const response = await apiClient.post("/upload/direct/complete", data)
// complete 内含 OSS 存在性检查 + 建库 + 派单,放宽到 60s;
// 超时不代表失败(记录可能已建成),调用方禁止超时后盲目重传整个文件
const response = await apiClient.post("/upload/direct/complete", data, { timeout: 60_000 })
return response.data
}
@@ -109,6 +121,10 @@ export interface DirectUploadHandle {
export const prepareDirectUploadHandle = async (data: {
file: File
library_id: string
/** 前端算好的文件内容哈希(SHA-256 hex),prepare/complete 均携带 */
fileHash?: string
/** 本次逻辑上传的幂等 tokenprepare/complete 一致、重试复用 */
clientUploadId?: string
}): Promise<DirectUploadHandle> => {
const project = await getOrCreateDefaultProject()
@@ -118,6 +134,8 @@ export const prepareDirectUploadHandle = async (data: {
filename: data.file.name,
content_type: data.file.type || "application/octet-stream",
file_size: data.file.size,
file_hash: data.fileHash,
client_upload_id: data.clientUploadId,
})
return {
@@ -128,6 +146,8 @@ export const prepareDirectUploadHandle = async (data: {
project_id: project.id,
library_id: data.library_id,
storage_key: prepared.storage_key,
file_hash: data.fileHash,
client_upload_id: data.clientUploadId,
}),
}
}
@@ -137,8 +157,20 @@ export const uploadAssetDirect = async (data: {
file: File
library_id: string
onProgress?: (percent: number) => void
/** 文件内容哈希;未传时自动补算(配音/封面/克隆等非队列链路统一受益) */
fileHash?: string
/** 幂等 token;未传时自动生成 */
clientUploadId?: string
}): Promise<DirectUploadCompleteResult> => {
const handle = await prepareDirectUploadHandle({ file: data.file, library_id: data.library_id })
// 自动补算哈希与幂等 token:确保 file_hash 去重闸门对所有上传链路生效
const fileHash = data.fileHash ?? (await computeFileHash(data.file))
const clientUploadId = data.clientUploadId ?? makeClientUploadId()
const handle = await prepareDirectUploadHandle({
file: data.file,
library_id: data.library_id,
fileHash,
clientUploadId,
})
await handle.transfer(data.onProgress)
return handle.complete()
}
+143
View File
@@ -0,0 +1,143 @@
/**
* 上传去重 / 幂等工具(Issue #1714
*
* 背景:同一文件被反复入队、complete 超时后盲目重传,导致后端创建大量重复
* PROCESSING 素材记录。本模块提供两类纯函数:
*
* 1. 文件指纹:
* - makeFileFingerprint():文件名+大小+lastModified,入队去重用(同步、零开销)
* - computeFileHash()SHA-256 内容哈希(小文件全量、大文件抽样头尾),
* prepare/complete 时发给后端打开 file_hash 去重闸门
* 2. 队列去重:findDuplicateInQueue() 判断文件是否已在队列中
* 3. 幂等 tokenmakeClientUploadId() 生成上传幂等 ID(每次"一次逻辑上传"一个,
* 重试复用同一 ID,重新入队才生成新 ID)
*/
/** 大文件抽样阈值:超过此大小只哈希头尾片段,避免上传前长时间卡 UI */
export const HASH_FULL_READ_LIMIT = 256 * 1024 * 1024 // 256MB
/** 抽样读取的头尾片段大小(各 8MB) */
export const HASH_SAMPLE_CHUNK = 8 * 1024 * 1024
/** 计算指纹时,文件在队列中已存在的状态(已失败的可以重试,不算重复) */
export type DedupExcludeStatus = "error" | "done"
/**
* 文件入队指纹:同库 + 文件名 + 大小 + 修改时间。
* 同一文件(File 对象由 <input> 重选或拖拽重复触发时三个字段均一致)稳定复现;
* 不同文件极小概率碰撞时可由后端 file_hash 内容去重兜底。
*/
export function makeFileFingerprint(file: Pick<File, "name" | "size" | "lastModified">): string {
return `${file.name}::${file.size}::${file.lastModified}`
}
/**
* 在现有队列项中查找同一文件的在途记录。
* 已失败(error)的项允许重试路径复用、已完成(done)的可跳过;
* 处于 preparing/uploading/ingesting 的在途项一律视为重复,禁止重复入队。
*
* 返回命中的队列项 id(tempId),未命中返回 null。
*/
export function findDuplicateInQueue<T extends { fileKey: string; status: string }>(
queue: T[],
fileKey: string,
excludeStatuses: DedupExcludeStatus[] = [],
): T | null {
const exclude = new Set<string>(excludeStatuses)
return queue.find((it) => it.fileKey === fileKey && !exclude.has(it.status)) ?? null
}
/** 生成上传幂等 token:一次"逻辑上传"一个,重试复用、重新入队换新 */
export function makeClientUploadId(): string {
const rand =
typeof crypto !== "undefined" && "randomUUID" in crypto
? crypto.randomUUID()
: `${Date.now()}-${Math.random().toString(36).slice(2, 10)}-${Math.random()
.toString(36)
.slice(2, 10)}`
return `up_${Date.now().toString(36)}_${rand.replace(/-/g, "").slice(0, 16)}`
}
/** 读取 Blob/File 片段为 ArrayBuffer:优先 Blob.arrayBuffer(),老环境回退 FileReader */
function readAsArrayBuffer(blob: Blob): Promise<ArrayBuffer> {
if (typeof blob.arrayBuffer === "function") {
return blob.arrayBuffer()
}
return new Promise<ArrayBuffer>((resolve, reject) => {
const reader = new FileReader()
reader.onload = () => resolve(reader.result as ArrayBuffer)
reader.onerror = () => reject(reader.error ?? new Error("FileReader read failed"))
reader.readAsArrayBuffer(blob)
})
}
/**
* 把 buffer 复制到当前 JS realm 的 Uint8Array 再哈希。
* jsdom/测试环境中 Blob.arrayBuffer() 可能返回另一 realm 的 ArrayBuffer
* Node WebCrypto 的 WebIDL instanceof 校验会拒绝跨 realm 参数。
*/
async function digestSha256(buffer: ArrayBuffer): Promise<ArrayBuffer> {
const subtle =
typeof globalThis !== "undefined" && globalThis.crypto ? globalThis.crypto.subtle : null
if (!subtle) throw new Error("crypto.subtle unavailable")
const local = new Uint8Array(buffer.byteLength)
local.set(new Uint8Array(buffer))
return subtle.digest("SHA-256", local)
}
function toHex(buffer: ArrayBuffer): string {
const bytes = new Uint8Array(buffer)
let hex = ""
for (let i = 0; i < bytes.length; i += 1) {
hex += bytes[i].toString(16).padStart(2, "0")
}
return hex
}
/**
* 计算文件内容 SHA-256hex64 字符,与后端 file_hash 字段长度一致)。
* - ≤256MB:全量哈希,内容一致必然一致
* - >256MB:哈希「头部 8MB + 尾部 8MB + 文件大小」,视频素材体积大、
* 头部含 moov 元数据、尾部含 mdat 结尾,抽样碰撞概率可忽略,
* 且避免上传前对 2GB 文件全量读取造成长时间卡顿
*
* 运行环境不支持 crypto.subtle(非安全上下文/老浏览器)时返回空字符串,
* 调用方据此降级为不传 hash(后端仍有幂等 token + 同文件名兜底去重)。
*/
export async function computeFileHash(file: File): Promise<string> {
try {
const subtle =
typeof globalThis !== "undefined" &&
globalThis.crypto &&
typeof globalThis.crypto.subtle?.digest === "function"
? globalThis.crypto.subtle
: null
if (!subtle) return ""
if (file.size <= HASH_FULL_READ_LIMIT) {
const data = await readAsArrayBuffer(file.slice(0, file.size))
return toHex(await digestSha256(data))
}
// 大文件:头 8MB + 尾 8MB + 大小,拼成一段后哈希
const head = await readAsArrayBuffer(file.slice(0, HASH_SAMPLE_CHUNK))
const tail =
file.size > HASH_SAMPLE_CHUNK
? await readAsArrayBuffer(file.slice(Math.max(0, file.size - HASH_SAMPLE_CHUNK), file.size))
: new ArrayBuffer(0)
const merged = new Uint8Array(head.byteLength + tail.byteLength + 8)
merged.set(new Uint8Array(head), 0)
merged.set(new Uint8Array(tail), head.byteLength)
const sizeView = new DataView(merged.buffer, head.byteLength + tail.byteLength, 8)
// 文件大小以 64 位大端写入(BigInt 最稳;不支持 BigInt64 时手算高低位)
if (typeof sizeView.setBigUint64 === "function") {
sizeView.setBigUint64(0, BigInt(file.size), false)
} else {
sizeView.setUint32(0, Math.floor(file.size / 0x100000000), false)
sizeView.setUint32(4, file.size >>> 0, false)
}
return toHex(await digestSha256(merged.buffer))
} catch (err) {
console.warn("[uploadDedup] 计算文件哈希失败,降级为不传 file_hash:", err)
return ""
}
}
@@ -42,6 +42,7 @@ const AssetLibrary: React.FC = () => {
assetsError,
assetsErrorObj,
refetchAssets,
stalledAssetIds,
searchText,
setSearchText,
filterType,
@@ -77,6 +78,7 @@ const AssetLibrary: React.FC = () => {
removeUpload,
clearFinished,
uploading,
transferActive,
activeCount,
pendingCount,
} = useAssetUpload({ effectiveLibId })
@@ -168,6 +170,7 @@ const AssetLibrary: React.FC = () => {
{/* 上传区域 */}
<AssetUploadZone
uploading={uploading}
transferActive={transferActive}
activeCount={activeCount}
pendingCount={pendingCount}
onUpload={enqueueUploads}
@@ -214,6 +217,7 @@ const AssetLibrary: React.FC = () => {
selectedIds={selectedIds}
diagnosingId={diagnosingId}
uploadProgressMap={uploadProgressMap}
stalledAssetIds={stalledAssetIds}
onRetry={refetchAssets}
onToggleSelect={toggleSelect}
onDiagnose={handleDiagnose}
+22
View File
@@ -1041,3 +1041,25 @@
background: #fef2f2;
color: #dc2626;
}
/* 上传入口禁用态(直传进行中,防重复提交,Issue #1714 */
.xx-asset-upload-btn:disabled {
opacity: 0.6;
cursor: not-allowed;
}
.xx-asset-upload-btn:disabled:hover {
opacity: 0.6;
}
.xx-asset-upload-btn:disabled:active {
transform: none;
}
/* 处理超时遮罩:创建超过 10 分钟仍在处理中(疑似后端卡住),停止转圈并警示 */
.xx-asset-thumb-stalled {
background: rgba(217, 119, 6, 0.28);
color: #fde68a;
backdrop-filter: blur(2px);
}
.xx-asset-thumb-stalled :first-child {
font-size: var(--font-size-2xl);
}
@@ -22,6 +22,8 @@ export interface AssetCardProps {
diagnosing?: boolean
/** 上传中实时进度(仅 uploading 态有值;ingesting 后由后端状态接管) */
uploadProgress?: { progress: number; uploading: boolean }
/** 处理超过 10 分钟仍未就绪(疑似后端卡住):停止转圈并提示处理超时 */
stalled?: boolean
onToggle: () => void
onDiagnose: () => void
onPlay: () => void
@@ -33,6 +35,7 @@ const AssetCard: React.FC<AssetCardProps> = ({
selected,
diagnosing,
uploadProgress,
stalled,
onToggle,
onDiagnose,
onPlay,
@@ -67,11 +70,15 @@ const AssetCard: React.FC<AssetCardProps> = ({
</div>
)}
{/* 转码/处理中遮罩 */}
{/* 转码/处理中遮罩(卡死超过 10 分钟时停止转圈,提示超时) */}
{asset.loading && !isUploading && (
<div className="xx-asset-thumb-overlay xx-asset-thumb-processing">
<LoadingOutlined />
<span></span>
<div
className={`xx-asset-thumb-overlay ${
stalled ? "xx-asset-thumb-stalled" : "xx-asset-thumb-processing"
}`}
>
{stalled ? <CloseCircleOutlined /> : <LoadingOutlined />}
<span>{stalled ? "处理超时,可重试上传" : "转码处理中"}</span>
</div>
)}
@@ -129,7 +136,7 @@ const AssetCard: React.FC<AssetCardProps> = ({
</p>
<div className="xx-asset-meta">
<span className="xx-asset-meta-status">
<StatusPill status={asset.status} label={asset.statusLabel} />
<StatusPill status={asset.status} label={stalled ? "处理超时" : asset.statusLabel} />
</span>
{asset.duration && <span className="xx-asset-meta-duration">{asset.duration}</span>}
</div>
@@ -19,6 +19,8 @@ export interface AssetGridSectionProps {
selectedIds: Set<string>
diagnosingId: string | null
uploadProgressMap?: UploadProgressMap
/** 创建超过 10 分钟仍在处理中的素材 id(疑似后端卡住),卡片提示处理超时 */
stalledAssetIds?: Set<string>
onRetry?: () => void
onToggleSelect: (id: string) => void
onDiagnose: (asset: AssetItem) => void
@@ -34,6 +36,7 @@ export const AssetGridSection: React.FC<AssetGridSectionProps> = ({
selectedIds,
diagnosingId,
uploadProgressMap,
stalledAssetIds,
onRetry,
onToggleSelect,
onDiagnose,
@@ -76,6 +79,7 @@ export const AssetGridSection: React.FC<AssetGridSectionProps> = ({
selected={selectedIds.has(asset.id)}
diagnosing={diagnosingId === asset.id}
uploadProgress={uploadProgressMap?.get(asset.id)}
stalled={stalledAssetIds?.has(asset.id)}
onToggle={() => onToggleSelect(asset.id)}
onDiagnose={() => onDiagnose(asset)}
onPlay={() => onPlay(asset)}
@@ -4,10 +4,13 @@
* - 拖拽文件到内容区任意位置同样触发上传(不再占用大面积虚线框)
*/
import React, { useRef, useState } from "react"
import { message } from "antd"
import { PlusOutlined, CloudUploadOutlined } from "@ant-design/icons"
export interface AssetUploadZoneProps {
uploading: boolean
/** 有文件正在本地指纹/prepare/直传(非服务端转码),此时禁用入口防重复提交 */
transferActive: boolean
activeCount: number
pendingCount: number
onUpload: (files: File[]) => void
@@ -15,6 +18,7 @@ export interface AssetUploadZoneProps {
export const AssetUploadZone: React.FC<AssetUploadZoneProps> = ({
uploading,
transferActive,
activeCount,
pendingCount,
onUpload,
@@ -27,6 +31,11 @@ export const AssetUploadZone: React.FC<AssetUploadZoneProps> = ({
const pickFiles = (list: FileList | null) => {
if (!list || list.length === 0) return
// 直传进行中拦截重复触发:相同文件仍由入队去重兜底,这里先给明确反馈
if (transferActive) {
message.warning("文件正在上传中,请等待当前上传完成后再添加")
return
}
onUpload(Array.from(list))
}
@@ -58,10 +67,18 @@ export const AssetUploadZone: React.FC<AssetUploadZoneProps> = ({
<button
type="button"
className="xx-asset-upload-btn"
onClick={() => inputRef.current?.click()}
disabled={transferActive}
title={transferActive ? "文件上传中,暂不能添加新文件" : undefined}
onClick={() => {
if (transferActive) {
message.warning("文件正在上传中,请等待当前上传完成后再添加")
return
}
inputRef.current?.click()
}}
>
<PlusOutlined />
{transferActive ? "上传中…" : "上传素材"}
</button>
<span className="xx-asset-upload-status">
{uploading ? (
@@ -80,6 +80,7 @@ const UploadQueuePanel: React.FC<UploadQueuePanelProps> = ({
) : null}
<div className="xx-upload-queue-status">
{it.duplicated ? "素材已存在,已跳过" : STATUS_TEXT[it.status]}
{it.status === "preparing" && it.hint ? `${it.hint}` : ""}
{it.status === "uploading" ? ` ${it.progress}%` : ""}
{it.status === "error" && it.error ? `${it.error}` : ""}
</div>
@@ -89,7 +90,9 @@ const UploadQueuePanel: React.FC<UploadQueuePanelProps> = ({
<button
type="button"
className="xx-upload-queue-btn"
title="重试"
title={
it.failedStage === "complete" ? "安全重试(只确认,不重新上传)" : "重试上传"
}
onClick={() => onRetry(it.tempId)}
>
<ReloadOutlined />
+187 -32
View File
@@ -3,33 +3,57 @@ import { useQueryClient } from "@tanstack/react-query"
import { message } from "antd"
import { prepareDirectUploadHandle, type DirectUploadHandle } from "@/api/assets"
import { MAX_FILE_SIZE } from "../constants"
import {
computeFileHash,
findDuplicateInQueue,
makeClientUploadId,
makeFileFingerprint,
} from "@/api/assets/uploadDedup"
/** 单文件上传状态机 */
export type UploadItemStatus = "preparing" | "uploading" | "ingesting" | "done" | "error"
/** 失败发生的阶段:complete 阶段失败时记录可能已在后端建成,禁止盲目重传整个文件 */
export type UploadFailStage = "prepare" | "transfer" | "complete"
export interface UploadItem {
/** 前端临时 idprepare 前无 asset_id 时用) */
/** 前端临时 idprepare 前无 asset_id 时用),同时作为队列项 key */
tempId: string
file: File
fileName: string
/** 进度 0~100(仅直传阶段有真实进度) */
progress: number
status: UploadItemStatus
/** 文件指纹(name+size+lastModified),入队去重用 */
fileKey: string
/** 本次逻辑上传的幂等 token:重试复用、重新入队才换新 */
clientUploadId: string
/** 上传前算好的文件内容哈希(SHA-256),prepare/complete 都带上 */
fileHash?: string
/** 后端 prepare 预建的 asset id(旧后端可能为空) */
assetId?: string
/** 去重命中:complete 返回 duplicated,标记完成但不产生新素材 */
duplicated?: boolean
/** 失败发生的阶段;complete 阶段失败点重试只重发 complete,不重新上传文件 */
failedStage?: UploadFailStage
/** 状态行补充提示(如"正在计算文件指纹…" */
hint?: string
error?: string
}
/** 批量直传最大并发数,避免多文件瓜分上行带宽 */
const MAX_CONCURRENT = 3
/** complete 阶段失败后的错误提示:素材可能已在服务器处理中,重试不会重新上传 */
const COMPLETE_ERROR_HINT =
"确认请求失败,素材可能已在服务器处理中;点重试将安全确认,不会重新上传文件"
/**
* 素材批量上传 Hook
* - prepare 阶段后端预建 status=uploading 的 asset,前端拿到 asset_id 立即刷新列表
* - 入队按文件指纹(name+size+lastModified)去重:同一文件已在队列/上传中/处理中时不重复入队
* - 上传前计算文件 SHA-256prepare/complete 携带 file_hash + 幂等 tokenclientUploadId
* - complete 超时/失败不盲目重传:复用 handle 只重发 complete(幂等),prepare/transfer 失败才全量重跑
* - OSS 直传并发限制为 3,其余排队;每个文件独立进度/状态
* - complete 后素材进入转码(ingesting/processing),由列表轮询反映
* - 失败卡片支持重试/移除
*/
export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
@@ -39,6 +63,13 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
const itemsRef = useRef<UploadItem[]>([])
itemsRef.current = items
/**
* prepare 成功后的 handle 按 tempId 留存:
* complete 阶段失败(超时/网络)时 OSS 文件已存在、后端记录也可能已建成,
* 重试必须复用同一 handle 只重发 complete,绝不能重新 prepare+直传。
*/
const handlesRef = useRef<Map<string, DirectUploadHandle>>(new Map())
const updateItem = useCallback((tempId: string, patch: Partial<UploadItem>) => {
setItems((prev) => prev.map((it) => (it.tempId === tempId ? { ...it, ...patch } : it)))
}, [])
@@ -52,14 +83,57 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
queryClient.invalidateQueries({ queryKey: ["asset-libraries"] })
}, [queryClient, effectiveLibId])
/** 执行单个文件的完整上传流程(prepare→transfer→complete */
/**
* 执行单个文件的完整上传流程。
* @param completeOnly complete 阶段失败后的重试:跳过 hash/prepare/transfer,只重发 complete
* (OSS 文件已传完,重发由 file_hash + clientUploadId 保证幂等)
*/
const runUpload = useCallback(
async (item: UploadItem, handle?: DirectUploadHandle) => {
async (item: UploadItem, handle?: DirectUploadHandle, completeOnly = false) => {
let stage: UploadFailStage = "prepare"
try {
// 1. prepare(重试时复用已准备的 handle 也行,但签名可能过期,重新 prepare 最稳)
const h =
handle ??
(await prepareDirectUploadHandle({ file: item.file, library_id: effectiveLibId }))
let h = handle
if (completeOnly && h) {
// ── complete 重试:文件已在 OSS,直接幂等重发确认 ──
stage = "complete"
updateItem(item.tempId, {
status: "ingesting",
progress: 100,
error: undefined,
failedStage: undefined,
hint: undefined,
})
const result = await h.complete()
refreshList()
handlesRef.current.delete(item.tempId)
if (result.duplicated) {
updateItem(item.tempId, { status: "done", duplicated: true, assetId: result.asset_id })
message.info(`"${item.fileName}" 与素材库已有内容相同,已跳过`)
} else {
updateItem(item.tempId, { status: "done", assetId: result.asset_id || item.assetId })
message.success(`"${item.fileName}" 上传完成,正在转码处理`)
}
return
}
// 1. 计算文件内容哈希(失败不阻塞,降级为不传 hash;后端仍有幂等 token 兜底)
updateItem(item.tempId, { hint: "正在计算文件指纹…" })
const fileHash = item.fileHash || (await computeFileHash(item.file))
updateItem(item.tempId, { fileHash, hint: undefined })
// 2. prepare(携带 file_hash + 幂等 token;重试时复用同一 clientUploadId
stage = "prepare"
h =
h ??
(await prepareDirectUploadHandle({
file: item.file,
library_id: effectiveLibId,
fileHash,
clientUploadId: item.clientUploadId,
}))
handlesRef.current.set(item.tempId, h)
if (h.prepared.asset_id) {
updateItem(item.tempId, {
status: "uploading",
@@ -72,26 +146,50 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
updateItem(item.tempId, { status: "uploading", progress: 0 })
}
// 2. OSS 直传(真实进度)
// 3. OSS 直传(真实进度)
stage = "transfer"
await h.transfer((pct) => updateItem(item.tempId, { progress: pct }))
// 3. complete:后端创建 ingest job,素材进入转码
// 4. complete:后端确认入库并创建 ingest jobfile_hash + 幂等 token 已在 handle 闭包中)
stage = "complete"
updateItem(item.tempId, { status: "ingesting", progress: 100 })
const result = await h.complete()
refreshList()
handlesRef.current.delete(item.tempId)
if (result.duplicated) {
updateItem(item.tempId, { status: "done", duplicated: true, assetId: result.asset_id })
message.info(`"${item.fileName}" 与素材库已有内容相同,已跳过`)
} else {
updateItem(item.tempId, { status: "done" })
updateItem(item.tempId, { status: "done", assetId: result.asset_id || item.assetId })
message.success(`"${item.fileName}" 上传完成,正在转码处理`)
}
} catch (err: unknown) {
const detail = err instanceof Error ? err.message : "上传失败"
console.error("[useAssetUpload] 上传失败:", item.fileName, err)
updateItem(item.tempId, { status: "error", error: detail })
message.error(`"${item.fileName}" 上传失败:${detail}`)
console.error("[useAssetUpload] 上传失败:", item.fileName, stage, err)
if (stage === "complete") {
// complete 失败(超时/5xx/网络):后端记录可能已建成,handle 保留供幂等重试;
// 刷新列表让用户看到可能已创建的「处理中」素材,避免误以为没传上去而重复操作
refreshList()
updateItem(item.tempId, {
status: "error",
failedStage: "complete",
error: COMPLETE_ERROR_HINT,
hint: undefined,
})
message.error(`"${item.fileName}" ${COMPLETE_ERROR_HINT}`)
} else {
// prepare / transfer 失败:后端尚无素材记录,可安全全量重跑
handlesRef.current.delete(item.tempId)
updateItem(item.tempId, {
status: "error",
failedStage: stage === "transfer" ? "transfer" : "prepare",
error: detail,
hint: undefined,
})
message.error(`"${item.fileName}" 上传失败:${detail}`)
}
}
},
[effectiveLibId, refreshList, updateItem],
@@ -113,7 +211,9 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
if (!next) return
claimedRef.current.add(next.tempId)
inFlightRef.current += 1
void runUpload(next).finally(() => {
// complete 阶段失败的重试:复用留存的 handle,只重发 complete
const existingHandle = handlesRef.current.get(next.tempId)
void runUpload(next, existingHandle, existingHandle !== undefined).finally(() => {
inFlightRef.current -= 1
claimedRef.current.delete(next.tempId)
// 一个任务结束(成功/失败)后继续拉起排队任务
@@ -126,7 +226,11 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
pumpRef.current()
}, [items])
/** 入队一个或多个文件 */
/**
* 入队一个或多个文件(按文件指纹去重):
* - 同一文件已在队列且 preparing/uploading/ingesting/done → 跳过,不重复入队
* - 同一文件此前失败(error)→ 重新激活原队列项(复用 clientUploadId,保持幂等语义)
*/
const enqueueUploads = useCallback(
(files: File[]) => {
if (!effectiveLibId) {
@@ -143,37 +247,86 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
}
if (valid.length === 0) return
const newItems: UploadItem[] = valid.map((file, idx) => ({
tempId: `${Date.now()}-${idx}-${Math.random().toString(36).slice(2, 8)}`,
file,
fileName: file.name,
progress: 0,
status: "preparing",
}))
setItems((prev) => [...prev, ...newItems])
let skipped = 0
let rearmed = 0
const newItems: UploadItem[] = []
for (const file of valid) {
const fileKey = makeFileFingerprint(file)
// error 项允许重新激活;其余状态(preparing/uploading/ingesting/done)都算重复
const dup = findDuplicateInQueue([...itemsRef.current, ...newItems], fileKey, ["error"])
if (dup) {
skipped += 1
continue
}
// 失败项重新激活:复用 tempId/clientUploadId,由 pump 按留存 handle 决定重试方式
const failed = itemsRef.current.find(
(it) => it.fileKey === fileKey && it.status === "error",
)
if (failed) {
rearmed += 1
updateItem(failed.tempId, {
status: "preparing",
progress: 0,
error: undefined,
failedStage: undefined,
hint: undefined,
})
continue
}
newItems.push({
tempId: `${Date.now()}-${newItems.length}-${Math.random().toString(36).slice(2, 8)}`,
file,
fileName: file.name,
progress: 0,
status: "preparing",
fileKey,
clientUploadId: makeClientUploadId(),
})
}
if (newItems.length > 0) {
setItems((prev) => [...prev, ...newItems])
}
if (skipped > 0) {
message.warning(`已跳过 ${skipped} 个重复文件(已在上传队列、处理中或本页已上传)`)
}
if (rearmed > 0) {
message.info(`已重新加入 ${rearmed} 个此前失败的文件`)
}
},
[effectiveLibId],
[effectiveLibId, updateItem],
)
/** 重试失败任务 */
/**
* 重试失败任务(仅限 status=error):
* - complete 阶段失败:复用留存 handle 只重发 complete(幂等,不重新上传)
* - prepare/transfer 阶段失败:全量重跑(后端尚无记录,安全)
*/
const retryUpload = useCallback(
(tempId: string) => {
const target = itemsRef.current.find((it) => it.tempId === tempId)
if (!target) return
updateItem(tempId, { status: "preparing", progress: 0, error: undefined })
// 状态更新后由 useEffect 触发 pump
if (!target || target.status !== "error") return
updateItem(tempId, { status: "preparing", progress: 0, error: undefined, hint: undefined })
// 状态更新后由 useEffect 触发 pumppump 会按 handlesRef 自动选择 completeOnly / 全量
},
[updateItem],
)
/** 从上传列表移除(已进入转码的由素材网格管理;这里只移除上传面板记录) */
const removeUpload = useCallback((tempId: string) => {
handlesRef.current.delete(tempId)
setItems((prev) => prev.filter((it) => it.tempId !== tempId))
}, [])
/** 清空已完成/去重记录 */
const clearFinished = useCallback(() => {
setItems((prev) => prev.filter((it) => it.status !== "done"))
setItems((prev) => {
for (const it of prev) {
if (it.status === "done") handlesRef.current.delete(it.tempId)
}
return prev.filter((it) => it.status !== "done")
})
}, [])
const activeCount = items.filter(
@@ -188,8 +341,10 @@ export function useAssetUpload({ effectiveLibId }: { effectiveLibId: string }) {
retryUpload,
removeUpload,
clearFinished,
/** 是否有进行中的上传(用于上传区文案) */
/** 是否有进行中的上传(用于上传区文案/禁用入口 */
uploading: hasActive,
/** 是否有文件正在本地处理或直传(用于禁用上传入口,防重复提交) */
transferActive: activeCount > 0,
activeCount,
pendingCount,
}
@@ -10,6 +10,26 @@ import {
import { getOrCreateDefaultProject } from "@/api/projects"
import { mapLibrary, mapAsset, type AssetItem, type LibraryItem } from "../types"
/** 处理中素材快速轮询(3s)的最大持续时间:超过后停止快轮询,避免孤儿任务永久转圈 */
const PROCESSING_POLL_MAX_MS = 10 * 60 * 1000 // 10 分钟
const isProcessingStatus = (st?: string | null): boolean =>
st === "uploading" || st === "ingesting" || st === "processing" || st === "pending"
/** 判断列表是否存在「创建超过 maxMs 仍在处理中」的卡死素材 */
const hasStalledProcessing = (
list: ApiAssetItem[],
maxMs: number = PROCESSING_POLL_MAX_MS,
): boolean => {
const now = Date.now()
return list.some((a) => {
if (!isProcessingStatus(a.status ?? "")) return false
if (!a.created_at) return false
const created = new Date(a.created_at).getTime()
return Number.isFinite(created) && now - created > maxMs
})
}
/**
* 素材库数据 Hook
* 封装视频库列表、素材列表的数据查询,以及筛选、搜索状态管理
@@ -63,15 +83,15 @@ export function useAssetsData() {
}),
enabled: !!effectiveLibId,
staleTime: 30_000,
// 列表中存在上传中/转码中素材时每 3s 轮询;全部就绪后自动停止
// 列表中存在上传中/转码中素材时每 3s 轮询;全部就绪后自动停止
// 但若处理中素材创建已超过 10 分钟仍未就绪(疑似后端卡住/孤儿任务),
// 停止快轮询避免无限转圈——卡死素材在网格中显示「处理超时」提示。
refetchInterval: (query) => {
const data = query.state.data as { items: ApiAssetItem[] } | undefined
const items = data?.items ?? []
const processing = items.some((a) => {
const st = a.status ?? ""
return st === "uploading" || st === "ingesting" || st === "processing" || st === "pending"
})
return processing ? 3000 : false
const processing = items.some((a) => isProcessingStatus(a.status ?? ""))
if (!processing) return false
return hasStalledProcessing(items) ? false : 3000
},
})
@@ -80,6 +100,21 @@ export function useAssetsData() {
[apiAssets],
)
/** 创建超过 10 分钟仍在处理中的素材(后端可能卡住),网格提示「处理超时」 */
const stalledAssetIds = useMemo(() => {
const list = Array.isArray(apiAssets?.items) ? apiAssets.items : []
const ids = new Set<string>()
const now = Date.now()
for (const a of list) {
if (!isProcessingStatus(a.status ?? "") || !a.created_at) continue
const created = new Date(a.created_at).getTime()
if (Number.isFinite(created) && now - created > PROCESSING_POLL_MAX_MS) {
ids.add(a.id)
}
}
return ids
}, [apiAssets])
/* ── 筛选状态 ── */
const [searchText, setSearchText] = useState("")
const [filterType, setFilterType] = useState<string>("all")
@@ -129,6 +164,8 @@ export function useAssetsData() {
assetsError,
assetsErrorObj,
refetchAssets,
stalledAssetIds,
hasStalledAssets: stalledAssetIds.size > 0,
// 筛选
searchText,
setSearchText,
+101
View File
@@ -0,0 +1,101 @@
/**
* 上传去重/幂等工具单测(Issue #1714
*/
import { describe, it, expect } from "vitest"
import {
computeFileHash,
findDuplicateInQueue,
makeClientUploadId,
makeFileFingerprint,
} from "@/api/assets/uploadDedup"
const makeFile = (name: string, size = 100, lastModified = 1_700_000_000_000) =>
new File([new Uint8Array(size)], name, { type: "video/mp4", lastModified })
describe("makeFileFingerprint", () => {
it("同一文件(name+size+lastModified 相同)指纹一致", () => {
const a = makeFile("a.mp4", 1000, 12345)
const b = makeFile("a.mp4", 1000, 12345)
expect(makeFileFingerprint(a)).toBe(makeFileFingerprint(b))
})
it("文件名/大小/修改时间任一不同指纹即不同", () => {
const base = makeFile("a.mp4", 1000, 100)
expect(makeFileFingerprint(base)).not.toBe(makeFileFingerprint(makeFile("b.mp4", 1000, 100)))
expect(makeFileFingerprint(base)).not.toBe(makeFileFingerprint(makeFile("a.mp4", 1001, 100)))
expect(makeFileFingerprint(base)).not.toBe(makeFileFingerprint(makeFile("a.mp4", 1000, 101)))
})
})
describe("findDuplicateInQueue", () => {
const queue = [
{ fileKey: "k1", status: "preparing" },
{ fileKey: "k2", status: "uploading" },
{ fileKey: "k3", status: "ingesting" },
{ fileKey: "k4", status: "done" },
{ fileKey: "k5", status: "error" },
]
it("在途状态(preparing/uploading/ingesting/done)命中重复", () => {
expect(findDuplicateInQueue(queue, "k1")?.status).toBe("preparing")
expect(findDuplicateInQueue(queue, "k2")?.status).toBe("uploading")
expect(findDuplicateInQueue(queue, "k3")?.status).toBe("ingesting")
expect(findDuplicateInQueue(queue, "k4")?.status).toBe("done")
})
it("未命中返回 null", () => {
expect(findDuplicateInQueue(queue, "missing")).toBeNull()
})
it("排除 error 状态后,失败项不算重复(允许重新激活)", () => {
expect(findDuplicateInQueue(queue, "k5", ["error"])).toBeNull()
})
it("同时排除 done 后,已完成项也不算重复", () => {
expect(findDuplicateInQueue(queue, "k4", ["error", "done"])).toBeNull()
// 但在途的仍然命中
expect(findDuplicateInQueue(queue, "k1", ["error", "done"])).not.toBeNull()
})
})
describe("makeClientUploadId", () => {
it("生成带前缀且互不相同的幂等 token", () => {
const ids = new Set(Array.from({ length: 20 }, () => makeClientUploadId()))
expect(ids.size).toBe(20)
for (const id of ids) expect(id.startsWith("up_")).toBe(true)
})
})
describe("computeFileHash", () => {
it("相同内容 hash 一致、不同内容 hash 不同", async () => {
const f1 = makeFile("a.mp4", 4096)
const f2 = makeFile("b.mp4", 4096)
// 两个文件都是 0 填充,内容相同 → hash 一致
expect(await computeFileHash(f1)).toBe(await computeFileHash(f2))
const f3 = new File([new Uint8Array(4096).fill(7)], "c.mp4", { type: "video/mp4" })
expect(await computeFileHash(f1)).not.toBe(await computeFileHash(f3))
})
it("返回 64 位十六进制(SHA-256,与后端 file_hash 长度一致)", async () => {
const hash = await computeFileHash(makeFile("a.mp4", 1024))
expect(hash).toMatch(/^[0-9a-f]{64}$/)
})
})
describe("computeFileHash 大文件抽样(>256MB", () => {
it("抽样路径正常返回 64 位 hex,且大小不同则 hash 不同", async () => {
// mock 一个「声称」300MB 的 Fileslice 返回小 buffer 即可,不真分配 300MB
const makeBig = (declaredSize: number, head: number) => {
const f = new File([new Uint8Array([head, 2, 3])], "big.mov", { type: "video/quicktime" })
Object.defineProperty(f, "size", { value: declaredSize, configurable: true })
// slice 仍按真实内容返回小片段(头尾片段内容由底层小 buffer 决定)
return f
}
const h1 = await computeFileHash(makeBig(300 * 1024 * 1024, 1))
const h2 = await computeFileHash(makeBig(301 * 1024 * 1024, 1))
expect(h1).toMatch(/^[0-9a-f]{64}$/)
// 声明大小不同 → 写入的 64 位 size 字段不同 → hash 必须不同(锁定 setBigUint64 路径)
expect(h1).not.toBe(h2)
})
})
@@ -31,17 +31,23 @@ interface FakeHandle {
complete: ReturnType<typeof vi.fn>
/** 手动结束传输(transfer 被调用后挂载);finish(true) 以失败结束 */
finish: (fail?: boolean) => void
/** complete 已被调用的次数 */
completeCalls: { resolve: () => void; reject: (err: unknown) => void }[]
}
let activeTransfers = 0
let maxConcurrent = 0
/**
* 创建一个假 handletransfer 返回挂起的 promise
* finish 槽位在 transfer executor 同步执行时挂载,测试中调用 finish() 控制成败
* 创建一个假 handle
* - transfer 返回挂起的 promisefinish()/finish(true) 控制成败
* - complete 每次调用返回独立的挂起 promise,由 completeCalls 记录控制,
* 成功调 resolve(idx) / 失败调 reject(idx)(模拟超时)
*/
const makeFakeHandle = (opts: { id: string; duplicated?: boolean; failTransfer?: boolean }) => {
const h = {
const makeFakeHandle = (opts: {
id: string
duplicated?: boolean
failTransfer?: boolean
completeAuto?: boolean
}) => {
const h: FakeHandle = {
prepared: {
upload_url: "https://oss.example.com/u",
method: "POST",
@@ -52,15 +58,39 @@ const makeFakeHandle = (opts: { id: string; duplicated?: boolean; failTransfer?:
asset_id: opts.id,
},
transfer: vi.fn(),
complete: vi.fn().mockResolvedValue({
storage_key: "uploads/x/y.mp4",
ingest_job_id: opts.duplicated ? "" : "job-1",
url: "https://oss.example.com/u",
duplicated: opts.duplicated,
asset_id: opts.id,
}),
finish: (() => {}) as (fail?: boolean) => void,
complete: vi.fn(),
finish: () => {},
completeCalls: [],
}
h.complete.mockImplementation(
() =>
new Promise<{
storage_key: string
ingest_job_id: string
url: string
duplicated: boolean
asset_id: string
}>((resolve, reject) => {
h.completeCalls.push({
resolve: () =>
resolve({
storage_key: "uploads/x/y.mp4",
ingest_job_id: opts.duplicated ? "" : `job-${opts.id}`,
url: "https://oss.example.com/u",
duplicated: !!opts.duplicated,
asset_id: opts.id,
}),
reject,
})
// 默认立即成功,保持旧用例简单
if (opts.completeAuto !== false) {
const idx = h.completeCalls.length - 1
Promise.resolve().then(() => h.completeCalls[idx]?.resolve())
}
}),
)
h.transfer.mockImplementation(
() =>
new Promise<void>((_resolve, reject) => {
@@ -78,17 +108,24 @@ const makeFakeHandle = (opts: { id: string; duplicated?: boolean; failTransfer?:
type FakeHandleLike = ReturnType<typeof makeFakeHandle>
let activeTransfers = 0
let maxConcurrent = 0
/** prepare mock:调用序号生成稳定 id,立即把 handle(含 finish 槽位)推入数组 */
const installPrepareMock = (
handles: FakeHandleLike[],
optOverrides?: (id: string) => { duplicated?: boolean; failTransfer?: boolean },
optOverrides?: (id: string) => {
duplicated?: boolean
failTransfer?: boolean
completeAuto?: boolean
},
) => {
let callNo = 0
;(prepareDirectUploadHandle as unknown as ReturnType<typeof vi.fn>).mockImplementation(
async () => {
const id = `asset-${callNo++}`
const overrides = optOverrides?.(id) ?? {}
const h = makeFakeHandle({ id, ...overrides })
const h = makeFakeHandle({ id, completeAuto: true, ...overrides })
handles.push(h)
await new Promise((r) => setTimeout(r, 10))
return h
@@ -223,4 +260,96 @@ describe("useAssetUpload", () => {
expect(result.current.uploadItems[0].status).toBe("done")
})
})
it("同一文件多次选择不重复入队(指纹去重)", async () => {
const handles: FakeHandleLike[] = []
installPrepareMock(handles)
const { result } = renderHook(() => useAssetUpload({ effectiveLibId: "lib-1" }), {
wrapper: createWrapper(),
})
// 同一文件(name+size+lastModified 完全一致)第一次入队
const sameFile = mp4("same.mp4")
await act(async () => {
result.current.enqueueUploads([sameFile])
})
await waitFor(() => expect(handles.length).toBe(1))
expect(result.current.uploadItems).toHaveLength(1)
// transfer 挂起期间,再次选择同一文件(模拟用户反复点选/拖拽)
await act(async () => {
result.current.enqueueUploads([sameFile])
})
await act(async () => {
result.current.enqueueUploads([sameFile])
})
// 队列表只有 1 项、prepare 只有 1 次
expect(result.current.uploadItems).toHaveLength(1)
expect(handles.length).toBe(1)
// 完成后再次重复选择(已 done):仍然不新增
await act(async () => {
handles[0].finish()
})
await waitFor(() => expect(result.current.uploadItems[0].status).toBe("done"))
await act(async () => {
result.current.enqueueUploads([sameFile])
})
expect(result.current.uploadItems).toHaveLength(1)
expect(handles.length).toBe(1)
// 不同文件正常入队
await act(async () => {
result.current.enqueueUploads([mp4("other.mp4")])
})
await waitFor(() => expect(handles.length).toBe(2))
expect(result.current.uploadItems).toHaveLength(2)
})
it("complete 失败(超时)后重试:只重发 complete,不重新 prepare/直传", async () => {
const handles: FakeHandleLike[] = []
installPrepareMock(handles, () => ({ completeAuto: false }))
const { result } = renderHook(() => useAssetUpload({ effectiveLibId: "lib-1" }), {
wrapper: createWrapper(),
})
await act(async () => {
result.current.enqueueUploads([mp4("slow.mp4")])
})
await waitFor(() => expect(handles.length).toBe(1))
await waitFor(() => expect(handles[0].transfer).toHaveBeenCalled())
await act(async () => {
handles[0].finish()
})
// complete 被调用但挂起
await waitFor(() => expect(handles[0].complete).toHaveBeenCalledTimes(1))
// 模拟 complete 超时(后端记录可能已建成)
await act(async () => {
handles[0].completeCalls[0]?.reject(new Error("complete timeout (ECONNABORTED)"))
})
const tempId = result.current.uploadItems[0].tempId
await waitFor(() => {
const it = result.current.uploadItems.find((x) => x.tempId === tempId)
expect(it?.status).toBe("error")
expect(it?.failedStage).toBe("complete")
})
// 点重试:pump 复用 handle,只再调一次 completetransfer/prepare 不重复)
await act(async () => {
result.current.retryUpload(tempId)
})
await waitFor(() => expect(handles[0].complete).toHaveBeenCalledTimes(2))
expect(handles.length).toBe(1) // 没有重新 prepare
expect(handles[0].transfer).toHaveBeenCalledTimes(1) // 没有重新直传
// 第二次 complete 成功
await act(async () => {
handles[0].completeCalls[1]?.resolve()
})
await waitFor(() => {
expect(result.current.uploadItems.find((x) => x.tempId === tempId)?.status).toBe("done")
})
})
})