"""Lambda audit function module.""" import base64 import json import aws_kinesis_agg.aggregator import boto3 from lambdacommon.common_config import logger from neoparse import NeoMaxwells from neoparse import ParserError import config kinesis_client = boto3.client('kinesis', region_name=config.AWS_REGION) kinesis_agg = aws_kinesis_agg.aggregator.RecordAggregator() neo_parser = NeoMaxwells(config.NEO4J_DATABASE_NAME) def send_record(agg_record): """Send an aggregated Kinesis record to the stream. Args: agg_record: Kinesis aggregated record """ if agg_record: global kinesis_client pk, ehk, data = agg_record.get_contents() kinesis_client.put_record( StreamName=config.KINESIS_STREAM_NAME, Data=data, PartitionKey=pk) # Add kinesis aggregator callback to flush records to the stream. kinesis_agg.on_record_complete(send_record) def handler(event, context): """Lambda entry point.""" try: for partition, records in event['records'].items(): logger.info(f'Received {len(records)} from MSK.') for record in records: record_str = base64.b64decode(record['value']) record_json = json.loads(record_str.decode()) try: maxwells_record = neo_parser.parse(record_json) except ParserError: logger.exception( f'Failed to parse neo4j record. Record: {record_str}.') continue pk = maxwells_record['table'] kinesis_agg.add_user_record( pk, json.dumps(maxwells_record).encode()) logger.info('Records processed.') send_record(kinesis_agg.clear_and_get()) return {'status': 'OK'} except Exception as e: logger.exception('Records processing failed.') raise e