import base64 import json import os import boto3 import botocore import psycopg2 from botocore.exceptions import ClientError from sme_logger import get_logger from const import APP_NAME, ENV, aws class AwsUtils(object): def __init__(self): self.logger = get_logger( APP_NAME, os.environ.get("ENVIRONMENT", ENV).lower() != ENV) def save_json_data_to_s3(self, *, bucket_name: str, path: str, file_name_to_save: str, data): """ :param bucket_name: Bucket name where data is to be dumped :param path: Prefix of the bucket :param file_name_to_save: File name which is to be saved :param data: Data to be written in the file :return: success response """ s3 = boto3.resource("s3") s3object = s3.Object(bucket_name, f"{path}/{file_name_to_save}") s3object.put(Body=(bytes(json.dumps(data).encode("UTF-8")))) self.logger.info("Data Dumped to Bucket") # Request to create a new client session def create_aws_client(self, *, service_name: str, region_name: str) -> classmethod: """ :param service_name: Name of service for which client needs to be connected :param region_name: AWS region name :return: class object of the client created """ try: # Create a aws service client session = boto3.session.Session() self.logger.info("AWS CLIENT CREATED SUCCESSFULLY") return session.client(service_name=service_name, region_name=region_name) except botocore.exceptions.ClientError as e: self.logger.error( f"Error creating AWS client: {e.response['error']}") def get_s3_files(self, *, bucket: str, prefix: str) -> list: """ :param bucket: Bucket name from which files are to be read :param prefix: Prefix of bucket from where files are to be read :return: list of file names """ try: s3 = self.create_aws_client(service_name="s3", region_name=aws.get( "AWS_DEFAULT_REGION", None)) files = [] """Get a list of keys in an S3 bucket.""" resp = s3.list_objects_v2(Bucket=bucket, Prefix=prefix) for obj in resp["Contents"]: file_name = os.path.basename(obj["Key"]) if file_name != "": files.append(file_name) return files except botocore.exceptions.ClientError as e: self.logger.error( f"Error getting files from {bucket}: {e.response['Error']}") # Function to pass in the data to Firehose delivery stream (Pass client using create_aws_client) def deliver_to_firehose(self, *, client, delivery_stream_name: str, data) -> dict: """ :param client: AWS service client :param delivery_stream_name: Delivery stream name to be used for delivery :param data: data that needs to be sent ot firehose :return: Dictionary of AWS response if success """ try: return client.put_record( DeliveryStreamName=delivery_stream_name, Record={"Data": json.dumps(data).encode()}, ) except botocore.exceptions.ClientError as e: self.logger.error( f"Error delivering data to {delivery_stream_name}: {e.response['Error']}" ) # Pass in the secret name for which you need the secret data def get_secret(self, *, required_secret: str) -> str: """ :param required_secret: Name of the secret value that needs to be fetched :return: Secret value in string format """ secret_name = required_secret client = self.create_aws_client( service_name="secretsmanager", region_name=aws.get("AWS_DEFAULT_REGION", None), ) try: get_secret_value_response = client.get_secret_value( SecretId=secret_name) except ClientError as e: if e.response["Error"]["Code"] in [ "DecryptionFailureException", "InternalServiceErrorException", "InvalidParameterException", "InvalidRequestException", "ResourceNotFoundException", ]: raise e else: if "SecretString" in get_secret_value_response: return get_secret_value_response["SecretString"] else: return base64.b64decode( get_secret_value_response["SecretBinary"]) def generate_rds_token( self, rds_hostname=aws.get("RDS", None).get("DB_HOST"), db_user=aws.get("RDS", None).get("USERNAME"), rds_port=aws.get("RDS", None).get("POSTGRESQL_PORT"), aws_region=aws.get("RDS", None).get("AWS_DEFAULT_REGION"), ) -> str: """ :param rds_hostname: Hostname for RDS instance :param db_user: Database user name :param rds_port: Port used for connection :param aws_region: AWS region name :return: string of response if success """ rds_client = boto3.client("rds") try: return rds_client.generate_db_auth_token( Region=aws_region, DBHostname=rds_hostname, Port=rds_port, DBUsername=db_user, ) except botocore.exceptions.ClientError as e: self.logger.error( f"Error generating RDS token{e.response['error']}") # AWS POSTGRES UTILITY FUNCTIONS def connect_to_database(self) -> classmethod: """ :return: Class object once connection is successful """ user = aws.get("RDS", None).get("USERNAME") password = self.generate_rds_token() try: return psycopg2.connect( database=aws.get("RDS", None).get("DB_NAME"), user=user, password=password, host=aws.get("RDS", None).get("DB_HOST"), port=aws.get("RDS", None).get("POSTGRESQL_PORT"), ) except psycopg2.OperationalError as e: self.logger.error(f"Error connecting to Postgres: {e}") if __name__ == '__main__': AwsUtils()