4
File size: 10,993 Bytes
28d7780
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
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}`)
})