"""Main Handler.""" from bulk_metadata_ingester_common.utils.jitter import jitter from bulk_metadata_ingester_common.utils.logging import get_current_logger from bulk_metadata_ingester_common.utils.release_utils import \ get_all_release_models import config from constants.data_sources import JSON_FULL from ddex_ingester_common.constants.catalog_ingestion import REJECT from ddex_ingester_common.helpers.catalog_ingestion import ( CatalogIngestionValidationResult ) from rules import ( RULES_PER_INGESTION_SOURCE ) class ValidationRuleException(Exception): """ValidationRuleException exception.""" def handler(event, context): """Lambda entry point.""" # TODO remove when Lambda is implemented and working # This key won't be in the event unless we explicitly add it # It skips some Lambdas to allow testing of other parts of the sfn if 'pre_validate' not in event: return event """Lambda Entrypoint.""" key = event.get('key') bucket = event.get('bucket') execution_name = event.get('execution_name') state_machine_name = event.get('state_machine_name') correlation_id = event.get('correlation_id') config.graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id } ) logger = get_current_logger( config.ENVIRONMENT, config.LAMBDA_NAME, logging_level=config.LOGGING_LEVEL, correlation_id=correlation_id) # Jitter calls - Sleep for randomness jitter(logger) # Log lambda begins message logger.info(f'pre_validate received: {event}') # Get releases from JSON doc all_release_models = get_all_release_models(bucket, key) results = None # TODO: Make this conditional based on datafile? rules = RULES_PER_INGESTION_SOURCE.get(JSON_FULL, []) # Execute validation against the given data file validation_results = validate(all_release_models, rules, logger) results = save_validation_result( validation_results, state_machine_name, execution_name ) rejected_rules = [ result.message for result in results if result.response == REJECT ] if rejected_rules: raise ValidationRuleException(rejected_rules) return event def validate(release_track_list, rules, logger): """Validate incoming data against a set of rules.""" validation_results = [] for rule in rules: try: validation_result = rule(release_track_list, logger) if validation_result: validation_results.append(validation_result) except Exception as e: msg = ( 'Exception encountered while executing validation rule: ' f'{rule.__name__}' ) raise RuntimeError(msg) from e return validation_results def save_validation_result( rule_results, state_machine_name, execution_name): """Save validation results to S3.""" results = [] for rule_result in rule_results: results.append( CatalogIngestionValidationResult( state_machine_name=state_machine_name, state_machine_execution_name=execution_name, validation_rule_id=rule_result.validation_rule_id, response=rule_result.response, message=rule_result.message ) ) config.catalog_ingestion_session.add(results) config.catalog_ingestion_session.save() return results