Refactor service account implementation
Signed-off-by: lzzy12 <jhashivam2020@gmail.com>
This commit is contained in:
parent
7ad7d9388a
commit
2bbfca0496
|
|
@ -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'<a href="{self.__G_DRIVE_DIR_BASE_DOWNLOAD_URL.format(dir_id)}">{meta.get("name")}</a> ({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'<a href="{self.__G_DRIVE_BASE_DOWNLOAD_URL.format(file.get("id"))}">{meta.get("name")}</a> ({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' | <a href="{url}"> Index URL</a>'
|
||||
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
|
||||
|
|
|
|||
Loading…
Reference in New Issue