import crc32 from './crc32.js' import http from 'axios' import express from 'express' const app = express() const port = 7860 app.use(express.json()) app.use(express.urlencoded({ extended: true })) // 带取消令牌的请求函数 const createRequestWithCancel = (cancelToken) => { return { get: (url, config = {}) => http.get(url, { ...config, cancelToken }), post: (url, data, config = {}) => http.post(url, data, { ...config, cancelToken }), head: (url, config = {}) => http.head(url, { ...config, cancelToken }) } } // 带重试功能的分片下载函数 async function downloadChunkWithRetry(request, finalUrl, start, end, partNumber, maxRetries = 30) { let lastError = null for (let attempt = 1; attempt <= maxRetries; attempt++) { try { const response = await request.get(finalUrl, { headers: { 'Range': `bytes=${start}-${end}` }, responseType: 'arraybuffer' }) if (response.status !== 206 && response.status !== 200) { throw new Error(`HTTP ${response.status}: ${response.statusText}`) } const arrayBuffer = response.data if (arrayBuffer.byteLength !== (end - start + 1)) { throw new Error(`分片大小不匹配: 期望 ${end - start + 1} 字节, 实际 ${arrayBuffer.byteLength} 字节`) } console.log(`分片 ${partNumber} 下载成功 (尝试 ${attempt})`) return arrayBuffer } catch (error) { // 如果是取消请求,直接抛出 if (http.isCancel(error)) { throw error } lastError = error console.warn(`分片 ${partNumber} 下载失败 (尝试 ${attempt}/${maxRetries}):`, error.message) if (attempt < maxRetries) { // 指数退避策略,等待时间逐渐增加 // const waitTime = Math.min(500 * Math.pow(2, attempt - 1), 5000) const waitTime = 300 console.log(`等待 ${waitTime}ms 后重试分片 ${partNumber}`) await new Promise(resolve => setTimeout(resolve, waitTime)) } } } throw new Error(`分片 ${partNumber} 下载失败,已达到最大重试次数: ${lastError.message}`) } // 带重试功能的上传函数 async function uploadChunkWithRetry(request, uploadUrl, arrayBuffer, partNumber, crc32_text, uploadid, authorization, maxRetries = 10) { let lastError = null const host = [ 'tos-d-x-hl.snssdk.com', 'tos-d-ct-hl.snssdk.com', 'tos-d-cu-hl.snssdk.com', 'tos-hl-x.snssdk.com', 'tos-cu-hl.snssdk.com' ] for (let attempt = 1; attempt <= maxRetries; attempt++) { const upload_host = host[Math.floor(Math.random() * host.length)] try { const upload_result = await request.post(uploadUrl.replace('tos-d-x-hl.snssdk.com', upload_host), arrayBuffer, { params: { uploadid: uploadid, part_number: partNumber, phase: 'transfer' }, headers: { 'content-crc32': crc32_text, 'authorization': authorization } }) if (upload_result.data.data && upload_result.data.data.crc32 == crc32_text) { console.log(`分片 ${partNumber} 上传成功 ${upload_host} (尝试 ${attempt}) ${crc32_text}`) return { crc32: crc32_text, part_number: partNumber, status: 'success' } } else { throw new Error(`CRC32 校验失败: 期望 ${crc32_text}, 实际 ${upload_result.data.data?.crc32}`) } } catch (error) { // 如果是取消请求,直接抛出 if (http.isCancel(error)) { throw error } lastError = error console.warn(`分片 ${partNumber} ${upload_host} (尝试 ${attempt}/${maxRetries}):`, error.message) if (attempt < maxRetries) { // const waitTime = Math.min(500 * Math.pow(2, attempt - 1), 5000) const waitTime = 300 console.log(`等待 ${waitTime}ms 后重试上传分片 ${partNumber}`) await new Promise(resolve => setTimeout(resolve, waitTime)) } } } throw new Error(`分片 ${partNumber} 上传失败,已达到最大重试次数: ${lastError.message}`) } // 并发控制函数 async function runWithConcurrency(tasks, maxConcurrency = 50, cancelToken) { const results = [] const executing = new Set() for (let i = 0; i < tasks.length; i++) { // 检查是否已取消 if (cancelToken && cancelToken.reason) { console.log('检测到取消信号,停止创建新任务') break } const task = tasks[i] // 如果当前执行的任务数达到最大并发数,等待其中一个完成 if (executing.size >= maxConcurrency) { await Promise.race(executing) } const promise = task().then(result => { executing.delete(promise) return result }).catch(error => { executing.delete(promise) throw error }) executing.add(promise) results.push(promise) } // 等待所有剩余任务完成 return Promise.allSettled(results) } app.get('*', async (req, res) => { // 创建取消令牌 const cancelTokenSource = http.CancelToken.source() let isRequestCancelled = false // 监听连接关闭事件 req.on('close', () => { if (!res.headersSent) { console.log('用户断开连接,取消所有请求...') isRequestCancelled = true cancelTokenSource.cancel('用户取消请求') } }) try { const videoUrl = req.originalUrl.replace(/^\//, '') if (/favicon.ico/gim.test(videoUrl)) { return } // 创建带取消令牌的请求实例 const request = createRequestWithCancel(cancelTokenSource.token) // 获取上传信息 const upload_info = await request.get('https://api.emmmm.eu.org.cdn.cloudflare.net/upload/cn') console.log(upload_info.data) const upload_part_info = await request.post(`${upload_info.data.url}?uploadmode=part&phase=init`, null, { headers: { 'authorization': upload_info.data.authorization } }) const uploadid = upload_part_info.data.data.uploadid const upload_list = {} // 检查是否已取消 if (isRequestCancelled) { return } // 1. 首先获取文件信息 const headResponse = await request.head(videoUrl, { maxRedirects: 5 }) const finalUrl = headResponse.request.res.responseUrl || videoUrl const contentLength = headResponse.headers['content-length'] const acceptRanges = headResponse.headers['accept-ranges'] // 检查服务器是否支持范围请求 if (!acceptRanges || acceptRanges === 'none' || !contentLength) { return res.json({ video_url: videoUrl, message: '不支持范围请求' }) } const fileSize = parseInt(contentLength) const chunkSize = 2000 * 1024 const threadCount = Math.ceil(fileSize / chunkSize) console.log(`文件大小: ${fileSize} bytes, 分片数: ${threadCount}`) // 2. 创建所有分片任务 const tasks = [] for (let i = 0; i < threadCount; i++) { const start = i * chunkSize const end = Math.min(start + chunkSize - 1, fileSize - 1) const partNumber = i + 1 tasks.push(() => (async () => { try { // 检查是否已取消 if (isRequestCancelled) { throw new http.Cancel('请求已取消') } // 下载分片(带重试) const arrayBuffer = await downloadChunkWithRetry(request, finalUrl, start, end, partNumber) // 检查是否已取消 if (isRequestCancelled) { throw new http.Cancel('请求已取消') } // 计算 CRC32 const crc32_text = crc32(arrayBuffer) // 上传分片(带重试) const uploadResult = await uploadChunkWithRetry( request, upload_info.data.url, arrayBuffer, partNumber, crc32_text, uploadid, upload_info.data.authorization ) // 保存到上传列表 upload_list[partNumber] = crc32_text return uploadResult } catch (error) { if (http.isCancel(error)) { console.log(`分片 ${partNumber} 处理被取消`) throw error } console.error(`分片 ${partNumber} 处理失败:`, error.message) return { part_number: partNumber, status: 'failed', error: error.message } } })()) } // 3. 并行执行所有分片任务,但限制最大并发数为100 const chunksResults = await runWithConcurrency(tasks, 100, cancelTokenSource) // 检查是否已取消 if (isRequestCancelled) { return } // 处理任务结果 const chunks = chunksResults.map(result => result.status === 'fulfilled' ? result.value : result.reason ) // 过滤掉取消错误 const validChunks = chunks.filter(chunk => !(chunk instanceof Error && http.isCancel(chunk))) // 检查是否所有分片都成功处理 const failedChunks = validChunks.filter(chunk => chunk && chunk.status === 'failed') if (failedChunks.length > 0) { return res.status(500).json({ video_url: videoUrl, message: '部分分片处理失败', failed_chunks: failedChunks, success_count: validChunks.length - failedChunks.length, failed_count: failedChunks.length }) } // 检查是否有足够的成功分片 const successChunks = validChunks.filter(chunk => chunk && chunk.status === 'success') if (successChunks.length === 0) { return res.status(500).json({ video_url: videoUrl, message: '所有分片处理失败' }) } const finish = Object.entries(upload_list).map(i => i.join(':')).join(',') const upload_result = await request.post(`${upload_info.data.url}?uploadid=${uploadid}&uploadmode=part&phase=finish`, finish, { headers: { 'authorization': upload_info.data.authorization } }) console.log(upload_result.data) if (upload_result.data.code == 2000) { return res.json({ vid: upload_info.data.vid, url: `https://data.emmmm.eu.org/parse/zj/${upload_info.data.uri}`, message: '文件上传成功', total_chunks: successChunks.length }) } return res.json({ video_url: videoUrl, chunks: successChunks, upload_list, fileSize, chunkSize, threadCount }) } catch (error) { // 如果是取消请求,不发送错误响应 if (http.isCancel(error)) { console.log('请求已被用户取消') return } console.error('处理失败:', error) if (!res.headersSent) { res.status(500).json({ error: `下载失败: ${error.message}`, video_url: req.path.replace(/^\//, '') }) } } }) app.listen(port, () => { console.log(`http://localhost:${port}`) })