diff --git a/bot/helper/mirror_utils/upload_utils/gdriveTools.py b/bot/helper/mirror_utils/upload_utils/gdriveTools.py index 62bd6ef..646d17e 100644 --- a/bot/helper/mirror_utils/upload_utils/gdriveTools.py +++ b/bot/helper/mirror_utils/upload_utils/gdriveTools.py @@ -19,50 +19,20 @@ from bot.helper.ext_utils.bot_utils import * from bot.helper.ext_utils.fs_utils import get_mime_type logging.getLogger('googleapiclient.discovery').setLevel(logging.ERROR) - -G_DRIVE_TOKEN_FILE = "token.pickle" -# Check https://developers.google.com/drive/scopes for all available scopes -OAUTH_SCOPE = ["https://www.googleapis.com/auth/drive"] - SERVICE_ACCOUNT_INDEX = 0 -def authorize(): - # Get credentials - credentials = None - if not USE_SERVICE_ACCOUNTS: - if os.path.exists(G_DRIVE_TOKEN_FILE): - with open(G_DRIVE_TOKEN_FILE, 'rb') as f: - credentials = pickle.load(f) - if credentials is None or not credentials.valid: - if credentials and credentials.expired and credentials.refresh_token: - credentials.refresh(Request()) - else: - flow = InstalledAppFlow.from_client_secrets_file( - 'credentials.json', OAUTH_SCOPE) - LOGGER.info(flow) - credentials = flow.run_console(port=0) - - # Save the credentials for the next run - with open(G_DRIVE_TOKEN_FILE, 'wb') as token: - pickle.dump(credentials, token) - else: - credentials = service_account.Credentials \ - .from_service_account_file(f'accounts/{SERVICE_ACCOUNT_INDEX}.json', - scopes=OAUTH_SCOPE) - return build('drive', 'v3', credentials=credentials, cache_discovery=False) - - -service = authorize() - - class GoogleDriveHelper: - # Redirect URI for installed apps, can be left as is - REDIRECT_URI = "urn:ietf:wg:oauth:2.0:oob" - G_DRIVE_DIR_MIME_TYPE = "application/vnd.google-apps.folder" - G_DRIVE_BASE_DOWNLOAD_URL = "https://drive.google.com/uc?id={}&export=download" - def __init__(self, name=None, listener=None): + self.__G_DRIVE_TOKEN_FILE = "token.pickle" + # Check https://developers.google.com/drive/scopes for all available scopes + self.__OAUTH_SCOPE = ['https://www.googleapis.com/auth/drive'] + # Redirect URI for installed apps, can be left as is + self.__REDIRECT_URI = "urn:ietf:wg:oauth:2.0:oob" + self.__G_DRIVE_DIR_MIME_TYPE = "application/vnd.google-apps.folder" + self.__G_DRIVE_BASE_DOWNLOAD_URL = "https://drive.google.com/uc?id={}&export=download" + self.__listener = listener + self.__service = self.authorize() self.__listener = listener self._file_uploaded_bytes = 0 self.uploaded_bytes = 0 @@ -119,8 +89,8 @@ class GoogleDriveHelper: } if parent_id is not None: file_metadata['parents'] = [parent_id] - return service.files().create(supportsTeamDrives=True, - body=file_metadata, media_body=media_body).execute() + return self.__service.files().create(supportsTeamDrives=True, + body=file_metadata, media_body=media_body).execute() @retry(wait=wait_exponential(multiplier=2, min=3, max=6), stop=stop_after_attempt(5), retry=retry_if_exception_type(HttpError), before=before_log(LOGGER, logging.DEBUG)) @@ -131,13 +101,11 @@ class GoogleDriveHelper: 'value': None, 'withLink': True } - return service.permissions().create(supportsTeamDrives=True, fileId=drive_id, body=permissions).execute() + return self.__service.permissions().create(supportsTeamDrives=True, fileId=drive_id, body=permissions).execute() @retry(wait=wait_exponential(multiplier=2, min=3, max=6), stop=stop_after_attempt(5), retry=retry_if_exception_type(HttpError), before=before_log(LOGGER, logging.DEBUG)) def upload_file(self, file_path, file_name, mime_type, parent_id): - global SERVICE_ACCOUNT_INDEX - global service # File body description file_metadata = { 'name': file_name, @@ -151,13 +119,14 @@ class GoogleDriveHelper: media_body = MediaFileUpload(file_path, mimetype=mime_type, resumable=False) - response = service.files().create(supportsTeamDrives=True, - body=file_metadata, media_body=media_body).execute() + response = self.__service.files().create(supportsTeamDrives=True, + body=file_metadata, media_body=media_body).execute() if not IS_TEAM_DRIVE: self.__set_permission(response['id']) - drive_file = service.files().get(supportsTeamDrives=True, - fileId=response['id']).execute() - download_url = self.G_DRIVE_BASE_DOWNLOAD_URL.format(drive_file.get('id')) + + drive_file = self.__service.files().get(supportsTeamDrives=True, + fileId=response['id']).execute() + download_url = self.__G_DRIVE_BASE_DOWNLOAD_URL.format(drive_file.get('id')) return download_url media_body = MediaFileUpload(file_path, mimetype=mime_type, @@ -165,8 +134,8 @@ class GoogleDriveHelper: chunksize=50 * 1024 * 1024) # Insert a file - drive_file = service.files().create(supportsTeamDrives=True, - body=file_metadata, media_body=media_body) + drive_file = self.__service.files().create(supportsTeamDrives=True, + body=file_metadata, media_body=media_body) response = None while response is None: if self.is_cancelled: @@ -177,16 +146,18 @@ class GoogleDriveHelper: if err.resp.get('content-type', '').startswith('application/json'): reason = json.loads(err.content).get('error').get('errors')[0].get('reason') if reason == 'userRateLimitExceeded': + global SERVICE_ACCOUNT_INDEX SERVICE_ACCOUNT_INDEX += 1 - service = authorize() + LOGGER.info(f"Switching to {SERVICE_ACCOUNT_INDEX}.json service account") + self.__service = self.authorize() raise err self._file_uploaded_bytes = 0 # Insert new permissions if not IS_TEAM_DRIVE: self.__set_permission(response['id']) # Define file instance and get url for download - drive_file = service.files().get(supportsTeamDrives=True, fileId=response['id']).execute() - download_url = self.G_DRIVE_BASE_DOWNLOAD_URL.format(drive_file.get('id')) + drive_file = self.__service.files().get(supportsTeamDrives=True, fileId=response['id']).execute() + download_url = self.__G_DRIVE_BASE_DOWNLOAD_URL.format(drive_file.get('id')) return download_url def upload(self, file_name: str): @@ -204,9 +175,13 @@ class GoogleDriveHelper: raise Exception('Upload has been manually cancelled') LOGGER.info("Uploaded To G-Drive: " + file_path) except Exception as e: - LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") - LOGGER.error(e.last_attempt.exception()) - self.__listener.onUploadError(e) + if isinstance(e, RetryError): + LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") + err = e.last_attempt.exception() + else: + err = e + LOGGER.error(err) + self.__listener.onUploadError(str(err)) return finally: self.updater.cancel() @@ -219,9 +194,13 @@ class GoogleDriveHelper: LOGGER.info("Uploaded To G-Drive: " + file_name) link = f"https://drive.google.com/folderview?id={dir_id}" except Exception as e: - LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") - LOGGER.error(e.last_attempt.exception()) - self.__listener.onUploadError(e) + if isinstance(e, RetryError): + LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") + err = e.last_attempt.exception() + else: + err = e + LOGGER.error(err) + self.__listener.onUploadError(str(err)) return finally: self.updater.cancel() @@ -230,71 +209,72 @@ class GoogleDriveHelper: LOGGER.info("Deleting downloaded file/folder..") return link - @retry(wait=wait_exponential(multiplier=2, min=3, max=6),stop=stop_after_attempt(5),retry=retry_if_exception_type(HttpError),before=before_log(LOGGER,logging.DEBUG)) - def copyFile(self,file_id,dest_id): + @retry(wait=wait_exponential(multiplier=2, min=3, max=6), stop=stop_after_attempt(5), + retry=retry_if_exception_type(HttpError), before=before_log(LOGGER, logging.DEBUG)) + def copyFile(self, file_id, dest_id): body = { 'parents': [dest_id] } - return self.__service.files().copy(supportsAllDrives=True,fileId=file_id,body=body).execute() + return self.__service.files().copy(supportsAllDrives=True, fileId=file_id, body=body).execute() - def clone(self,link): + def clone(self, link): self.transferred_size = 0 file_id = self.getIdFromUrl(link) msg = "" LOGGER.info(f"File ID: {file_id}") try: - meta = self.__service.files().get(supportsAllDrives=True,fileId=file_id,fields="name,id,mimeType,size").execute() + meta = self.__service.files().get(supportsAllDrives=True, fileId=file_id, + fields="name,id,mimeType,size").execute() except Exception as e: - return f"{str(e).replace('>','').replace('<','')}" + return f"{str(e).replace('>', '').replace('<', '')}" if meta.get("mimeType") == self.__G_DRIVE_DIR_MIME_TYPE: - dir_id = self.create_directory(meta.get('name'),parent_id) + dir_id = self.create_directory(meta.get('name'), parent_id) try: - result = self.cloneFolder(meta.get('name'),meta.get('name'),meta.get('id'),dir_id) + result = self.cloneFolder(meta.get('name'), meta.get('name'), meta.get('id'), dir_id) except Exception as e: - if isinstance(e,RetryError): + if isinstance(e, RetryError): LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") err = e.last_attempt.exception() else: - err = str(e).replace('>','').replace('<','') + err = str(e).replace('>', '').replace('<', '') LOGGER.error(err) return err msg += f'{meta.get("name")} ({get_readable_file_size(self.transferred_size)})' else: - file = self.copyFile(meta.get('id'),parent_id) + file = self.copyFile(meta.get('id'), parent_id) msg += f'{meta.get("name")} ({get_readable_file_size(int(meta.get("size")))})' return msg - - def cloneFolder(self,name,local_path,folder_id,parent_id): + def cloneFolder(self, name, local_path, folder_id, parent_id): page_token = None - q =f"'{folder_id}' in parents" + q = f"'{folder_id}' in parents" files = [] LOGGER.info(f"Syncing: {local_path}") new_id = None while True: response = self.__service.files().list(q=q, - spaces='drive', - fields='nextPageToken, files(id, name, mimeType,size)', - pageToken=page_token).execute() + spaces='drive', + fields='nextPageToken, files(id, name, mimeType,size)', + pageToken=page_token).execute() for file in response.get('files', []): files.append(file) page_token = response.get('nextPageToken', None) if page_token is None: - break + break if len(files) == 0: return parent_id for file in files: if file.get('mimeType') == self.__G_DRIVE_DIR_MIME_TYPE: - file_path = os.path.join(local_path,file.get('name')) - current_dir_id = self.create_directory(file.get('name'),parent_id) - new_id = self.cloneFolder(file.get('name'),file_path,file.get('id'),current_dir_id) + file_path = os.path.join(local_path, file.get('name')) + current_dir_id = self.create_directory(file.get('name'), parent_id) + new_id = self.cloneFolder(file.get('name'), file_path, file.get('id'), current_dir_id) else: self.transferred_size += int(file.get('size')) try: - self.copyFile(file.get('id'),parent_id) + self.copyFile(file.get('id'), parent_id) new_id = parent_id except Exception as e: - if isinstance(e,RetryError): + if isinstance(e, RetryError): LOGGER.info(f"Total Attempts: {e.last_attempt.attempt_number}") err = e.last_attempt.exception() else: @@ -307,11 +287,11 @@ class GoogleDriveHelper: def create_directory(self, directory_name, parent_id): file_metadata = { "name": directory_name, - "mimeType": self.G_DRIVE_DIR_MIME_TYPE + "mimeType": self.__G_DRIVE_DIR_MIME_TYPE } if parent_id is not None: file_metadata["parents"] = [parent_id] - file = service.files().create(supportsTeamDrives=True, body=file_metadata).execute() + file = self.__service.files().create(supportsTeamDrives=True, body=file_metadata).execute() file_id = file.get("id") if not IS_TEAM_DRIVE: self.__set_permission(file_id) @@ -338,22 +318,48 @@ class GoogleDriveHelper: new_id = parent_id return new_id + def authorize(self): + # Get credentials + credentials = None + if not USE_SERVICE_ACCOUNTS: + if os.path.exists(G_DRIVE_TOKEN_FILE): + with open(G_DRIVE_TOKEN_FILE, 'rb') as f: + credentials = pickle.load(f) + if credentials is None or not credentials.valid: + if credentials and credentials.expired and credentials.refresh_token: + credentials.refresh(Request()) + else: + flow = InstalledAppFlow.from_client_secrets_file( + 'credentials.json', self.__OAUTH_SCOPE) + LOGGER.info(flow) + credentials = flow.run_console(port=0) + + # Save the credentials for the next run + with open(G_DRIVE_TOKEN_FILE, 'wb') as token: + pickle.dump(credentials, token) + else: + LOGGER.info(f"Authorizing with {SERVICE_ACCOUNT_INDEX}.json service account") + credentials = service_account.Credentials.from_service_account_file( + f'accounts/{SERVICE_ACCOUNT_INDEX}.json', + scopes=self.__OAUTH_SCOPE) + return build('drive', 'v3', credentials=credentials, cache_discovery=False) + def drive_list(self, fileName): msg = "" # Create Search Query for API request. query = f"'{parent_id}' in parents and (name contains '{fileName}')" page_token = None - results = [] + count = 0 while True: - response = service.files().list(supportsTeamDrives=True, - includeTeamDriveItems=True, - q=query, - spaces='drive', - fields='nextPageToken, files(id, name, mimeType, size)', - pageToken=page_token, - orderBy='modifiedTime desc').execute() + response = self.__service.files().list(supportsTeamDrives=True, + includeTeamDriveItems=True, + q=query, + spaces='drive', + fields='nextPageToken, files(id, name, mimeType, size)', + pageToken=page_token, + orderBy='modifiedTime desc').execute() for file in response.get('files', []): - if len(results) >= 20: + if count >= 20: break if file.get( 'mimeType') == "application/vnd.google-apps.folder": # Detect Whether Current Entity is a Folder or File. @@ -369,9 +375,8 @@ class GoogleDriveHelper: url = requests.utils.requote_uri(f'{INDEX_URL}/{file.get("name")}') msg += f' | Index URL' msg += '\n' - results.append(file) + count += 1 page_token = response.get('nextPageToken', None) - if page_token is None: + if page_token is None or count >= 20: break - del results return msg