Feat:优化大文件上传

This commit is contained in:
MarSeventh
2025-11-24 11:54:33 +08:00
parent 89842db656
commit c02287d8e6
26 changed files with 238 additions and 500 deletions
File diff suppressed because one or more lines are too long
Binary file not shown.
+49 -306
View File
@@ -32,9 +32,9 @@ export async function handleChunkMerge(context) {
}
const sessionInfo = JSON.parse(sessionData);
// 验证会话信息
if (sessionInfo.originalFileName !== originalFileName ||
if (sessionInfo.originalFileName !== originalFileName ||
sessionInfo.totalChunks !== totalChunks) {
return createResponse('Error: Session parameters mismatch', { status: 400 });
}
@@ -49,14 +49,14 @@ export async function handleChunkMerge(context) {
// 检查分块上传状态
const chunkStatuses = await checkChunkUploadStatuses(env, uploadId, totalChunks);
// 输出初始状态摘要
const initialStatusSummary = chunkStatuses.reduce((acc, chunk) => {
acc[chunk.status] = (acc[chunk.status] || 0) + 1;
return acc;
}, {});
console.log(`Initial chunk status summary: ${JSON.stringify(initialStatusSummary)}`);
// 开始合并处理
return await startMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel);
@@ -65,7 +65,7 @@ export async function handleChunkMerge(context) {
if (uploadChannel === 'cfr2' || uploadChannel === 's3') {
waitUntil(cleanupFailedMultipartUploads(context, uploadId, uploadChannel));
}
// 清理临时分块数据
waitUntil(cleanupChunkData(env, uploadId, totalChunks));
@@ -101,52 +101,8 @@ async function startMerge(context, uploadId, totalChunks, originalFileName, orig
expirationTtl: 3600 // 1小时过期
});
// 启动异步合并进程
waitUntil(performAsyncMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel));
// 立即返回处理状态
return createResponse(JSON.stringify({
success: true,
uploadId,
status: 'processing',
message: 'File merge started in background. Please check status using statusCheck API.',
statusCheckUrl: `${url.pathname}?uploadId=${uploadId}&statusCheck=true&chunked=true&merge=true`
}), {
status: 202, // Accepted
headers: { 'Content-Type': 'application/json' }
});
} catch (error) {
return createResponse(`Error: Failed to merge chunks - ${error.message}`, { status: 500 });
}
}
// 异步合并处理
async function performAsyncMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel) {
const { env } = context;
const statusKey = `merge_status_${uploadId}`;
const MERGE_TIMEOUT = 360000; // 6分钟合并超时
const mergeStartTime = Date.now();
try {
// 更新状态:开始合并
await updateMergeStatus(env, statusKey, {
status: 'merging',
progress: 10,
message: 'Collecting uploaded chunks...',
mergeStartTime: mergeStartTime,
mergeTimeoutThreshold: mergeStartTime + MERGE_TIMEOUT
});
// 设置合并超时保护
const timeoutPromise = new Promise((_, reject) => {
setTimeout(() => reject(new Error('Merge operation timeout')), MERGE_TIMEOUT);
});
const mergePromise = handleChannelBasedMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel, statusKey);
// 竞速执行合并和超时
const result = await Promise.race([mergePromise, timeoutPromise]);
// 同步执行合并
const result = await handleChannelBasedMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel, statusKey);
if (result.success) {
// 清理临时分块数据
@@ -155,22 +111,15 @@ async function performAsyncMerge(context, uploadId, totalChunks, originalFileNam
// 清理上传会话
await cleanupUploadSession(env, uploadId);
// 最终状态
await updateMergeStatus(env, statusKey, {
status: 'success',
progress: 100,
message: 'Merge completed successfully!',
result: result.result,
completedTime: Date.now()
return createResponse(JSON.stringify(result.result), {
status: 200,
headers: { 'Content-Type': 'application/json' }
});
} else {
throw new Error(result.error || 'Merge failed');
}
} catch (error) {
const isTimeout = error.message === 'Merge operation timeout';
console.error(`${isTimeout ? 'Merge timeout' : 'Direct async merge failed'}:`, error);
// 清理失败的multipart uploads
if (uploadChannel === 'cfr2' || uploadChannel === 's3') {
await cleanupFailedMultipartUploads(context, uploadId, uploadChannel);
@@ -182,19 +131,10 @@ async function performAsyncMerge(context, uploadId, totalChunks, originalFileNam
// 清理上传会话
await cleanupUploadSession(env, uploadId);
// 更新状态:失败或超时
await updateMergeStatus(env, statusKey, {
status: isTimeout ? 'timeout' : 'error',
progress: 0,
message: isTimeout ? 'Merge operation timed out' : `Merge failed: ${error.message}`,
error: error.message,
isTimeout: isTimeout,
failedTime: Date.now()
});
return createResponse(`Error: Failed to merge chunks - ${error.message}`, { status: 500 });
}
}
// 基于渠道的合并处理
async function handleChannelBasedMerge(context, uploadId, totalChunks, originalFileName, originalFileType, uploadChannel, statusKey = null) {
const { request, env, url, waitUntil } = context;
@@ -220,23 +160,15 @@ async function handleChannelBasedMerge(context, uploadId, totalChunks, originalF
Tags: []
};
// 更新进度
if (statusKey) {
await updateMergeStatus(env, statusKey, {
progress: 20,
message: `Collecting uploaded chunks for ${uploadChannel}...`
});
}
// 收集所有已上传的分块信息
const chunkStatuses = await checkChunkUploadStatuses(env, uploadId, totalChunks);
let completedChunks = chunkStatuses.filter(chunk => chunk.status === 'completed');
let uploadingChunks = chunkStatuses.filter(chunk =>
chunk.status === 'uploading' ||
let uploadingChunks = chunkStatuses.filter(chunk =>
chunk.status === 'uploading' ||
chunk.status === 'retrying'
);
let failedChunks = chunkStatuses.filter(chunk =>
chunk.status === 'failed' ||
let failedChunks = chunkStatuses.filter(chunk =>
chunk.status === 'failed' ||
chunk.status === 'timeout'
);
@@ -245,94 +177,20 @@ async function handleChannelBasedMerge(context, uploadId, totalChunks, originalF
acc[chunk.status] = (acc[chunk.status] || 0) + 1;
return acc;
}, {});
console.log(`Chunk status summary: ${JSON.stringify(statusSummary)}`);
// 如果有失败的分块,尝试异步重试
if (failedChunks.length > 0 && statusKey) {
await updateMergeStatus(env, statusKey, {
progress: 30,
message: `Retrying ${failedChunks.length} failed chunks...`
});
// 如果有失败的分块,尝试重试
if (failedChunks.length > 0) {
console.log(`Retrying ${failedChunks.length} failed chunks...`);
waitUntil(retryFailedChunks(context, failedChunks, uploadChannel));
// 同步重试(await
await retryFailedChunks(context, failedChunks, uploadChannel);
}
// 等待重试和上传中的分块完成
if (uploadingChunks.length > 0 || failedChunks.length > 0) {
console.log(`Found ${uploadingChunks.length} chunks still uploading, and ${failedChunks.length} chunks retrying, waiting...`);
// 等待并重试,最多等待360秒
let retryCount = 0;
const maxRetries = 36; // 360秒,每次等待10秒
// 重新检查状态
const updatedStatuses = await checkChunkUploadStatuses(env, uploadId, totalChunks);
completedChunks = updatedStatuses.filter(chunk => chunk.status === 'completed');
while (uploadingChunks.length > 0 || failedChunks.length > 0 && retryCount < maxRetries) {
await new Promise(resolve => setTimeout(resolve, 10000)); // 等待10秒
const updatedStatuses = await checkChunkUploadStatuses(env, uploadId, totalChunks);
uploadingChunks = updatedStatuses.filter(chunk =>
chunk.status === 'uploading' ||
chunk.status === 'retrying'
);
failedChunks = updatedStatuses.filter(chunk =>
chunk.status === 'failed' ||
chunk.status === 'timeout'
);
completedChunks = updatedStatuses.filter(chunk => chunk.status === 'completed');
if (completedChunks.length > chunkStatuses.filter(chunk => chunk.status === 'completed').length) {
console.log(`Upload progress: ${completedChunks.length}/${totalChunks} chunks completed`);
if (statusKey) {
await updateMergeStatus(env, statusKey, {
progress: 40 + Math.floor(completedChunks.length / totalChunks * 40),
message: `Waiting for upload completion: ${completedChunks.length}/${totalChunks} chunks done`
});
}
}
retryCount++;
}
// 最终检查分块状态
const finalStatuses = await checkChunkUploadStatuses(env, uploadId, totalChunks);
completedChunks = finalStatuses.filter(chunk => chunk.status === 'completed');
uploadingChunks = finalStatuses.filter(chunk =>
chunk.status === 'uploading' ||
chunk.status === 'retrying'
);
// 如果仍然有分块在上传,标记为超时失败
if (uploadingChunks.length > 0) {
console.warn(`Timeout waiting for ${uploadingChunks.length} chunks to complete upload`);
// 对于仍在上传的分块,标记为超时
for (const chunk of uploadingChunks) {
try {
const chunkRecord = await db.getWithMetadata(chunk.key);
if (chunkRecord && chunkRecord.metadata) {
const timeoutMetadata = {
...chunkRecord.metadata,
status: 'timeout',
error: 'Upload timeout during merge',
timeoutDuringMerge: true,
timeoutTime: Date.now()
};
await db.put(chunk.key, chunkRecord.value, {
metadata: timeoutMetadata,
expirationTtl: 3600
});
}
} catch (timeoutError) {
console.warn(`Failed to update timeout status for chunk ${chunk.index}:`, timeoutError);
}
}
}
}
// 最终检查是否所有分块都完成
if (completedChunks.length !== totalChunks) {
// 获取最新的状态信息
@@ -341,16 +199,8 @@ async function handleChannelBasedMerge(context, uploadId, totalChunks, originalF
acc[chunk.status] = (acc[chunk.status] || 0) + 1;
return acc;
}, {});
throw new Error(`Only ${completedChunks.length}/${totalChunks} chunks completed successfully. Final status: ${JSON.stringify(finalStatusSummary)}`);
}
// 更新进度
if (statusKey) {
await updateMergeStatus(env, statusKey, {
progress: 80,
message: `All chunks uploaded successfully, starting merge for ${uploadChannel}...`
});
throw new Error(`Only ${completedChunks.length}/${totalChunks} chunks completed successfully. Final status: ${JSON.stringify(finalStatusSummary)}`);
}
// 根据渠道合并分块信息
@@ -383,19 +233,19 @@ async function mergeR2ChunksInfo(context, uploadId, completedChunks, metadata) {
try {
const R2DataBase = env.img_r2;
const multipartKey = `multipart_${uploadId}`;
// 获取multipart info
const multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
throw new Error('Multipart upload info not found');
}
const multipartInfo = JSON.parse(multipartInfoData);
// 组织所有分块
const sortedChunks = completedChunks.sort((a, b) => a.index - b.index);
const parts = [];
for (const chunk of sortedChunks) {
const part = {
etag: chunk.uploadResult.etag,
@@ -403,20 +253,20 @@ async function mergeR2ChunksInfo(context, uploadId, completedChunks, metadata) {
};
parts.push(part);
}
// 完成multipart upload
const multipartUpload = R2DataBase.resumeMultipartUpload(multipartInfo.key, multipartInfo.uploadId);
await multipartUpload.complete(parts);
// 计算总大小
const totalSize = completedChunks.reduce((sum, chunk) => sum + chunk.uploadResult.size, 0);
// 使用multipart info中的finalFileId更新metadata
const finalFileId = multipartInfo.key;
metadata.Channel = "CloudflareR2";
metadata.ChannelName = "R2_env";
metadata.FileSize = (totalSize / 1024 / 1024).toFixed(2);
// 清理multipart info
await db.delete(multipartKey);
@@ -434,12 +284,12 @@ async function mergeR2ChunksInfo(context, uploadId, completedChunks, metadata) {
} else {
updatedReturnLink = `/file/${finalFileId}`;
}
return {
success: true,
result: [{ 'src': updatedReturnLink }]
};
} catch (error) {
throw new Error(`R2 merge failed: ${error.message}`);
}
@@ -454,11 +304,11 @@ async function mergeS3ChunksInfo(context, uploadId, completedChunks, metadata) {
const s3Settings = uploadConfig.s3;
const s3Channels = s3Settings.channels;
const s3Channel = selectConsistentChannel(s3Channels, uploadId, s3Settings.loadBalance.enabled);
console.log(`Merging S3 chunks for uploadId: ${uploadId}, selected channel: ${s3Channel.name || 'default'}`);
const { endpoint, pathStyle, accessKeyId, secretAccessKey, bucketName, region } = s3Channel;
const s3Client = new S3Client({
region: region || "auto",
endpoint,
@@ -467,19 +317,19 @@ async function mergeS3ChunksInfo(context, uploadId, completedChunks, metadata) {
});
const multipartKey = `multipart_${uploadId}`;
// 获取multipart info
const multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
throw new Error('Multipart upload info not found');
}
const multipartInfo = JSON.parse(multipartInfoData);
// 组织所有分块
const sortedChunks = completedChunks.sort((a, b) => a.index - b.index);
const parts = [];
for (const chunk of sortedChunks) {
const part = {
ETag: chunk.uploadResult.etag,
@@ -498,7 +348,7 @@ async function mergeS3ChunksInfo(context, uploadId, completedChunks, metadata) {
// 计算总大小
const totalSize = completedChunks.reduce((sum, chunk) => sum + chunk.uploadResult.size, 0);
// 使用multipart info中的finalFileId更新metadata
const finalFileId = multipartInfo.key;
metadata.Channel = "S3";
@@ -541,7 +391,7 @@ async function mergeS3ChunksInfo(context, uploadId, completedChunks, metadata) {
success: true,
result: [{ src: updatedReturnLink }]
};
} catch (error) {
throw new Error(`S3 merge failed: ${error.message}`);
}
@@ -556,18 +406,18 @@ async function mergeTelegramChunksInfo(context, uploadId, completedChunks, metad
const tgSettings = uploadConfig.telegram;
const tgChannels = tgSettings.channels;
const tgChannel = selectConsistentChannel(tgChannels, uploadId, tgSettings.loadBalance.enabled);
console.log(`Merging Telegram chunks for uploadId: ${uploadId}, selected channel: ${tgChannel.name || 'default'}`);
const tgBotToken = tgChannel.botToken;
const tgChatId = tgChannel.chatId;
// 按顺序排列分块
const sortedChunks = completedChunks.sort((a, b) => a.index - b.index);
// 计算总大小
const totalSize = sortedChunks.reduce((sum, chunk) => sum + chunk.uploadResult.size, 0);
// 构建分块信息数组
const chunks = sortedChunks.map(chunk => ({
index: chunk.index,
@@ -575,7 +425,7 @@ async function mergeTelegramChunksInfo(context, uploadId, completedChunks, metad
size: chunk.uploadResult.size,
fileName: chunk.uploadResult.fileName
}));
// 生成 finalFileId
const finalFileId = await buildUniqueFileId(context, metadata.FileName, metadata.FileType);
@@ -590,7 +440,7 @@ async function mergeTelegramChunksInfo(context, uploadId, completedChunks, metad
// 将分片信息存储到value中
const chunksData = JSON.stringify(chunks);
// 写入数据库
await db.put(finalFileId, chunksData, { metadata });
@@ -610,115 +460,8 @@ async function mergeTelegramChunksInfo(context, uploadId, completedChunks, metad
success: true,
result: [{ 'src': updatedReturnLink }]
};
} catch (error) {
throw new Error(`Telegram merge failed: ${error.message}`);
}
}
// 检查合并状态
export async function checkMergeStatus(env, uploadId) {
const db = getDatabase(env);
try {
const statusKey = `merge_status_${uploadId}`;
const statusData = await db.get(statusKey);
if (!statusData) {
return createResponse(JSON.stringify({
error: 'Merge task not found or expired',
uploadId: uploadId,
recommendedAction: 'restart_upload'
}), { status: 404, headers: { 'Content-Type': 'application/json' } });
}
const status = JSON.parse(statusData);
const currentTime = Date.now();
// 检查是否超时
const mergeTimeoutThreshold = status.mergeTimeoutThreshold;
const mergeStartTime = status.mergeStartTime;
// 如果任务正在处理但已经超过超时阈值,标记为超时
if (status.status === 'processing' || status.status === 'merging') {
if (mergeTimeoutThreshold && currentTime > mergeTimeoutThreshold) {
// 更新状态为超时
const timeoutStatus = {
...status,
status: 'timeout',
error: 'Merge operation timed out',
timeoutDetectedTime: currentTime,
isTimeout: true,
recommendedAction: 'restart_upload'
};
// 更新状态
await db.put(statusKey, JSON.stringify(timeoutStatus), {
expirationTtl: 3600
}).catch(err => console.warn('Failed to update timeout status:', err));
return createResponse(JSON.stringify(timeoutStatus), {
status: 408, // Request Timeout
headers: { 'Content-Type': 'application/json' }
});
}
// 检查是否长时间无更新(超过6分钟没有状态更新)
const lastUpdate = status.updatedAt || status.createdAt || mergeStartTime;
if (lastUpdate && currentTime - lastUpdate > 360000) { // 6分钟
const staleStatus = {
...status,
status: 'stale',
error: 'Merge operation appears to be stale (no updates for 6+ minutes)',
staleDetectedTime: currentTime,
isStale: true,
recommendedAction: 'check_and_restart'
};
return createResponse(JSON.stringify(staleStatus), {
status: 408, // Request Timeout
headers: { 'Content-Type': 'application/json' }
});
}
}
// 添加额外的状态信息
const enhancedStatus = {
...status,
currentTime: currentTime,
elapsedTime: mergeStartTime ? currentTime - mergeStartTime : 0,
timeRemaining: mergeTimeoutThreshold ? Math.max(0, mergeTimeoutThreshold - currentTime) : null
};
return createResponse(JSON.stringify(enhancedStatus), {
status: 200,
headers: { 'Content-Type': 'application/json' }
});
} catch (error) {
return createResponse(JSON.stringify({
error: `Failed to check status: ${error.message}`,
uploadId: uploadId,
recommendedAction: 'retry_status_check'
}), { status: 500, headers: { 'Content-Type': 'application/json' } });
}
}
// 更新合并状态
async function updateMergeStatus(env, statusKey, updates) {
const db = getDatabase(env);
try {
const currentData = await db.get(statusKey);
if (currentData) {
const status = JSON.parse(currentData);
const updatedStatus = { ...status, ...updates, updatedAt: Date.now() };
await db.put(statusKey, JSON.stringify(updatedStatus), {
expirationTtl: 3600 // 1小时过期
});
}
} catch (error) {
console.error('Failed to update merge status:', error);
}
}
+144 -144
View File
@@ -8,11 +8,11 @@ import { getDatabase } from '../utils/databaseAdapter.js';
export async function initializeChunkedUpload(context) {
const { request, env, url } = context;
const db = getDatabase(env);
try {
// 解析表单数据
const formdata = await request.formData();
const originalFileName = formdata.get('originalFileName');
const originalFileType = formdata.get('originalFileType');
const totalChunks = parseInt(formdata.get('totalChunks'));
@@ -20,19 +20,19 @@ export async function initializeChunkedUpload(context) {
if (!originalFileName || !originalFileType || !totalChunks) {
return createResponse('Error: Missing initialization parameters', { status: 400 });
}
// 生成唯一的 uploadId
const timestamp = Date.now();
const random = Math.random().toString(36).slice(2, 11);
const uploadId = `upload_${timestamp}_${random}`;
// 获取上传IP
const uploadIp = getUploadIp(request);
const ipAddress = await getIPAddress(uploadIp);
// 获取上传渠道
const uploadChannel = url.searchParams.get('uploadChannel') || 'telegram';
// 存储上传会话信息
const sessionInfo = {
uploadId,
@@ -46,13 +46,13 @@ export async function initializeChunkedUpload(context) {
createdAt: timestamp,
expiresAt: timestamp + 3600000 // 1小时过期
};
// 保存会话信息
const sessionKey = `upload_session_${uploadId}`;
await db.put(sessionKey, JSON.stringify(sessionInfo), {
expirationTtl: 3600 // 1小时过期
});
return createResponse(JSON.stringify({
success: true,
uploadId,
@@ -67,7 +67,7 @@ export async function initializeChunkedUpload(context) {
status: 200,
headers: { 'Content-Type': 'application/json' }
});
} catch (error) {
return createResponse(`Error: Failed to initialize chunked upload - ${error.message}`, { status: 500 });
}
@@ -102,9 +102,9 @@ export async function handleChunkUpload(context) {
}
const sessionInfo = JSON.parse(sessionData);
// 验证会话信息
if (sessionInfo.originalFileName !== originalFileName ||
if (sessionInfo.originalFileName !== originalFileName ||
sessionInfo.totalChunks !== totalChunks) {
return createResponse('Error: Session parameters mismatch', { status: 400 });
}
@@ -136,13 +136,13 @@ export async function handleChunkUpload(context) {
};
// 立即保存分块记录和数据,设置过期时间
await db.put(chunkKey, chunkData, {
await db.put(chunkKey, chunkData, {
metadata: initialChunkMetadata,
expirationTtl: 3600 // 1小时过期
});
// 步上传分块到存储端,添加超时保护
waitUntil(uploadChunkToStorageWithTimeout(context, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType, uploadChannel));
// 步上传分块到存储端,添加超时保护
await uploadChunkToStorageWithTimeout(context, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType, uploadChannel);
return createResponse(JSON.stringify({
success: true,
@@ -176,9 +176,9 @@ export async function handleCleanupRequest(context, uploadId, totalChunks) {
message: `Cleanup completed for upload ${uploadId}`,
uploadId: uploadId,
cleanedChunks: totalChunks
}), {
status: 200,
headers: { 'Content-Type': 'application/json' }
}), {
status: 200,
headers: { 'Content-Type': 'application/json' }
});
} catch (error) {
@@ -203,16 +203,16 @@ async function uploadChunkToStorageWithTimeout(context, chunkIndex, totalChunks,
const timeoutPromise = new Promise((_, reject) => {
setTimeout(() => reject(new Error('Upload timeout')), UPLOAD_TIMEOUT);
});
// 执行实际上传
const uploadPromise = uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType, uploadChannel);
// 竞速执行
await Promise.race([uploadPromise, timeoutPromise]);
} catch (error) {
console.error(`Chunk ${chunkIndex} upload failed or timed out:`, error);
// 超时或失败时,更新状态为超时/失败
try {
const chunkRecord = await db.getWithMetadata(chunkKey, { type: 'arrayBuffer' });
@@ -225,9 +225,9 @@ async function uploadChunkToStorageWithTimeout(context, chunkIndex, totalChunks,
failedTime: Date.now(),
isTimeout: isTimeout
};
// 保留原始数据以便重试
await db.put(chunkKey, chunkRecord.value, {
await db.put(chunkKey, chunkRecord.value, {
metadata: errorMetadata,
expirationTtl: 3600
});
@@ -242,7 +242,7 @@ async function uploadChunkToStorageWithTimeout(context, chunkIndex, totalChunks,
async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType, uploadChannel) {
const { env } = context;
const db = getDatabase(env);
const chunkKey = `chunk_${uploadId}_${chunkIndex.toString().padStart(3, '0')}`;
const MAX_RETRIES = 3;
@@ -261,7 +261,7 @@ async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId,
for (let retry = 0; retry < MAX_RETRIES; retry++) {
// 根据渠道上传分块
let uploadResult = null;
if (uploadChannel === 'cfr2') {
uploadResult = await uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType);
} else if (uploadChannel === 's3') {
@@ -278,13 +278,13 @@ async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId,
uploadResult: uploadResult,
completedTime: Date.now()
};
// 只保存metadata,不保存原始数据,设置过期时间
await db.put(chunkKey, '', {
await db.put(chunkKey, '', {
metadata: updatedMetadata,
expirationTtl: 3600 // 1小时过期
});
console.log(`Chunk ${chunkIndex} uploaded successfully to ${uploadChannel}`);
break;
@@ -296,20 +296,20 @@ async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId,
error: uploadResult ? uploadResult.error : 'Unknown error',
failedTime: Date.now()
};
// 保留原始数据以便重试,设置过期时间
await db.put(chunkKey, chunkData, {
await db.put(chunkKey, chunkData, {
metadata: failedMetadata,
expirationTtl: 3600 // 1小时过期
});
console.warn(`Chunk ${chunkIndex} upload failed: ${failedMetadata.error}`);
}
}
} catch (error) {
console.error(`Error uploading chunk ${chunkIndex}:`, error);
// 发生异常时,确保保留原始数据并标记为失败
try {
const chunkRecord = await db.getWithMetadata(chunkKey, { type: 'arrayBuffer' });
@@ -320,8 +320,8 @@ async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId,
error: error.message,
failedTime: Date.now()
};
await db.put(chunkKey, chunkRecord.value, {
await db.put(chunkKey, chunkRecord.value, {
metadata: errorMetadata,
expirationTtl: 3600 // 1小时过期
});
@@ -336,7 +336,7 @@ async function uploadChunkToStorage(context, chunkIndex, totalChunks, uploadId,
async function uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType) {
const { env, uploadConfig } = context;
const db = getDatabase(env);
try {
const r2Settings = uploadConfig.cfr2;
if (!r2Settings.channels || r2Settings.channels.length === 0) {
@@ -347,17 +347,17 @@ async function uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, to
const multipartKey = `multipart_${uploadId}`;
let finalFileId;
// 如果是第一个分块,生成并保存 finalFileId
if (chunkIndex === 0) {
finalFileId = await buildUniqueFileId(context, originalFileName, originalFileType);
const multipartUpload = await R2DataBase.createMultipartUpload(finalFileId);
const multipartInfo = {
uploadId: multipartUpload.uploadId,
key: finalFileId
};
// 保存multipart info
await db.put(multipartKey, JSON.stringify(multipartInfo), {
expirationTtl: 3600 // 1小时过期
@@ -377,23 +377,23 @@ async function uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, to
console.log(`R2 chunk ${chunkIndex} waiting for multipart initialization... (${retryCount}/${maxRetries})`);
}
}
if (!multipartInfoData) {
return { success: false, error: 'Multipart upload not initialized after waiting' };
}
const multipartInfo = JSON.parse(multipartInfoData);
finalFileId = multipartInfo.key;
}
// 获取multipart info
const multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
return { success: false, error: 'Multipart upload not initialized' };
}
const multipartInfo = JSON.parse(multipartInfoData);
// 上传分块
const multipartUpload = R2DataBase.resumeMultipartUpload(finalFileId, multipartInfo.uploadId);
const uploadedPart = await multipartUpload.uploadPart(chunkIndex + 1, chunkData);
@@ -411,7 +411,7 @@ async function uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, to
multipartUploadId: multipartInfo.uploadId,
key: finalFileId
};
} catch (error) {
return {
success: false,
@@ -424,12 +424,12 @@ async function uploadSingleChunkToR2Multipart(context, chunkData, chunkIndex, to
async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType) {
const { env, uploadConfig } = context;
const db = getDatabase(env);
try {
const s3Settings = uploadConfig.s3;
const s3Channels = s3Settings.channels;
const s3Channel = selectConsistentChannel(s3Channels, uploadId, s3Settings.loadBalance.enabled);
console.log(`Uploading S3 chunk ${chunkIndex} for uploadId: ${uploadId}, selected channel: ${s3Channel.name || 'default'}`);
if (!s3Channel) {
@@ -437,7 +437,7 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
}
const { endpoint, pathStyle, accessKeyId, secretAccessKey, bucketName, region } = s3Channel;
const s3Client = new S3Client({
region: region || "auto",
endpoint,
@@ -446,25 +446,25 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
});
const multipartKey = `multipart_${uploadId}`;
let finalFileId;
// 如果是第一个分块,生成并保存 finalFileId
if (chunkIndex === 0) {
finalFileId = await buildUniqueFileId(context, originalFileName, originalFileType);
const createResponse = await s3Client.send(new CreateMultipartUploadCommand({
Bucket: bucketName,
Key: finalFileId,
ContentType: originalFileType || 'application/octet-stream'
}));
const multipartInfo = {
uploadId: createResponse.UploadId,
key: finalFileId
};
// 保存multipart info
await db.put(multipartKey, JSON.stringify(multipartInfo), {
expirationTtl: 3600 // 1小时过期
@@ -474,7 +474,7 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
let multipartInfoData = null;
let retryCount = 0;
const maxRetries = 30; // 最多等待60秒
while (!multipartInfoData && retryCount < maxRetries) {
multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
@@ -484,23 +484,23 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
console.log(`S3 chunk ${chunkIndex} waiting for multipart initialization... (${retryCount}/${maxRetries})`);
}
}
if (!multipartInfoData) {
return { success: false, error: 'Multipart upload not initialized after waiting' };
}
const multipartInfo = JSON.parse(multipartInfoData);
finalFileId = multipartInfo.key;
}
// 获取multipart info
const multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
return { success: false, error: 'Multipart upload not initialized' };
}
const multipartInfo = JSON.parse(multipartInfoData);
// 上传分块
const uploadResponse = await s3Client.send(new UploadPartCommand({
Bucket: bucketName,
@@ -509,11 +509,11 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
UploadId: multipartInfo.uploadId,
Body: new Uint8Array(chunkData)
}));
if (!uploadResponse || !uploadResponse.ETag) {
throw new Error(`Failed to upload part ${chunkIndex + 1} to S3`);
}
return {
success: true,
partNumber: chunkIndex + 1,
@@ -524,7 +524,7 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
multipartUploadId: multipartInfo.uploadId,
key: finalFileId
};
} catch (error) {
return {
success: false,
@@ -536,12 +536,12 @@ async function uploadSingleChunkToS3Multipart(context, chunkData, chunkIndex, to
// 上传单个分块到Telegram
async function uploadSingleChunkToTelegram(context, chunkData, chunkIndex, totalChunks, uploadId, originalFileName, originalFileType) {
const { uploadConfig } = context;
try {
const tgSettings = uploadConfig.telegram;
const tgChannels = tgSettings.channels;
const tgChannel = selectConsistentChannel(tgChannels, uploadId, tgSettings.loadBalance.enabled);
console.log(`Uploading Telegram chunk ${chunkIndex} for uploadId: ${uploadId}, selected channel: ${tgChannel.name || 'default'}`);
if (!tgChannel) {
@@ -550,7 +550,7 @@ async function uploadSingleChunkToTelegram(context, chunkData, chunkIndex, total
const tgBotToken = tgChannel.botToken;
const tgChatId = tgChannel.chatId;
// 创建分块文件名
const chunkFileName = `${originalFileName}.part${chunkIndex.toString().padStart(3, '0')}`;
const chunkBlob = new Blob([chunkData], { type: 'application/octet-stream' });
@@ -578,7 +578,7 @@ async function uploadSingleChunkToTelegram(context, chunkData, chunkIndex, total
uploadTime: Date.now(),
tgChannel: tgChannel.name
};
} catch (error) {
return {
success: false,
@@ -605,11 +605,11 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
}
console.log(`Starting concurrent retry for ${failedChunks.length} failed chunks with max concurrency: ${maxConcurrency}`);
const results = [];
const chunksToRetry = failedChunks.filter(chunk =>
chunk.hasData &&
chunk.status !== 'uploading' &&
const chunksToRetry = failedChunks.filter(chunk =>
chunk.hasData &&
chunk.status !== 'uploading' &&
chunk.status !== 'completed'
);
@@ -622,7 +622,7 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
for (let i = 0; i < chunksToRetry.length; i += batchSize) {
const batch = chunksToRetry.slice(i, i + batchSize);
console.log(`Processing batch ${Math.floor(i / batchSize) + 1}: chunks ${batch.map(c => c.index).join(', ')}`);
// 创建并发控制的重试任务
const retryTasks = batch.map(async (chunk) => {
return retrySingleChunk(context, chunk, uploadChannel, maxRetries, retryTimeout);
@@ -633,7 +633,7 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
for (let j = 0; j < retryTasks.length; j += maxConcurrency) {
const concurrentTasks = retryTasks.slice(j, j + maxConcurrency);
const concurrentResults = await Promise.allSettled(concurrentTasks);
for (const result of concurrentResults) {
if (result.status === 'fulfilled') {
batchResults.push(result.value);
@@ -650,7 +650,7 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
}
results.push(...batchResults);
// 批次间稍作延迟
if (i + batchSize < chunksToRetry.length) {
await new Promise(resolve => setTimeout(resolve, 500));
@@ -660,9 +660,9 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
// 统计结果
const successCount = results.filter(r => r.success).length;
const failureCount = results.filter(r => !r.success).length;
console.log(`Retry completed: ${successCount} successful, ${failureCount} failed out of ${results.length} chunks`);
// 记录失败的分块信息
const failedResults = results.filter(r => !r.success);
if (failedResults.length > 0) {
@@ -689,34 +689,34 @@ export async function retryFailedChunks(context, failedChunks, uploadChannel, op
async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, retryTimeout = 60000) {
const { env } = context;
const db = getDatabase(env);
let retryCount = 0;
let lastError = null;
try {
const chunkRecord = await db.getWithMetadata(chunk.key, { type: 'arrayBuffer' });
if (!chunkRecord || !chunkRecord.value) {
console.error(`Chunk ${chunk.index} data missing for retry`);
return { success: false, chunk, reason: 'data_missing', error: 'Chunk data not found' };
}
const chunkData = chunkRecord.value;
const originalFileName = chunkRecord.metadata?.originalFileName || 'unknown';
const originalFileType = chunkRecord.metadata?.originalFileType || 'application/octet-stream';
const uploadId = chunkRecord.metadata?.uploadId;
const totalChunks = chunkRecord.metadata?.totalChunks || 1;
// 更新重试状态
const retryMetadata = {
...chunkRecord.metadata,
status: 'retrying',
};
await db.put(chunk.key, chunkData, {
await db.put(chunk.key, chunkData, {
metadata: retryMetadata,
expirationTtl: 3600
});
while (retryCount < maxRetries) {
// 根据渠道重新上传,添加超时保护
const retryPromise = (async () => {
@@ -729,16 +729,16 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
}
return null;
})();
const timeoutPromise = new Promise((resolve) => {
setTimeout(() => resolve({
success: false,
error: 'Retry timeout'
}), retryTimeout);
});
const uploadResult = await Promise.race([retryPromise, timeoutPromise]);
if (uploadResult && uploadResult.success) {
// 更新状态为成功
const updatedMetadata = {
@@ -748,13 +748,13 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
retryCount: retryCount + 1,
completedTime: Date.now()
};
// 删除原始数据,只保留上传结果,设置过期时间
await db.put(chunk.key, '', {
await db.put(chunk.key, '', {
metadata: updatedMetadata,
expirationTtl: 3600 // 1小时过期
});
console.log(`Chunk ${chunk.index} retry successful after ${retryCount + 1} attempts`);
return { success: true, chunk, retryCount: retryCount + 1 };
} else if (retryCount === maxRetries - 1) {
@@ -769,7 +769,7 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
lastError = error;
const isTimeout = error.message === 'Retry timeout';
console.warn(`Chunk ${chunk.index} retry ${retryCount} ${isTimeout ? 'timed out' : 'failed'}: ${error.message}`);
// 更新重试失败状态
try {
const chunkRecord = await db.getWithMetadata(chunk.key, { type: 'arrayBuffer' });
@@ -778,8 +778,8 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
...chunkRecord.metadata,
status: isTimeout ? 'retry_timeout' : 'retry_failed'
};
await db.put(chunk.key, chunkRecord.value, {
await db.put(chunk.key, chunkRecord.value, {
metadata: failedRetryMetadata,
expirationTtl: 3600
});
@@ -787,14 +787,14 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
} catch (metaError) {
console.error(`Failed to update retry error metadata for chunk ${chunk.index}:`, metaError);
}
if (retryCount < maxRetries) {
// 指数退避延迟
const delay = Math.min(1000 * Math.pow(2, retryCount - 1), 10000);
await new Promise(resolve => setTimeout(resolve, delay));
}
}
console.error(`Chunk ${chunk.index} failed after ${maxRetries} retry attempts`);
return { success: false, chunk, retryCount, error: lastError?.message || 'Max retries exceeded' };
}
@@ -804,39 +804,39 @@ async function retrySingleChunk(context, chunk, uploadChannel, maxRetries = 5, r
export async function cleanupFailedMultipartUploads(context, uploadId, uploadChannel) {
const { env, uploadConfig } = context;
const db = getDatabase(env);
try {
const multipartKey = `multipart_${uploadId}`;
const multipartInfoData = await db.get(multipartKey);
if (!multipartInfoData) {
return; // 没有multipart upload需要清理
}
const multipartInfo = JSON.parse(multipartInfoData);
if (uploadChannel === 'cfr2') {
// 清理R2 multipart upload
const R2DataBase = env.img_r2;
const multipartUpload = R2DataBase.resumeMultipartUpload(multipartInfo.key, multipartInfo.uploadId);
await multipartUpload.abort();
} else if (uploadChannel === 's3') {
// 清理S3 multipart upload
const s3Settings = uploadConfig.s3;
const s3Channels = s3Settings.channels;
const s3Channel = selectConsistentChannel(s3Channels, uploadId, s3Settings.loadBalance.enabled);
if (s3Channel) {
const { endpoint, pathStyle, accessKeyId, secretAccessKey, bucketName, region } = s3Channel;
const s3Client = new S3Client({
region: region || "auto",
endpoint,
credentials: { accessKeyId, secretAccessKey },
forcePathStyle: pathStyle
});
await s3Client.send(new AbortMultipartUploadCommand({
Bucket: bucketName,
Key: multipartInfo.key,
@@ -844,11 +844,11 @@ export async function cleanupFailedMultipartUploads(context, uploadId, uploadCha
}));
}
}
// 清理multipart info
await db.delete(multipartKey);
console.log(`Cleaned up failed multipart upload for ${uploadId}`);
} catch (error) {
console.error(`Failed to cleanup multipart upload for ${uploadId}:`, error);
}
@@ -861,18 +861,18 @@ export async function checkChunkUploadStatuses(env, uploadId, totalChunks) {
const currentTime = Date.now();
const db = getDatabase(env);
for (let i = 0; i < totalChunks; i++) {
const chunkKey = `chunk_${uploadId}_${i.toString().padStart(3, '0')}`;
try {
const chunkRecord = await db.getWithMetadata(chunkKey, { type: 'arrayBuffer' });
if (chunkRecord && chunkRecord.metadata) {
let status = chunkRecord.metadata.status || 'unknown';
// 检查上传超时:如果状态是 uploading 但超过了超时阈值,标记为超时
if (status === 'uploading' && chunkRecord.metadata.timeoutThreshold && currentTime > chunkRecord.metadata.timeoutThreshold) {
status = 'timeout';
// 更新状态为超时
const timeoutMetadata = {
...chunkRecord.metadata,
@@ -880,13 +880,13 @@ export async function checkChunkUploadStatuses(env, uploadId, totalChunks) {
error: 'Upload timeout detected',
timeoutDetectedTime: currentTime
};
await db.put(chunkKey, chunkRecord.value, {
await db.put(chunkKey, chunkRecord.value, {
metadata: timeoutMetadata,
expirationTtl: 3600
}).catch(err => console.warn(`Failed to update timeout status for chunk ${i}:`, err));
}
let hasData = false;
if (status === 'completed') {
// 已完成的分块,不存储原始数据
@@ -898,7 +898,7 @@ export async function checkChunkUploadStatuses(env, uploadId, totalChunks) {
// 其他状态也检查是否有数据
hasData = (chunkRecord.value && chunkRecord.value.byteLength > 0);
}
chunkStatuses.push({
index: i,
key: chunkKey,
@@ -931,7 +931,7 @@ export async function checkChunkUploadStatuses(env, uploadId, totalChunks) {
});
}
}
return chunkStatuses;
}
@@ -943,11 +943,11 @@ export async function cleanupChunkData(env, uploadId, totalChunks) {
for (let i = 0; i < totalChunks; i++) {
const chunkKey = `chunk_${uploadId}_${i.toString().padStart(3, '0')}`;
// 删除数据库中的分块记录
await db.delete(chunkKey);
}
// 清理multipart info(如果存在)
const multipartKey = `multipart_${uploadId}`;
await db.delete(multipartKey);
@@ -985,31 +985,31 @@ export async function forceCleanupUpload(context, uploadId, totalChunks) {
await cleanupFailedMultipartUploads(context, uploadId, uploadChannel);
const cleanupPromises = [];
// 清理所有分块
for (let i = 0; i < totalChunks; i++) {
const chunkKey = `chunk_${uploadId}_${i.toString().padStart(3, '0')}`;
cleanupPromises.push(db.delete(chunkKey).catch(err =>
cleanupPromises.push(db.delete(chunkKey).catch(err =>
console.warn(`Failed to delete chunk ${i}:`, err)
));
}
// 清理相关的键
const keysToCleanup = [
`upload_session_${uploadId}`,
`multipart_${uploadId}`,
`merge_status_${uploadId}`
];
keysToCleanup.forEach(key => {
cleanupPromises.push(db.delete(key).catch(err =>
cleanupPromises.push(db.delete(key).catch(err =>
console.warn(`Failed to delete key ${key}:`, err)
));
});
await Promise.allSettled(cleanupPromises);
console.log(`Force cleanup completed for ${uploadId}`);
} catch (cleanupError) {
console.warn('Failed to force cleanup upload:', cleanupError);
}
@@ -1023,59 +1023,59 @@ export async function uploadLargeFileToTelegram(context, file, fullId, metadata,
const CHUNK_SIZE = 20 * 1024 * 1024; // 20MB
const fileSize = file.size;
const totalChunks = Math.ceil(fileSize / CHUNK_SIZE);
// 为了避免CPU超时,限制最大分片数(考虑Cloudflare Worker的CPU时间限制)
if (totalChunks > 50) {
return createResponse('Error: File too large (exceeds 1GB limit)', { status: 413 });
}
const chunks = [];
const uploadedChunks = [];
try {
// 分片上传,每10个分片做一次微小延迟以避免CPU超时
for (let i = 0; i < totalChunks; i++) {
const start = i * CHUNK_SIZE;
const end = Math.min(start + CHUNK_SIZE, fileSize);
const chunkBlob = file.slice(start, end);
// 生成分片文件名
const chunkFileName = `${fileName}.part${i.toString().padStart(3, '0')}`;
// 上传分片(带重试机制)
const chunkInfo = await uploadChunkToTelegramWithRetry(
tgBotToken,
tgChatId,
chunkBlob,
chunkFileName,
i,
tgBotToken,
tgChatId,
chunkBlob,
chunkFileName,
i,
totalChunks
);
if (!chunkInfo) {
throw new Error(`Failed to upload chunk ${i + 1}/${totalChunks} after retries`);
}
// 验证分片信息完整性
if (!chunkInfo.file_id || !chunkInfo.file_size) {
throw new Error(`Invalid chunk info for chunk ${i + 1}/${totalChunks}`);
}
chunks.push({
index: i,
fileId: chunkInfo.file_id,
size: chunkInfo.file_size,
fileName: chunkFileName
});
uploadedChunks.push(chunkInfo.file_id);
// 每10个分片检查一下,添加微小延迟避免CPU限制
if (i > 0 && i % 10 === 0) {
await new Promise(resolve => setTimeout(resolve, 50)); // 50ms延迟
}
}
// 所有分片上传成功,更新metadata
metadata.Channel = "TelegramNew";
metadata.ChannelName = tgChannel.name;
@@ -1085,15 +1085,15 @@ export async function uploadLargeFileToTelegram(context, file, fullId, metadata,
metadata.TotalChunks = totalChunks;
metadata.FileSize = (fileSize / 1024 / 1024).toFixed(2);
// 将分片信息存储到value中
const chunksData = JSON.stringify(chunks);
// 验证分片完整性
if (chunks.length !== totalChunks) {
throw new Error(`Chunk count mismatch: expected ${totalChunks}, got ${chunks.length}`);
}
// 写入最终的数据库记录,分片信息作为value
await db.put(fullId, chunksData, { metadata });
@@ -1104,12 +1104,12 @@ export async function uploadLargeFileToTelegram(context, file, fullId, metadata,
JSON.stringify([{ 'src': returnLink }]),
{
status: 200,
headers: {
headers: {
'Content-Type': 'application/json',
}
}
);
} catch (error) {
return createResponse(`Telegram Channel Error: Large file upload failed - ${error.message}`, { status: 500 });
}
@@ -1132,20 +1132,20 @@ async function uploadChunkToTelegramWithRetry(tgBotToken, tgChatId, chunkBlob, c
if (!fileInfo) {
throw new Error('Failed to extract file info from response');
}
return fileInfo;
} catch (error) {
console.warn(`Chunk ${chunkIndex} upload attempt ${attempt + 1} failed:`, error.message);
if (attempt === maxRetries - 1) {
return null; // 最后一次尝试也失败了
}
// 减少重试等待时间以节省CPU时间
await new Promise(resolve => setTimeout(resolve, 500 * (attempt + 1)));
}
}
return null;
}
+33 -38
View File
@@ -1,9 +1,11 @@
import { userAuthCheck, UnauthorizedResponse } from "../utils/userAuth";
import { fetchUploadConfig, fetchSecurityConfig } from "../utils/sysConfig";
import { createResponse, getUploadIp, getIPAddress, isExtValid,
moderateContent, purgeCDNCache, isBlockedUploadIp, buildUniqueFileId, endUpload } from "./uploadTools";
import { initializeChunkedUpload, handleChunkUpload, uploadLargeFileToTelegram, handleCleanupRequest} from "./chunkUpload";
import { handleChunkMerge, checkMergeStatus } from "./chunkMerge";
import {
createResponse, getUploadIp, getIPAddress, isExtValid,
moderateContent, purgeCDNCache, isBlockedUploadIp, buildUniqueFileId, endUpload
} from "./uploadTools";
import { initializeChunkedUpload, handleChunkUpload, uploadLargeFileToTelegram, handleCleanupRequest } from "./chunkUpload";
import { handleChunkMerge } from "./chunkMerge";
import { TelegramAPI } from "../utils/telegramAPI";
import { S3Client, PutObjectCommand } from "@aws-sdk/client-s3";
import { getDatabase } from '../utils/databaseAdapter.js';
@@ -19,7 +21,7 @@ export async function onRequest(context) { // Contents of context object
// 读取各项配置,存入 context
const securityConfig = await fetchSecurityConfig(env);
const uploadConfig = await fetchUploadConfig(env);
context.securityConfig = securityConfig;
context.uploadConfig = uploadConfig;
@@ -37,13 +39,6 @@ export async function onRequest(context) { // Contents of context object
return createResponse('Error: Your IP is blocked', { status: 403 });
}
// 检查是否为状态查询请求
const statusCheck = url.searchParams.get('statusCheck') === 'true';
if (statusCheck) {
const uploadId = url.searchParams.get('uploadId');
return await checkMergeStatus(env, uploadId);
}
// 检查是否为清理请求
const cleanupRequest = url.searchParams.get('cleanup') === 'true';
if (cleanupRequest) {
@@ -61,7 +56,7 @@ export async function onRequest(context) { // Contents of context object
// 检查是否为分块上传
const isChunked = url.searchParams.get('chunked') === 'true';
const isMerge = url.searchParams.get('merge') === 'true';
if (isChunked) {
if (isMerge) {
return await handleChunkMerge(context);
@@ -81,7 +76,7 @@ async function processFileUpload(context, formdata = null) {
// 解析表单数据
formdata = formdata || await request.formData();
// 将 formdata 存储在 context 中
context.formdata = formdata;
@@ -119,18 +114,18 @@ async function processFileUpload(context, formdata = null) {
const fileType = formdata.get('file').type;
let fileName = formdata.get('file').name;
const fileSize = (formdata.get('file').size / 1024 / 1024).toFixed(2); // 文件大小,单位MB
// 检查fileType和fileName是否存在
if (fileType === null || fileType === undefined || fileName === null || fileName === undefined) {
return createResponse('Error: fileType or fileName is wrong, check the integrity of this file!', { status: 400 });
}
// 如果上传文件夹路径为空,尝试从文件名中获取
if (uploadFolder === '' || uploadFolder === null || uploadFolder === undefined) {
uploadFolder = fileName.split('/').slice(0, -1).join('/');
}
// 处理文件夹路径格式,确保没有开头的/
const normalizedFolder = uploadFolder
const normalizedFolder = uploadFolder
? uploadFolder.replace(/^\/+/, '') // 移除开头的/
.replace(/\/{2,}/g, '/') // 替换多个连续的/为单个/
.replace(/\/$/, '') // 移除末尾的/
@@ -229,7 +224,7 @@ async function uploadFileToCloudflareR2(context, fullId, metadata, returnLink) {
}
const r2Channel = r2Settings.channels[0];
const R2DataBase = env.img_r2;
// 写入R2数据库
@@ -258,10 +253,10 @@ async function uploadFileToCloudflareR2(context, fullId, metadata, returnLink) {
// 成功上传,将文件ID返回给客户端
return createResponse(
JSON.stringify([{ 'src': `${returnLink}` }]),
JSON.stringify([{ 'src': `${returnLink}` }]),
{
status: 200,
headers: {
headers: {
'Content-Type': 'application/json',
}
}
@@ -364,7 +359,7 @@ async function uploadFileToS3(context, fullId, metadata, returnLink) {
return createResponse(JSON.stringify([{ src: returnLink }]), {
status: 200,
headers: {
headers: {
"Content-Type": "application/json",
},
});
@@ -382,7 +377,7 @@ async function uploadFileToTelegram(context, fullId, metadata, fileExt, fileName
// 选择一个 Telegram 渠道上传,若负载均衡开启,则随机选择一个;否则选择第一个
const tgSettings = uploadConfig.telegram;
const tgChannels = tgSettings.channels;
const tgChannel = tgSettings.loadBalance.enabled? tgChannels[Math.floor(Math.random() * tgChannels.length)] : tgChannels[0];
const tgChannel = tgSettings.loadBalance.enabled ? tgChannels[Math.floor(Math.random() * tgChannels.length)] : tgChannels[0];
if (!tgChannel) {
return createResponse('Error: No Telegram channel provided', { status: 400 });
}
@@ -393,10 +388,10 @@ async function uploadFileToTelegram(context, fullId, metadata, fileExt, fileName
const fileSize = file.size;
const telegramAPI = new TelegramAPI(tgBotToken);
// 20MB 分片阈值
const CHUNK_SIZE = 20 * 1024 * 1024; // 20MB
if (fileSize > CHUNK_SIZE) {
// 大文件分片上传
return await uploadLargeFileToTelegram(env, file, fullId, metadata, fileName, fileType, url, returnLink, tgBotToken, tgChatId, tgChannel);
@@ -415,28 +410,28 @@ async function uploadFileToTelegram(context, fullId, metadata, fileExt, fileName
// 选择对应的发送接口
const fileTypeMap = {
'image/': {'url': 'sendPhoto', 'type': 'photo'},
'video/': {'url': 'sendVideo', 'type': 'video'},
'audio/': {'url': 'sendAudio', 'type': 'audio'},
'application/pdf': {'url': 'sendDocument', 'type': 'document'},
'image/': { 'url': 'sendPhoto', 'type': 'photo' },
'video/': { 'url': 'sendVideo', 'type': 'video' },
'audio/': { 'url': 'sendAudio', 'type': 'audio' },
'application/pdf': { 'url': 'sendDocument', 'type': 'document' },
};
const defaultType = {'url': 'sendDocument', 'type': 'document'};
const defaultType = { 'url': 'sendDocument', 'type': 'document' };
let sendFunction = Object.keys(fileTypeMap).find(key => fileType.startsWith(key))
? fileTypeMap[Object.keys(fileTypeMap).find(key => fileType.startsWith(key))]
let sendFunction = Object.keys(fileTypeMap).find(key => fileType.startsWith(key))
? fileTypeMap[Object.keys(fileTypeMap).find(key => fileType.startsWith(key))]
: defaultType;
// GIF 发送接口特殊处理
if (fileType === 'image/gif' || fileType === 'image/webp' || fileExt === 'gif' || fileExt === 'webp') {
sendFunction = {'url': 'sendAnimation', 'type': 'animation'};
sendFunction = { 'url': 'sendAnimation', 'type': 'animation' };
}
// 根据服务端压缩设置处理接口:从参数中获取serverCompress,如果为false,则使用sendDocument接口
if (url.searchParams.get('serverCompress') === 'false') {
sendFunction = {'url': 'sendDocument', 'type': 'document'};
sendFunction = { 'url': 'sendDocument', 'type': 'document' };
}
// 上传文件到 Telegram
let res = createResponse('upload error, check your environment params about telegram channel!', { status: 400 });
try {
@@ -452,7 +447,7 @@ async function uploadFileToTelegram(context, fullId, metadata, fileExt, fileName
JSON.stringify([{ 'src': `${returnLink}` }]),
{
status: 200,
headers: {
headers: {
'Content-Type': 'application/json',
}
}
@@ -517,10 +512,10 @@ async function uploadFileToExternal(context, fullId, metadata, returnLink) {
// 返回结果
return createResponse(
JSON.stringify([{ 'src': `${returnLink}` }]),
JSON.stringify([{ 'src': `${returnLink}` }]),
{
status: 200,
headers: {
headers: {
'Content-Type': 'application/json',
}
}
@@ -535,7 +530,7 @@ async function tryRetry(err, context, uploadChannel, fullId, metadata, fileExt,
const channelList = ['CloudflareR2', 'TelegramNew', 'S3'];
const errMessages = {};
errMessages[uploadChannel] = 'Error: ' + uploadChannel + err;
for (let i = 0; i < channelList.length; i++) {
if (channelList[i] !== uploadChannel) {
let res = null;
+1 -1
View File
@@ -1,4 +1,4 @@
<!doctype html><html lang=""><head><meta charset="utf-8"><meta http-equiv="X-UA-Compatible" content="IE=edge"><meta name="viewport" content="width=device-width,initial-scale=1"><link rel="icon" href="/logo.png"><link rel="apple-touch-icon" href="/logo.png"><link rel="mask-icon" href="/logo.png" color="#f4b400"><meta name="description" content="Sanyue ImgHub - A modern file hosting platform"><meta name="keywords" content="Sanyue, ImgHub, file hosting, image hosting, cloud storage"><meta name="author" content="SanyueQi"><title>Sanyue ImgHub</title><script defer="defer" src="/js/app.07ad4eb5.js"></script><link href="/css/app.dab02a4f.css" rel="stylesheet"></head><body><noscript><strong>We're sorry but sanyue_imghub doesn't work properly without JavaScript enabled. Please enable it to continue.</strong></noscript><div id="app"></div></body></html><style>/* 下拉菜单样式 */
<!doctype html><html lang=""><head><meta charset="utf-8"><meta http-equiv="X-UA-Compatible" content="IE=edge"><meta name="viewport" content="width=device-width,initial-scale=1"><link rel="icon" href="/logo.png"><link rel="apple-touch-icon" href="/logo.png"><link rel="mask-icon" href="/logo.png" color="#f4b400"><meta name="description" content="Sanyue ImgHub - A modern file hosting platform"><meta name="keywords" content="Sanyue, ImgHub, file hosting, image hosting, cloud storage"><meta name="author" content="SanyueQi"><title>Sanyue ImgHub</title><script defer="defer" src="/js/app.4e994b97.js"></script><link href="/css/app.dab02a4f.css" rel="stylesheet"></head><body><noscript><strong>We're sorry but sanyue_imghub doesn't work properly without JavaScript enabled. Please enable it to continue.</strong></noscript><div id="app"></div></body></html><style>/* 下拉菜单样式 */
.el-dropdown__popper.el-popper {
border-radius: 12px;
border: none;
BIN
View File
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
File diff suppressed because one or more lines are too long
Binary file not shown.
File diff suppressed because one or more lines are too long