"""Module to apply event source mapping configurations after a dynamodb table is created.""" import logging from typing import Any logger = logging.getLogger(__name__) class NoDynamodbStreamFailure(Exception): """Indicates the Dynamodb Table needs streams to be enabled.""" pass class NewEventSourceMappingFailure(Exception): """Indicates EventSourceMapping could not be created.""" pass def _get_dynamodb_stream_arn(dynamodb_client: Any, dynamodb_table: str) -> str: """Get the latest dynamodb stream arn.""" logger.info('Fetching latest dynamodb stream arn') try: table_description = dynamodb_client.describe_table(TableName=dynamodb_table) dynamodb_stream_arn = table_description.get('Table', {}).get('LatestStreamArn') if not dynamodb_stream_arn: logger.warning(f'payload missing LatestStreamArn for {dynamodb_table}', extra={ 'describe_table': table_description, }) raise NoDynamodbStreamFailure(f'missing stream arn for {dynamodb_table}') except Exception as e: logger.warning(f'missing stream arn for {dynamodb_table}', exc_info=e) raise NoDynamodbStreamFailure(f'missing stream arn for {dynamodb_table}') from e logger.info(f'Found latest dynamodb stream arn: {dynamodb_stream_arn}') return dynamodb_stream_arn def _get_relevant_event_source_mapping_uuids( lambda_client: Any, lambda_arn: str, dynamodb_table: str, ) -> list[str]: """Find current event source mappings that need to be deleted.""" logger.info('Fetching current event source mappings') relevant_mapping_uuids = [] try: current_mappings = lambda_client.list_event_source_mappings(FunctionName=lambda_arn) # There could be multiple triggers. We only want the ones where # the stream arn is for dynamodb + our dynamodb_table relevant_mapping_uuids = [ event_source_mapping['UUID'] for event_source_mapping in current_mappings.get('EventSourceMappings', []) if 'arn:aws:dynamodb' in event_source_mapping[ 'EventSourceArn'] and dynamodb_table in event_source_mapping[ 'EventSourceArn'] ] except Exception as e: logger.warning(f'Could not obtain current event source mappings for {lambda_arn}') logger.warning(e) logger.info(f'Found relevant event source mappings {relevant_mapping_uuids}') return relevant_mapping_uuids def _delete_event_source_mappings( lambda_client: Any, uuids: list[str], ) -> None: """Delete event source mappings.""" logger.info(f'Deleting old event source mappings: {uuids}') for delete_uuid in uuids: try: lambda_client.delete_event_source_mapping(UUID=delete_uuid) except Exception as e: logger.warning(f'Failed to delete an event source mapping: {delete_uuid}') logger.warning(e) def _create_event_source_mapping( lambda_client: Any, target_lambda_arn: str, dynamodb_stream_arn: str, **kwargs: Any, ) -> None: """Create new event source mapping.""" logger.info(f'Creating new event source mapping for {dynamodb_stream_arn}') try: event_source_mapping = lambda_client.create_event_source_mapping( EventSourceArn=dynamodb_stream_arn, FunctionName=target_lambda_arn, Enabled=True, **kwargs, ) except Exception as e: raise NewEventSourceMappingFailure( f'could not create event source mapping between {dynamodb_stream_arn} and {target_lambda_arn}') from e logger.info(f'Created event source mapping {event_source_mapping.get("UUID")}') def set_lambda_event_source_mapping( dynamodb_client: Any, dynamodb_table: str, lambda_client: Any, target_lambda_arn: str, **kwargs: Any, ) -> None: """Upsert the lambda event source mapping. Any additional kwargs are passed as keyword arguments for creating the event source mapping * Determines the dynamodb table's stream * Fetches the current event source mappings on the target_lambda fn * Creates a new event source mapping on the target_lambda fn, triggered by the dynamodb table's stream * Deletes the existing event source mapping on the target lambda fn """ logger.info('Setting lambda event source mapping') try: dynamodb_stream_arn = _get_dynamodb_stream_arn(dynamodb_client, dynamodb_table) except NoDynamodbStreamFailure: logger.error(f'Check stream enable configuration for {dynamodb_table} and start task again') return relevant_mapping_uuids = _get_relevant_event_source_mapping_uuids(lambda_client, target_lambda_arn, dynamodb_table) _create_event_source_mapping(lambda_client, target_lambda_arn, dynamodb_stream_arn, **kwargs) _delete_event_source_mappings(lambda_client, relevant_mapping_uuids) logger.info('Done setting lambda event source mapping')