"""CodeEngine transform code.""" import copy import json import boto3 import awsconfig ERROR_MESSAGE = 'Failed to send track event to CloudSearch' SUCCESS_LAMBDA_STATUSES = [200, 202, 204] # CloudSearch upload disabled now CLOUDSEARCH_TABLES = ['track'] IGNORE_EVENT_TYPES = ['update-delete'] client = boto3.client( 'lambda', aws_access_key_id=awsconfig.ACCESS_KEY, aws_secret_access_key=awsconfig.ACCESS_SECRET, region_name=awsconfig.AWS_REGION) def send_to_cloudsearch(event): """Send event to CloudSearch by invoking lambda function. Args: event (dict): Alooma event Returns: int: HTTP status code for invocation request """ converted_event = convert_event(event) res = client.invoke( FunctionName=awsconfig.LAMBDA_NAME, InvocationType='Event', Payload=json.dumps(converted_event)) status_code = res.get('StatusCode') return status_code def convert_event(event): """Converts an event to the required output format. - Cleanup unncecessary metadata - Change event type format Args: event (dict): Json event Returns: dict: converted event """ converted_event = {} event_data = copy.deepcopy(event) del event_data['_metadata'] converted_event['table'] = event['_metadata']['table'] converted_event['data'] = event_data converted_event['type'] = convert_event_type(event['_metadata']) return converted_event def convert_event_type(event_metadata): """Convert alooma event_metadata type to a more generic format. Args: event_metadata (dict): alooma event metadata Returns: str: event_type """ source_event_type = event_metadata.get('type', 'insert') if source_event_type == 'insert': converted_event_type = 'insert' elif (source_event_type == 'deleted' and event_metadata.get('deleted')): converted_event_type = 'delete' elif source_event_type in ('update-insert', 'update'): converted_event_type = 'update' return converted_event_type def transform(event): """Filter, transform, drop Alooma events. Args: event (dict): Event, reflecting a change in DB table. Returns: dict: transformed event (or None, for dropped events) """ if (event['_metadata']['event_type'] in CLOUDSEARCH_TABLES and event['_metadata']['type'] not in IGNORE_EVENT_TYPES): status_code = send_to_cloudsearch(event) if status_code not in SUCCESS_LAMBDA_STATUSES: # In case of failure raise an exception to send the event # to restream queue. raise Exception(ERROR_MESSAGE) return None