"""Lambda aws_break_glass_audit function module.""" import datetime from typing import Any from lambdacommon.common_config import logger import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration import config from src.utils import athena from src.utils import jira from src.utils import sqs from src.utils import s3 sentry_sdk.init( dsn=config.SENTRY_DSN, integrations=[AwsLambdaIntegration()], traces_sample_rate=1.0, ) def process_message(msg_body: dict[str, Any], receipt_handle: str) -> None: """Process the sqs message. Runs athena query for given user and date from the msg_body, upload results to s3, attach it to Jira ticket and delete sqs message from the queue. :param msg_body: SQS message body :type msg_body: dict[str, Any] :param receipt_handle: SQS message receipt handle :type receipt_handle: str """ athena_instance = athena.AthenaUtil() s3_instance = s3.S3Util() account_id = msg_body['account_id'] # Run Athena steps to generate audit log query_statement = athena_instance.get_named_query_statements(account_id) query_execution_id = athena_instance.start_audit_query_execution( msg_body['user'], msg_body['date'], query_statement) output_file = athena_instance.get_query_execution_details(query_execution_id) # Download output file from S3 output_file_local_path = s3_instance.download_audit_log(output_file) # Verify that jira issue exists jira_issue = msg_body['jira_issue'] issue_url = jira.verify_jira_issue(jira_issue) logger.info(f'Jira issue {jira_issue} found with url {issue_url}') # Attach audit log to Jira issue attachment_timestamp = jira.attach_to_jira_issue( jira_issue, output_file_local_path, msg_body['user']) logger.info(f'Attached audit log to issue {jira_issue} ' f'at time {attachment_timestamp}') sqs_response = sqs.delete_message(receipt_handle) logger.info('SQS Message with receipt handle ' f'{receipt_handle} deleted with ' f'status code {sqs_response}') def handler(event, context): """Lambda entry point.""" try: messages = sqs.get_queue_message_count() if messages > 0: for message in range(messages): receipt_handle, body = sqs.receive_message() if receipt_handle and body: # Check policy expiration date prior to processing policy_creation = datetime.datetime.strptime( body['datetime'], config.DATETIME_STR_FORMAT) policy_expiration = policy_creation + datetime.timedelta(seconds=config.IAM_ROLE_SESSION_DURATION) if datetime.datetime.now() > policy_expiration: process_message(body, receipt_handle) else: now = datetime.datetime.now().strftime(config.DATETIME_STR_FORMAT) expiration_str = policy_expiration.strftime(config.DATETIME_STR_FORMAT) logger.info('Skipping message because break-glass policy has not yet expired. ' f'Current time: {now}. Policy expiration time: {expiration_str}') else: logger.info('Queue had available messages when queried ' 'but there are now no messages in the queue.') else: logger.info('No messages found in queue.') logger.info('Break-Glass audit lambda finished') return {'status': 'OK'} except Exception as error: logger.exception(str(error)) sentry_sdk.capture_exception(error)