"""Lambda upload_audit_logs_to_s3 function module.""" import base64 import datetime import decimal import json import os import random import string from aws_kinesis_agg import deaggregator import boto3 import config from src.constants import general from lambdacommon.common_config import logger as app_logger import pandas as pd import pyarrow as pa import pyarrow.parquet as pq import pytz eastern = pytz.timezone('America/New_York') s3_client = boto3.client('s3') release_approval_logging_tables = ['release_logging', 'release_status', 'release_approval_queue', 'rejection_notes', 'release_correction', 'release_approval_comments'] timestamp_tables = ['manual_adjustment', 'release_manual_adjustment', 'release_logging', 'release_status', 'vendor_restricted_features', 'orchadmin_users', 'orchadmin_user_roles'] timestamp_columns = ['updated_timestamp', 'last_updated'] table_date_column = { 'check_receivable': ['check_date'], 'owner': ['agreement_start_date', 'agreement_end_date'], 'publishers_checkspaid': ['cut_date', 'cash_date'], 'vendor_chargeback': ['chargeback_date', 'label_approval_date', 'supervisor_approval_date'], 'release_status': ['date'] } table_datetime_column = { 'check_receivable': ['entry_date'], 'checkspaid': ['entry_date', 'cut_date', 'cash_date'], 'publishers_checkspaid': ['entry_date'], 'contact': ['address_last_updated'], 'manual_adjustment': ['date_added', 'updated_timestamp'], 'release_manual_adjustment': ['date_created', 'last_updated'], 'travelex': ['last_updated'], 'vendor_checkspaid': ['entry_date', 'cut_date', 'cash_date'], 'release_status': ['updated_timestamp'], 'release_logging': ['updated_timestamp'], 'release_approval_queue': ['last_updated', 'date_submitted'], 'release_correction': ['last_updated'], 'rejection_notes': ['date_added'], 'release_approval_comments': ['datetime'], 'vendor_restricted_features': ['updated_timestamp'], 'project': ['created_date_utc', 'updated_date_utc'], 'orchadmin_users': ['password_date_changed', 'password_reset_datetime', 'updated_timestamp'], 'orchadmin_user_roles': ['updated_timestamp'] } table_decimal_column = { 'checkspaid': ['check_amt', 'amount_in_original_currency'], 'dig_sales': ['activity_rate'], 'manual_adjustment': ['amount', 'amount_in_original_currency'], 'publishers_checkspaid': ['check_amt'], 'release_manual_adjustment': ['amount'], 'vendor_chargeback': ['amount'], 'vendor_checkspaid': ['check_amt'] } table_integer_column = { 'release_status': ['upc_REMOVED'], 'release_approval_queue': ['release_correction_id', 'checked_out_by', 'approved_by', 'submitted_by', 'changed_by'], 'release_approval_comments': ['release_status_id', 'release_approval_id', 'changed_by'], 'rejection_notes': ['key_id', 'changed_by'], 'vendor_restricted_features': ['vendor_restricted_features_id', 'vendor_id', 'feature_id', 'preview', 'last_modified_by'], 'project': ['project_id', 'vendor_id', 'subaccount_id', 'artist_id', 'last_modified_by'], 'orchadmin_users': ['role', 'extension', 'last_modified_by'], 'orchadmin_user_roles': ['orchadmin_user_id', 'orchadmin_role_id', 'last_modified_by'] } neo4j_tables = [ 'Profile', 'Identity', 'HAS_ACCESS_TO', ] def handler(event, context): """Lambda entry point.""" tables_to_log = [ 'check_receivable', 'check_receivable_detail', 'checkspaid', 'contact', 'dig_sales', 'manual_adjustment', 'owner_contact', 'owner_payment_level', 'owner', 'publishers_checkspaid', 'release_manual_adjustment', 'travelex', 'vendor_chargeback', 'vendor_checkspaid', 'vendor_restricted_features', 'project', 'orchadmin_users', 'orchadmin_user_roles'] tables_to_log.extend(release_approval_logging_tables) tables_to_log.extend(neo4j_tables) table_records = {table: [] for table in tables_to_log} app_logger.debug(event) kinesis_records = deaggregator.iter_deaggregate_records(event['Records']) for raw_record in kinesis_records: record_str = base64.b64decode( raw_record['kinesis']['data']).decode('utf-8') record_json = json.loads(record_str) app_logger.debug('Json data from kinesis {}'.format(record_json)) table_name = record_json.get('table') shouldProjectTableLog = project_table_log(record_json) \ if table_name == 'project' else True if table_name in tables_to_log and \ table_should_log(table_name, record_json['data']) and \ shouldProjectTableLog: record_json['data']['database'] = record_json['database'] record_json['data']['ts'] = record_json['ts'] record_json['data']['log_action'] = record_json['type'] dt = datetime.datetime.fromtimestamp( record_json['data']['ts'], eastern) record_json['data']['log_timestamp'] = \ pd.Timestamp(dt.strftime('%Y-%m-%d %H:%M:%S')) if 'last_modified_by' in record_json['data'] and \ record_json['data']['last_modified_by'] is None: record_json['data']['last_modified_by'] = 0 record_json['data'] = sanitize_column_values_for_athena( record_json['data'], table_name) if (table_name == 'owner' and record_json['data']['aggregator_owner_id'] is None): record_json['data']['aggregator_owner_id'] = 0 if (table_name == 'contact' and record_json['data']['orchard_country'] is None): record_json['data']['orchard_country'] = 0 table_records[table_name].append(record_json['data']) else: app_logger.debug( 'We dont process records for table {}'.format(table_name)) for table_name, records in table_records.items(): if records: app_logger.info( '{} records found for table {}'.format( len(records), table_name)) database = records[0]['database'] for record in records: del record['database'] try: app_logger.info('Creating parquet file...') filename = create_parquet_file(records) app_logger.info('Uploading parquet file to S3...') upload_to_s3(filename, database, table_name) temp_file_path = config.TEMP_FILE_PATH.format(filename) os.remove(temp_file_path) app_logger.info( 'Saved {} successfully to S3'.format(temp_file_path)) except Exception as e: error_msg = 'Error when trying to upload to s3 - {}'.format( str(e)) app_logger.error(error_msg) raise else: app_logger.debug('No records for table {}'.format(table_name)) app_logger.info('Completed process...') def sanitize_datetime_for_athena(record_json, table_name): """Sanitize datetime column value for athena. Args: record_json(dict): Kinesis JSON entry for a table record. table_name(string): Name of the table. """ if table_name in table_datetime_column: for table_datetime_col in table_datetime_column[table_name]: try: if (table_name in timestamp_tables and table_datetime_col in timestamp_columns): dt = datetime.datetime.fromtimestamp( record_json['ts'], eastern) record_json[table_datetime_col] = \ pd.Timestamp( dt.strftime('%Y-%m-%d %H:%M:%S')) else: record_json[table_datetime_col] = pd.Timestamp( record_json[table_datetime_col]) except ValueError: record_json[table_datetime_col] = None def sanitize_date_for_athena(record_json, table_name): """Sanitize date column value for athena. Args: record_json(dict): Kinesis JSON entry for a table record. table_name(string): Name of the table. """ if table_name in table_date_column: for table_date_col in table_date_column[table_name]: try: date_format = datetime.datetime.strptime( record_json[table_date_col], '%Y-%m-%d') record_json[table_date_col] = datetime.date( date_format.year, date_format.month, date_format.day) except Exception: record_json[table_date_col] = None def sanitize_decimals_for_athena(record_json, table_name): """Sanitize decimal column value for athena. Args: record_json(dict): Kinesis JSON entry for a table record. table_name(string): Name of the table. """ if table_name in table_decimal_column: for table_decimal_col in table_decimal_column[table_name]: if record_json[table_decimal_col]: stringed_decimal = str.format( '{0:.' + str(config.DECIMAL_DIGITS) + 'f}', decimal.Decimal(record_json[table_decimal_col])) record_json[table_decimal_col] = decimal.Decimal( stringed_decimal) else: record_json[table_decimal_col] = None def sanitize_column_values_for_athena(record_json, table_name): """Sanitize column values for athena. Args: record_json(dict): Kinesis JSON entry for a table record. table_name(string): Name of the table. """ sanitize_datetime_for_athena(record_json, table_name) sanitize_date_for_athena(record_json, table_name) sanitize_decimals_for_athena(record_json, table_name) sanitize_integers_for_athena(record_json, table_name) return record_json def upload_to_s3(filename, database, table_name): """Upload parquet file to s3. Args: filename(string): Randomly generated filename of parquet file. database(string): Database as part of S3 filepath of parquet file. table_name(string): Table name as part of S3 filepath of parquet file. """ s3_file_path = '{}/{}/{}'.format(database, table_name, filename) s3_client.upload_file( config.TEMP_FILE_PATH.format(filename), config.AWS_S3_BUCKET_NAME, s3_file_path) def generate_file_name(): """Generate filename for parquet file.""" chars = string.ascii_uppercase + string.digits random_str = ''.join( random.choice(chars) for _ in range(general.PARQUET_FILENAME_LENGTH)) return '{}_{}.parquet'.format( datetime.datetime.now().strftime('%Y_%m_%d_%H:%I:%S'), random_str) def create_parquet_file(table_records): """Create parquet file from Kinesis stream JSON data. Args: table_records(List): List of dict having record details from Kinesis. """ table_dframe = pd.DataFrame( table_records, index=list(range(1, len(table_records) + 1))) if 'currency_id' in table_dframe.columns: table_dframe['currency_id'] = \ table_dframe['currency_id'].fillna(0).astype(int) app_logger.info('Generated table data frame.') table = pa.Table.from_pandas(table_dframe) app_logger.info('Generated table.') file_name = generate_file_name() app_logger.info('Generated filename.') writer = pq.ParquetWriter( config.TEMP_FILE_PATH.format(file_name), table.schema, use_deprecated_int96_timestamps=True) if not writer: raise Exception('Error when writing to parquet file') writer.write_table(table=table) app_logger.info('Completed writing parquet file.') writer.close() return file_name def table_should_log(table_name, raw_record): """Check table should get log or not. Args: table_name(string): Name of the table. record_json(dict): Kinesis JSON entry for a table record. """ if table_name not in release_approval_logging_tables: return True changed_by_type = '' if 'changed_by_type' in raw_record: changed_by_type = raw_record['changed_by_type'] elif 'last_updated_type' in raw_record: changed_by_type = raw_record['last_updated_type'] return changed_by_type == 'oa' def project_table_log(record_json): """Check project table should only log when artist will be updated. Args: record_json(dict): Kinesis JSON entry for table records. """ if 'old' in record_json and 'artist_id' in record_json['old']: return True def sanitize_integers_for_athena(record_json, table_name): """Sanitize integer column value for athena. Args: record_json(dict): Kinesis JSON entry for a table record. table_name(string): Name of the table. """ if table_name in table_integer_column: for table_integer_col in table_integer_column[table_name]: if table_integer_col in record_json \ and record_json[table_integer_col] is None: record_json[table_integer_col] = 0