"""Transfer from Google Drive to S3.""" from contextlib import contextmanager from itertools import chain import json import threading import requests from video import config from video.connectors import s3 from video.constants import assets from video.constants import job_io_fields from video.logic.activity_task_logger import activity_task_logger from video.models import ows_video from video.utils import exception as exception_utils class FileTransferProgressLogger(object): """Log S3 Transfer Progress.""" def __init__( self, job_id=None, file_size_bytes=None, input_video_s3_key=None): """Init FileTransferProgressLogger.""" self._job_id = job_id self._input_video_s3_key = input_video_s3_key self._progress_bytes = 0 self._progress_update_lock = threading.Lock() self._logging_lock = threading.Lock() self._progress_percent = 0 self._logged_total_bytes = file_size_bytes self._logged_progress_percent = 0 self._logged_progress_bytes = 0 # Example cutoff (33, 2) means: Starting at 33 % log every 2 %. self._cutoffs = [ (0, 1), (2, 2), (6, 3), (12, 5), (22, 8), (38, 13), (63, 8), (79, 5), (89, 3), (95, 2), (98, 1), (101, 1), ] self._logged_progress_percent_to_progress_logging_interval_percent = { progress: interval for progress, interval in chain(*[ [ (progress, interval) for progress in range(cutoff_start, cutoff_end) ] for cutoff_start, interval, cutoff_end in [ (*cs, ce[0]) for cs, ce in zip(self._cutoffs, self._cutoffs[1:]) ] ]) } ows_video.set_job_fields({ 'id': self._job_id, 'outputs': { job_io_fields.FILE_TRANSFER_PROGRESS_PERCENT: self._logged_progress_percent, job_io_fields.FILE_TRANSFER_PROGRESS_BYTES: self._logged_progress_bytes, job_io_fields.FILE_TRANSFER_TOTAL_BYTES: self._logged_total_bytes }, }) self._log_to_s3() @contextmanager def logging_lock(self): """Non-Blocking Lock.""" is_locked = self._logging_lock.acquire(blocking=False) yield is_locked if is_locked: self._logging_lock.release() def _log_to_s3(self): s3.get_s3_client().put_object( Body=json.dumps({ job_io_fields.FILE_TRANSFER_PROGRESS_PERCENT: self._progress_percent, job_io_fields.FILE_TRANSFER_PROGRESS_BYTES: self._progress_bytes, job_io_fields.FILE_TRANSFER_TOTAL_BYTES: self._logged_total_bytes }).encode(), Bucket=assets.VIDEO_S3_BUCKET, Key='{}.progress.json'.format(self._input_video_s3_key), ) def _should_log_to_ows_video(self): return ( self._progress_percent == 100 or (self._progress_percent - self._logged_progress_percent > self._logged_progress_percent_to_progress_logging_interval_percent[self._progress_percent]) # noqa ) def __call__(self, bytes_amount): """Call FileTransferProgressLogger.""" with self._progress_update_lock: if not bytes_amount: return self._progress_bytes += bytes_amount self._progress_percent = int( (self._progress_bytes / self._logged_total_bytes) * 100) with self.logging_lock() as is_locked: if not is_locked: return if self._should_log_to_ows_video(): self._logged_progress_percent = self._progress_percent self._logged_progress_bytes = self._progress_bytes ows_video.set_job_fields({ 'id': self._job_id, 'outputs': { job_io_fields.FILE_TRANSFER_PROGRESS_PERCENT: self._logged_progress_percent, job_io_fields.FILE_TRANSFER_PROGRESS_BYTES: self._logged_progress_bytes }, }) self._log_to_s3() def is_done(self): """Indicate if transfer is done.""" return self._logged_total_bytes == self._logged_progress_bytes def _raise_exception_if_transfer_not_done( file_transfer_progress_logger, ): """Raise exception if transfer not done.""" if not file_transfer_progress_logger.is_done(): exception_utils.raise_exception( Exception, json.dumps({ 'error': 'File transfer stopped unexpectedly before completing.', })) @activity_task_logger([ job_io_fields.TRANSFER_FROM_GOOGLE_DRIVE_TO_S3_JOB_ID, job_io_fields.GOOGLE_DRIVE_FILES, job_io_fields.GOOGLE_DRIVE_AUTHORIZATION, job_io_fields.WORKFLOW_JOB_ID, ]) def transfer_from_google_drive_to_s3(inputs): """Transfer from Google Drive to S3. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ file, = inputs[job_io_fields.GOOGLE_DRIVE_FILES] authorization = inputs[job_io_fields.GOOGLE_DRIVE_AUTHORIZATION] workflow_job_id = inputs[job_io_fields.WORKFLOW_JOB_ID] job_id = inputs[job_io_fields.TRANSFER_FROM_GOOGLE_DRIVE_TO_S3_JOB_ID] input_video_s3_key = ( assets.RAW_S3_KEY_TEMPLATE.format( workflow_job_id=workflow_job_id, file_name=file['name'], )) request_params = { 'application': config.SERVICE_NAME, 'environment': config.OWS_ENVIRONMENT, 'method': 'GET', 'url': 'https://www.googleapis.com/drive/v3/files/{}'.format(file['id']), 'headers': { 'Authorization': 'Bearer {}'.format(authorization['access_token'])}, 'query_params': { 'key': config.GOOGLE_API_KEY, 'alt': 'media'}, 'stream': True } response = requests.request( request_params['method'], request_params['url'], headers=request_params['headers'], params=request_params['query_params'], stream=request_params['stream'], ) exception_utils.raise_exception_if_request_failed( request_params={ **request_params, 'query_params': { **request_params['query_params'], 'key': 'REDACTED', } }, response=response, ) file_transfer_progress_logger = FileTransferProgressLogger( job_id=job_id, file_size_bytes=file['size_bytes'], input_video_s3_key=input_video_s3_key, ) s3_client = s3.get_s3_client() s3_client.upload_fileobj( response.raw, assets.VIDEO_S3_BUCKET, input_video_s3_key, Callback=file_transfer_progress_logger, ) _raise_exception_if_transfer_not_done( file_transfer_progress_logger ) return { job_io_fields.INPUT_VIDEO_S3_BUCKET: assets.VIDEO_S3_BUCKET, job_io_fields.INPUT_VIDEO_S3_KEY: input_video_s3_key, } @activity_task_logger([ job_io_fields.TRANSFER_FROM_DROPBOX_TO_S3_JOB_ID, job_io_fields.DROPBOX_FILES, job_io_fields.WORKFLOW_JOB_ID, ]) def transfer_from_dropbox_to_s3(inputs): """Transfer from Dropbox to S3. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ file, = inputs[job_io_fields.DROPBOX_FILES] workflow_job_id = inputs[job_io_fields.WORKFLOW_JOB_ID] job_id = inputs[job_io_fields.TRANSFER_FROM_DROPBOX_TO_S3_JOB_ID] input_video_s3_key = ( assets.RAW_S3_KEY_TEMPLATE.format( workflow_job_id=workflow_job_id, file_name=file['name'], )) response = requests.request( 'GET', file['link'], stream=True, ) s3_client = s3.get_s3_client() s3_client.upload_fileobj( response.raw, assets.VIDEO_S3_BUCKET, input_video_s3_key, Callback=FileTransferProgressLogger( job_id=job_id, file_size_bytes=file['bytes'], input_video_s3_key=input_video_s3_key, ), ) return { job_io_fields.INPUT_VIDEO_S3_BUCKET: assets.VIDEO_S3_BUCKET, job_io_fields.INPUT_VIDEO_S3_KEY: input_video_s3_key, }