"""Transfer from browser to S3.""" from datetime import datetime from datetime import timedelta from botocore.errorfactory import ClientError from video import config from video.connectors import s3 from video.constants import job_io_fields from video.logic.activity_task_logger import activity_task_logger from video.models import ows_video @activity_task_logger([], is_first_handler=True, is_last_handler=False) def transfer_from_browser_to_s3(inputs): """Start transfer from browser to s3. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return {} @activity_task_logger( [ job_io_fields.TRANSFER_FROM_BROWSER_TO_S3_JOB_ID, job_io_fields.INPUT_VIDEO_S3_BUCKET, job_io_fields.INPUT_VIDEO_S3_KEY, job_io_fields.S3_TOKEN_EXPIRATION, ], is_first_handler=False, is_last_handler=False, job_type=transfer_from_browser_to_s3.__name__) def is_transfer_from_browser_to_s3_done(inputs): """Check if transfer from browser to s3 is done. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ is_done = False job_id = inputs[job_io_fields.TRANSFER_FROM_BROWSER_TO_S3_JOB_ID] input_s3_bucket = inputs.get(job_io_fields.INPUT_VIDEO_S3_BUCKET) input_s3_key = inputs.get(job_io_fields.INPUT_VIDEO_S3_KEY) if not input_s3_bucket or not input_s3_key: job = ows_video.get_job(job_id) input_s3_bucket = job['inputs'][job_io_fields.INPUT_VIDEO_S3_BUCKET] input_s3_key = job['inputs'][job_io_fields.INPUT_VIDEO_S3_KEY] try: s3.get_s3_client().head_object( Bucket=input_s3_bucket, Key=input_s3_key) is_done = True except ClientError: pass has_transfer_from_browser_to_s3_timed_out = not is_done and ( _has_s3_token_expired(inputs.get(job_io_fields.S3_TOKEN_EXPIRATION))) return { job_io_fields.IS_TRANSFER_FROM_BROWSER_TO_S3_DONE: is_done, job_io_fields.INPUT_VIDEO_S3_BUCKET: input_s3_bucket, job_io_fields.INPUT_VIDEO_S3_KEY: input_s3_key, job_io_fields.HAS_TRANSFER_FROM_BROWSER_TO_S3_TIMED_OUT: has_transfer_from_browser_to_s3_timed_out, } @activity_task_logger( [], is_first_handler=False, is_last_handler=True, job_type=transfer_from_browser_to_s3.__name__) def log_transfer_completed(inputs): """Log successful transfer from browser to S3. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return {} @activity_task_logger( [], is_first_handler=False, is_last_handler=True, job_type=transfer_from_browser_to_s3.__name__) def log_transfer_timed_out(inputs): """Log timeout of transfer from browser to S3. Args: inputs (dict): Inputs. Returns: dict: Outputs. """ return { job_io_fields.ERROR_TRANSFER_FROM_BROWSER_TO_S3_EXPIRED: 'Transfer from browser to S3 timed out.'} def _has_s3_token_expired(expiration_date): """Check if S3 token will expire soon. Args: expiration_date (datetime): S3 token expiration time by UTC. Returns: bool: True if the current time is within the threshold of the timeout, otherwise False. """ if expiration_date: upload_expiration = datetime.strptime( expiration_date, config.S3_TOKEN_EXPIRATION_FORMAT) threshold = timedelta( seconds=config.UPLOAD_FROM_BROWSER_TO_S3_TIMEOUT_THRESHOLD) return (upload_expiration - threshold) < datetime.utcnow() return False