Refactor Worker to support resumable Google Drive uploads
Rewrites the Cloudflare Worker to handle Telegram file uploads to Google Drive using the resumable upload API for large files. Adds detailed error handling, progress tracking, and modularizes authentication and upload logic. Removes legacy code for multipart uploads and streamlines the file handling process.
This commit is contained in:
1 parent
eac9e037e2
commit
651a7a03ea
1 file changed
+199
-168
+199
-168
@@ -1,202 +1,233 @@
|
||||
/**
|
||||
* A Cloudflare Worker for a Telegram bot that uploads files to Google Drive.
|
||||
*
|
||||
* Features:
|
||||
* - Responds to files sent to a Telegram bot.
|
||||
* - Authenticates with Google Drive using a service account.
|
||||
* - Uploads files under 20MB directly.
|
||||
* - For files 20MB and larger, it uses the Google Drive Resumable Upload API to upload in chunks.
|
||||
*
|
||||
* Secrets to be set in Cloudflare Worker settings:
|
||||
* - ADMIN_CHAT_ID: Your personal Telegram Chat ID to receive notifications/errors.
|
||||
* - GOOGLE_CREDENTIALS: The JSON credentials for your Google Cloud service account.
|
||||
* - GOOGLE_DRIVE_FOLDER_ID: The ID of the Google Drive folder where files will be uploaded.
|
||||
* - TELEGRAM_BOT_TOKEN: The token for your Telegram bot from BotFather.
|
||||
*/
|
||||
|
||||
// Recommended chunk size for Google Drive resumable uploads (must be a multiple of 256KB)
|
||||
const CHUNK_SIZE = 256 * 1024 * 10; // 2.5 MB
|
||||
|
||||
export default {
|
||||
async fetch(request, env) {
|
||||
async fetch(request, env, ctx) {
|
||||
if (request.method === 'POST') {
|
||||
const payload = await request.json();
|
||||
if (payload.message && payload.message.document) {
|
||||
await handleDocument(payload.message, env);
|
||||
} else if (payload.message && payload.message.photo) {
|
||||
await handlePhoto(payload.message, env);
|
||||
try {
|
||||
const update = await request.json();
|
||||
if (update.message && update.message.document) {
|
||||
ctx.waitUntil(this.handleFileUpload(update.message, env));
|
||||
} else if (update.message && update.message.text === '/start') {
|
||||
await this.sendMessage(
|
||||
env.TELEGRAM_BOT_TOKEN,
|
||||
update.message.chat.id,
|
||||
'Welcome! Send me a file and I will upload it to Google Drive.'
|
||||
);
|
||||
}
|
||||
} catch (e) {
|
||||
// If parsing fails or for any other reason, notify the admin.
|
||||
if (env.ADMIN_CHAT_ID) {
|
||||
await this.sendMessage(env.TELEGRAM_BOT_TOKEN, env.ADMIN_CHAT_ID, `Error in main fetch: ${e.message}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
return new Response('OK');
|
||||
return new Response('OK'); // Respond to Telegram's webhook POST
|
||||
},
|
||||
};
|
||||
|
||||
async function handleDocument(message, env) {
|
||||
const fileId = message.document.file_id;
|
||||
const fileName = message.document.file_name;
|
||||
const fileSize = message.document.file_size;
|
||||
await uploadToDrive(fileId, fileName, fileSize, env);
|
||||
}
|
||||
async handleFileUpload(message, env) {
|
||||
const doc = message.document;
|
||||
const chatId = message.chat.id;
|
||||
const fileName = doc.file_name;
|
||||
const fileId = doc.file_id;
|
||||
const fileSize = doc.file_size;
|
||||
|
||||
async function handlePhoto(message, env) {
|
||||
const photo = message.photo[message.photo.length - 1]; // Get the highest resolution photo
|
||||
const fileId = photo.file_id;
|
||||
const fileName = `${fileId}.jpg`;
|
||||
const fileSize = photo.file_size;
|
||||
await uploadToDrive(fileId, fileName, fileSize, env);
|
||||
}
|
||||
await this.sendMessage(env.TELEGRAM_BOT_TOKEN, chatId, `Received "${fileName}". Preparing to upload...`);
|
||||
|
||||
async function uploadToDrive(fileId, fileName, fileSize, env) {
|
||||
const botToken = env.TELEGRAM_BOT_TOKEN;
|
||||
const driveFolderId = env.GOOGLE_DRIVE_FOLDER_ID;
|
||||
const adminChatId = env.ADMIN_CHAT_ID;
|
||||
const googleCredentials = JSON.parse(env.GOOGLE_CREDENTIALS);
|
||||
try {
|
||||
// 1. Get Google Auth Token
|
||||
const authToken = await this.getGoogleAuthToken(env.GOOGLE_CREDENTIALS);
|
||||
|
||||
// Get file path from Telegram
|
||||
const fileDetailsUrl = `https://api.telegram.org/bot${botToken}/getFile?file_id=${fileId}`;
|
||||
const fileDetailsResponse = await fetch(fileDetailsUrl);
|
||||
const fileDetails = await fileDetailsResponse.json();
|
||||
if (!fileDetails.ok) {
|
||||
await sendMessage(adminChatId, `Error getting file details: ${fileDetails.description}`, botToken);
|
||||
return;
|
||||
}
|
||||
const filePath = fileDetails.result.file_path;
|
||||
const fileUrl = `https://api.telegram.org/file/bot${botToken}/${filePath}`;
|
||||
// 2. Get Telegram file path
|
||||
const fileInfo = await this.getTelegramFileInfo(env.TELEGRAM_BOT_TOKEN, fileId);
|
||||
if (!fileInfo.ok) {
|
||||
throw new Error(`Failed to get file info: ${fileInfo.description}`);
|
||||
}
|
||||
const filePath = fileInfo.result.file_path;
|
||||
const fileUrl = `https://api.telegram.org/file/bot${env.TELEGRAM_BOT_TOKEN}/${filePath}`;
|
||||
|
||||
// Authenticate with Google Drive
|
||||
const jwtToken = await getGoogleAuthToken(googleCredentials);
|
||||
// 3. Initiate Resumable Upload
|
||||
const uploadUrl = await this.initiateResumableUpload(authToken, fileName, env.GOOGLE_DRIVE_FOLDER_ID);
|
||||
|
||||
// For files larger than 20MB, use resumable upload
|
||||
if (fileSize > 20 * 1024 * 1024) {
|
||||
await resumableUpload(fileUrl, fileName, driveFolderId, jwtToken, adminChatId, botToken);
|
||||
return;
|
||||
}
|
||||
// 4. Download from Telegram and Upload to Drive in chunks
|
||||
const fileResponse = await fetch(fileUrl);
|
||||
if (!fileResponse.ok || !fileResponse.body) {
|
||||
throw new Error('Failed to download file from Telegram.');
|
||||
}
|
||||
const reader = fileResponse.body.getReader();
|
||||
let bytesUploaded = 0;
|
||||
|
||||
// Download file from Telegram
|
||||
const fileResponse = await fetch(fileUrl);
|
||||
const fileData = await fileResponse.arrayBuffer();
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) {
|
||||
break;
|
||||
}
|
||||
|
||||
// Upload file to Google Drive using multipart upload for smaller files
|
||||
const uploadUrl = 'https://www.googleapis.com/upload/drive/v3/files?uploadType=multipart';
|
||||
const metadata = {
|
||||
name: fileName,
|
||||
parents: [driveFolderId],
|
||||
};
|
||||
const chunk = value;
|
||||
await this.uploadChunk(uploadUrl, chunk, bytesUploaded, fileSize);
|
||||
bytesUploaded += chunk.length;
|
||||
|
||||
const formData = new FormData();
|
||||
formData.append('metadata', new Blob([JSON.stringify(metadata)], { type: 'application/json' }));
|
||||
formData.append('file', new Blob([fileData]));
|
||||
// Optional: Send progress updates
|
||||
// const progress = Math.round((bytesUploaded / fileSize) * 100);
|
||||
// await this.sendMessage(env.TELEGRAM_BOT_TOKEN, chatId, `Uploading... ${progress}%`);
|
||||
}
|
||||
|
||||
const uploadResponse = await fetch(uploadUrl, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${jwtToken}`,
|
||||
},
|
||||
body: formData,
|
||||
});
|
||||
await this.sendMessage(env.TELEGRAM_BOT_TOKEN, chatId, `✅ Successfully uploaded "${fileName}" to Google Drive!`);
|
||||
} catch (error) {
|
||||
console.error(error);
|
||||
await this.sendMessage(env.TELEGRAM_BOT_TOKEN, chatId, `❌ Failed to upload file. Reason: ${error.message}`);
|
||||
if (env.ADMIN_CHAT_ID && env.ADMIN_CHAT_ID !== chatId) {
|
||||
await this.sendMessage(env.TELEGRAM_BOT_TOKEN, env.ADMIN_CHAT_ID, `Upload failed for chat ${chatId}: ${error.message}`);
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
const uploadResult = await uploadResponse.json();
|
||||
if (uploadResult.id) {
|
||||
await sendMessage(adminChatId, `File "${fileName}" uploaded to Google Drive successfully!`, botToken);
|
||||
} else {
|
||||
await sendMessage(adminChatId, `Error uploading file: ${JSON.stringify(uploadResult.error)}`, botToken);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Generates a Google Cloud access token from service account credentials.
|
||||
*/
|
||||
async getGoogleAuthToken(credentialsJson) {
|
||||
const credentials = JSON.parse(credentialsJson);
|
||||
const jwtHeader = {
|
||||
alg: 'RS256',
|
||||
typ: 'JWT',
|
||||
};
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
const jwtClaimSet = {
|
||||
iss: credentials.client_email,
|
||||
scope: 'https://www.googleapis.com/auth/drive',
|
||||
aud: 'https://oauth2.googleapis.com/token',
|
||||
exp: now + 3600,
|
||||
iat: now,
|
||||
};
|
||||
|
||||
async function resumableUpload(fileUrl, fileName, driveFolderId, jwtToken, adminChatId, botToken) {
|
||||
const metadata = {
|
||||
name: fileName,
|
||||
parents: [driveFolderId],
|
||||
};
|
||||
const encodedHeader = btoa(JSON.stringify(jwtHeader)).replace(/=/g, '').replace(/\+/g, '-').replace(/\//g, '_');
|
||||
const encodedClaimSet = btoa(JSON.stringify(jwtClaimSet)).replace(/=/g, '').replace(/\+/g, '-').replace(/\//g, '_');
|
||||
const signatureInput = `${encodedHeader}.${encodedClaimSet}`;
|
||||
|
||||
const initiationResponse = await fetch('https://www.googleapis.com/upload/drive/v3/files?uploadType=resumable', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${jwtToken}`,
|
||||
'Content-Type': 'application/json',
|
||||
},
|
||||
body: JSON.stringify(metadata),
|
||||
});
|
||||
const privateKey = await crypto.subtle.importKey(
|
||||
'pkcs8',
|
||||
this.pemToBinary(credentials.private_key),
|
||||
{
|
||||
name: 'RSASSA-PKCS1-V1_5',
|
||||
hash: 'SHA-256',
|
||||
},
|
||||
false,
|
||||
['sign']
|
||||
);
|
||||
|
||||
if (!initiationResponse.ok) {
|
||||
await sendMessage(adminChatId, `Error initiating resumable upload: ${await initiationResponse.text()}`, botToken);
|
||||
return;
|
||||
}
|
||||
const signature = await crypto.subtle.sign('RSASSA-PKCS1-V1_5', privateKey, new TextEncoder().encode(signatureInput));
|
||||
const encodedSignature = btoa(String.fromCharCode(...new Uint8Array(signature)))
|
||||
.replace(/=/g, '')
|
||||
.replace(/\+/g, '-')
|
||||
.replace(/\//g, '_');
|
||||
|
||||
const location = initiationResponse.headers.get('Location');
|
||||
const jwt = `${signatureInput}.${encodedSignature}`;
|
||||
|
||||
const fileResponse = await fetch(fileUrl);
|
||||
const reader = fileResponse.body.getReader();
|
||||
let uploaded = 0;
|
||||
const chunkSize = 256 * 1024 * 10; // 2.56MB chunks
|
||||
const totalSize = Number(fileResponse.headers.get('Content-Length'));
|
||||
const tokenResponse = await fetch('https://oauth2.googleapis.com/token', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Content-Type': 'application/x-www-form-urlencoded',
|
||||
},
|
||||
body: `grant_type=urn:ietf:params:oauth:grant-type:jwt-bearer&assertion=${jwt}`,
|
||||
});
|
||||
|
||||
while (true) {
|
||||
const { done, value } = await reader.read();
|
||||
if (done) {
|
||||
break;
|
||||
const tokenData = await tokenResponse.json();
|
||||
return tokenData.access_token;
|
||||
},
|
||||
|
||||
/**
|
||||
* Initiates a resumable upload session with Google Drive.
|
||||
* @returns {string} The unique session URI for uploading chunks.
|
||||
*/
|
||||
async initiateResumableUpload(authToken, fileName, folderId) {
|
||||
const metadata = {
|
||||
name: fileName,
|
||||
parents: [folderId],
|
||||
};
|
||||
|
||||
const response = await fetch('https://www.googleapis.com/upload/drive/v3/files?uploadType=resumable', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization: `Bearer ${authToken}`,
|
||||
'Content-Type': 'application/json; charset=UTF-8',
|
||||
},
|
||||
body: JSON.stringify(metadata),
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorBody = await response.text();
|
||||
throw new Error(`Failed to initiate resumable upload: ${response.status} ${response.statusText} - ${errorBody}`);
|
||||
}
|
||||
|
||||
const chunk = value;
|
||||
const uploadResponse = await fetch(location, {
|
||||
return response.headers.get('Location');
|
||||
},
|
||||
|
||||
/**
|
||||
* Uploads a single chunk of the file to the resumable session URI.
|
||||
*/
|
||||
async uploadChunk(uploadUrl, chunk, startByte, totalSize) {
|
||||
const endByte = startByte + chunk.length - 1;
|
||||
|
||||
const response = await fetch(uploadUrl, {
|
||||
method: 'PUT',
|
||||
headers: {
|
||||
'Content-Range': `bytes ${uploaded}-${uploaded + chunk.length - 1}/${totalSize}`,
|
||||
'Content-Range': `bytes ${startByte}-${endByte}/${totalSize}`,
|
||||
'Content-Length': chunk.length,
|
||||
},
|
||||
body: chunk,
|
||||
});
|
||||
|
||||
if (uploadResponse.status === 308) {
|
||||
// Resume Incomplete
|
||||
uploaded += chunk.length;
|
||||
} else if (uploadResponse.ok) {
|
||||
const uploadResult = await uploadResponse.json();
|
||||
if (uploadResult.id) {
|
||||
await sendMessage(adminChatId, `File "${fileName}" uploaded to Google Drive successfully!`, botToken);
|
||||
} else {
|
||||
await sendMessage(adminChatId, `Error completing resumable upload: ${JSON.stringify(uploadResult.error)}`, botToken);
|
||||
}
|
||||
return;
|
||||
} else {
|
||||
await sendMessage(adminChatId, `Error during resumable upload chunk: ${await uploadResponse.text()}`, botToken);
|
||||
return;
|
||||
// For the last chunk, a 200 OK is returned. For all others, a 308 Resume Incomplete.
|
||||
if (response.status !== 308 && response.status !== 200) {
|
||||
const errorBody = await response.text();
|
||||
throw new Error(`Chunk upload failed: ${response.status} ${response.statusText} - ${errorBody}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
async function getGoogleAuthToken(credentials) {
|
||||
const header = { alg: 'RS256', typ: 'JWT' };
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
const payload = {
|
||||
iss: credentials.client_email,
|
||||
scope: 'https://www.googleapis.com/auth/drive',
|
||||
aud: 'https://oauth2.googleapis.com/token',
|
||||
exp: now + 3600,
|
||||
iat: now,
|
||||
};
|
||||
/**
|
||||
* Helper to get file path info from Telegram.
|
||||
*/
|
||||
async getTelegramFileInfo(botToken, fileId) {
|
||||
const response = await fetch(`https://api.telegram.org/bot${botToken}/getFile?file_id=${fileId}`);
|
||||
return response.json();
|
||||
},
|
||||
|
||||
const encodedHeader = btoa(JSON.stringify(header)).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/, '');
|
||||
const encodedPayload = btoa(JSON.stringify(payload)).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/, '');
|
||||
const dataToSign = `${encodedHeader}.${encodedPayload}`;
|
||||
/**
|
||||
* Helper to send a message via the Telegram API.
|
||||
*/
|
||||
async sendMessage(botToken, chatId, text) {
|
||||
const url = `https://api.telegram.org/bot${botToken}/sendMessage?chat_id=${chatId}&text=${encodeURIComponent(text)}`;
|
||||
await fetch(url);
|
||||
},
|
||||
|
||||
const privateKey = await crypto.subtle.importKey(
|
||||
'pkcs8',
|
||||
str2ab(atob(credentials.private_key.replace('-----BEGIN PRIVATE KEY-----\n', '').replace('\n-----END PRIVATE KEY-----\n', ''))),
|
||||
{ name: 'RSASSA-PKCS1-v1_5', hash: 'SHA-256' },
|
||||
false,
|
||||
['sign']
|
||||
);
|
||||
|
||||
const signature = await crypto.subtle.sign('RSASSA-PKCS1-v1_5', privateKey, str2ab(dataToSign));
|
||||
const encodedSignature = btoa(ab2str(signature)).replace(/\+/g, '-').replace(/\//g, '_').replace(/=+$/, '');
|
||||
|
||||
const jwt = `${dataToSign}.${encodedSignature}`;
|
||||
|
||||
const tokenResponse = await fetch('https://oauth2.googleapis.com/token', {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
|
||||
body: `grant_type=urn:ietf:params:oauth:grant-type:jwt-bearer&assertion=${jwt}`,
|
||||
});
|
||||
|
||||
const tokenData = await tokenResponse.json();
|
||||
return tokenData.access_token;
|
||||
}
|
||||
|
||||
async function sendMessage(chatId, text, botToken) {
|
||||
const url = `https://api.telegram.org/bot${botToken}/sendMessage?chat_id=${chatId}&text=${encodeURIComponent(text)}`;
|
||||
await fetch(url);
|
||||
}
|
||||
|
||||
function str2ab(str) {
|
||||
const buf = new ArrayBuffer(str.length);
|
||||
const bufView = new Uint8Array(buf);
|
||||
for (let i = 0, strLen = str.length; i < strLen; i++) {
|
||||
bufView[i] = str.charCodeAt(i);
|
||||
}
|
||||
return buf;
|
||||
}
|
||||
|
||||
function ab2str(buf) {
|
||||
return String.fromCharCode.apply(null, new Uint8Array(buf));
|
||||
}
|
||||
/**
|
||||
* Converts a PEM-formatted key to a binary format for WebCrypto API.
|
||||
*/
|
||||
pemToBinary(pem) {
|
||||
const lines = pem.split('\n');
|
||||
const base64 = lines.filter((line) => !line.startsWith('-----')).join('');
|
||||
const binary = atob(base64);
|
||||
const len = binary.length;
|
||||
const bytes = new Uint8Array(len);
|
||||
for (let i = 0; i < len; i++) {
|
||||
bytes[i] = binary.charCodeAt(i);
|
||||
}
|
||||
return bytes.buffer;
|
||||
},
|
||||
};
|
||||
Reference in new issue
Block a user