"""Utility functions for Bulk upload of pricing data for Physical Releases."""
import json
import os
import boto3
from boto3.session import Session
def send_email_with_ses(sender, subject, message, recipients):
"""Send an email using AWS SES.
Args:
sender (str): sender of the email
subject (str): subject of the email
message (str): body of the email
recipients (list): list of recipients of the email
"""
client = boto3.client(
'ses', region_name='us-east-1',
aws_access_key_id=os.environ.get('AWS_ACCESS_KEY_ID'),
aws_secret_access_key=os.environ.get('AWS_SECRET_ACCESS_KEY'))
client.send_email(
Source=sender,
Destination={
'ToAddresses': recipients,
},
Message={
'Subject': {
'Data': subject,
'Charset': 'utf-8'
},
'Body': {
'Html': {
'Data': message,
'Charset': 'utf-8'
}
}
},
ReplyToAddresses=[
sender
]
)
def send_email(subject, mail_body, receivers, release_ids, logger):
"""Send email to notify user.
Args:
subject (str): subject of the email
mail_body (str): body of the email
receivers (list): list of recipients of the email
logger (Logger object): Error logger object
"""
try:
sender = os.environ.get(
'SENDER_EMAIL', 'donotreply@theorchard.com')
messagebody = """
{}
Cheers,
Orchard Physical Pricing Robot
""".format(mail_body)
send_email_with_ses(sender, subject, messagebody, receivers)
print('Email sent successfully to {}'.format(receivers[0]))
logger.info('Email sent successfully to {} with message {} '.format(receivers[0], release_ids))
print('Email sent successfully to {} with message {} '.format(receivers[0], mail_body))
except Exception as e:
logger.info('Unable to send email {}'.format(str(e)))
print('Error: unable to send email',str(e))
def getS3connection():
"""Connect to AWS."""
access_key = os.environ.get('AWS_ACCESS_KEY_ID')
secret_key = os.environ.get('AWS_SECRET_ACCESS_KEY')
sessionconn = Session(
aws_access_key_id=access_key, aws_secret_access_key=secret_key)
session_res = sessionconn.resource('s3')
return session_res
def get_pricing_data(batch_file):
"""Read S3 JSON batch file.
Args:
batch_file (Object): Batch file to read
"""
return json.loads(batch_file.get()['Body'].read().decode('utf-8'))
def rename_ingest_batch(bucket, session_res, old_name, new_name):
"""Rename S3 JSON batch.
Args:
bucket (Object): Source bucket object
session_res (Object): S3 connection resource
old_name (str): Old filename of file to be renamed
new_name (str): New name of file to be renamed
"""
bucket_name = bucket.name
copy_source = {'Bucket': bucket_name, 'Key': old_name}
session_res.meta.client.copy(copy_source, bucket_name, new_name)
session_res.Object(bucket_name, old_name).delete()
def archive_batch(bucket, session_res, batch_filename, new_filename):
"""Archive S3 JSON batch.
Args:
bucket (Object): Source bucket object
session_res (Object): S3 connection resource
batch_filename (str): Old filename of file to be archived
new_filename (str): New name of file to be archived
"""
print('Archiving ', batch_filename, ' to ', new_filename)
rename_ingest_batch(bucket, session_res, batch_filename, new_filename)
def get_oldest_batch_to_ingest(bucket, ingest_folder):
"""Get oldest S3 JSON batch to ingest.
Args:
bucket (Object): Source bucket object
ingest_folder (str): Folder name for ingestion batch file.
"""
allfiles = bucket.objects.filter(Prefix='{}/ready'.format(ingest_folder))
allreadyfiles = [f for f in allfiles]
if not allreadyfiles:
return 0
return min(allreadyfiles, key=lambda k: k.last_modified)