import base64 import json import logging from typing import Dict, Any import sys import boto3 from botocore.exceptions import ClientError logger = logging.getLogger() logger.setLevel(logging.INFO) work_dir = '/tmp' def get_secret(secret_name: str): """Get secret from the AWS Secrets Manager Args: secret_name: name of a secret record aws_params: AWS connection parameters Returns: Secret value """ # Create a Secrets Manager client # Creating Session session = boto3.session.Session() # Defining client client = session.client(service_name='secretsmanager') secret: str = "" try: # Getting response from client get_secret_value_response = client.get_secret_value(SecretId=secret_name) logger.debug('Retreving Secret: %s', secret_name) except ClientError as c_error: # Get error messages if they arise logger.error(c_error) else: # Decrypts secret using the associated KMS CMK. # Depending on whether the secret is a string or binary, one of these # fields will be populated. if 'SecretString' in get_secret_value_response: secret = get_secret_value_response['SecretString'] # logs success logger.debug("Retrieval of Secret: [ %s ]" "[string] Successful!", secret_name) else: secret = base64.b64decode( get_secret_value_response['SecretBinary']) # logs success logger.debug("Retrieval of Secret: [ %s ]" "[binary] Successful!", secret_name) # Returns the Secret return secret def get_role_credentials(secret: Dict[str, Any]) -> Dict[str, any]: sts_client = boto3.client('sts') assumed_role_object = sts_client.assume_role( RoleArn=secret["ACCESS_ROLE"], RoleSessionName="AssumeRoleSession1" ) return assumed_role_object['Credentials'] def download(): response = source_client.list_objects_v2( Bucket=src_bucket, Prefix=file_prefix ) downloaded = set() if response['Contents']: for item in response['Contents']: key = item['Key'] dest_filename = key.replace(file_prefix, replacement).replace(".xlsx", '-01-01-01.xlsx') dest_filepath = work_dir + '/' + dest_filename if key.endswith(".xlsx"): logger.info('Downloading {} to {}'.format(key, dest_filepath)) source_client.download_file(src_bucket, key, dest_filepath) downloaded.add(dest_filename) # break else: logger.info('Skipping upsupported file type {}'.format(key)) else: logger.error("No files like {}/{}".format(src_bucket, file_prefix)) return downloaded def upload(downloaded): processed_dir = 'processed' processed_resp = dest_client.list_objects_v2( Bucket=dest_bucket, Prefix='{}/{}2020'.format(processed_dir, replacement) ) processed = set() if processed_resp['Contents']: for item in processed_resp['Contents']: key = item['Key'] processed.add(key.replace('{}/'.format(processed_dir), '')) new_files = downloaded - processed logger.info('New files %s', new_files) if new_files: for new_file in new_files: file_path = work_dir + '/' + new_file s3_key = 'incoming/' + new_file logger.info('Uploading %s to s3://%s/%s', file_path, dest_bucket, s3_key) dest_client.upload_file(file_path, dest_bucket, s3_key) else: logger.warning("No new files detected.") return new_files file_prefix = 'output/iTunesSpotifyDeezerProdCtlgPrevious14Days' replacement = 'Filtr_Admin_New_Releases_' s3_creds = json.loads(get_secret('delphi/dev/slz/sme_max')) src_bucket = s3_creds['BUCKET'] credentials = get_role_credentials(s3_creds) source_client = boto3.client( 's3', aws_access_key_id=credentials['AccessKeyId'], aws_secret_access_key=credentials['SecretAccessKey'], aws_session_token=credentials['SessionToken'], ) dest_bucket = 'filtr-new-releases-upc' dest_client = boto3.client('s3') def lambda_handler(event, context): try: downloaded = download() logger.info('Downloaded: %s', downloaded) uploaded = upload(downloaded) return 'Uploaded files {}'.format(uploaded) except Exception as e: logger.error(e) raise e