"""Logic for amazon_datapulse_trigger.""" import json import config from datadog_lambda.metric import lambda_metric from src import jenkins_utils logger = config.logger def process_message(body, message_id): """Process an SQS/SNS message and trigger the corresponding Jenkins job.""" if body.get('Type') != 'Notification': logger.warning( f'Skipping non-notification message (message_id={message_id}): ' f'{json.dumps(body)}' ) return message = json.loads(body['Message']) report_type = message['table'] if report_type not in config.REPORTS: raise ValueError(f'Unknown report type: {report_type}') report_config = config.REPORTS[report_type] jenkins_job = report_config['jenkins_job'] partition = message['partition'] for environment in config.JENKINS_ENVIRONMENTS: jenkins_args = _prepare_jenkins_args(environment, partition, report_type) jenkins_utils.trigger_job(jenkins_job, jenkins_args) lambda_metric( 'amazon_datapulse_trigger.routed', 1, tags=[ f'report_type:{report_type}', f'environment:{environment}', ], ) logger.info( f'Triggered Jenkins job for {report_type} on {environment} ' f'(message_id={message_id})' ) def _prepare_jenkins_args(environment: str, partition, report_type) -> dict[str, str]: report_config = config.REPORTS[report_type] args_templates = { **partition, # include all partition keys as potential Jenkins args 'report_type': report_type, 'environment': environment, 'partition': json.dumps(partition), } jenkins_args = { k: v.format(**args_templates) for k, v in report_config['jenkins_args'].items() } return jenkins_args