From 079beb6872132fe5edcfc03faf770e9f659f0c9c Mon Sep 17 00:00:00 2001 From: rohit kumar Date: Wed, 26 Aug 2026 10:51:52 +0000 Subject: [PATCH 1/5] Log per-input deletion validation metrics --- .../executor/app/executor/import_executor.py | 45 +++++++++++--- .../executor/test/import_executor_test.py | 62 ++++++++++++++++++- 2 files changed, 96 insertions(+), 11 deletions(-) diff --git a/import-automation/executor/app/executor/import_executor.py b/import-automation/executor/app/executor/import_executor.py index 47ec5f79e1..4f2a314aa4 100644 --- a/import-automation/executor/app/executor/import_executor.py +++ b/import-automation/executor/app/executor/import_executor.py @@ -666,6 +666,11 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, if not differ_status: differ_status = diff_found timer = Timer() + input_validation_status = False + validation_metrics = { + "stage": import_stage.name, + "import_input": import_prefix, + } try: config_file_path = self._get_validation_config_file( repo_dir, absolute_import_dir, import_spec, @@ -680,23 +685,43 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, lint_report=report_json, validation_output=validation_output_file) overall_status, current_results = validation.run_validations() + input_validation_status = overall_status validation_results.extend(current_results) if validation_status: validation_status = overall_status + + validator_by_rule_id = { + rule['rule_id']: rule['validator'] + for rule in validation.config.rules + if rule.get('enabled', True) + } + for result in current_results: + if (validator_by_rule_id.get( + result.name) == 'DELETED_RECORDS_PERCENT'): + validation_metrics.update({ + 'deleted_records_percent': + result.details.get('percent'), + 'deleted_records_count': + result.details.get('deleted_records_count'), + 'previous_obs_count': + result.details.get('previous_obs_count'), + }) + break except ValueError as e: logging.error('ValidationRunner failed: %s', e) validation_status = False + validation_metrics.update({ + "latency": + timer.time(), + "status": + ImportStatus.SUCCESS.name + if input_validation_status else ImportStatus.FAILURE.name, + }) log_metric( - AUTO_IMPORT_JOB_STAGE, "INFO" if validation_status else "ERROR", - f"Import: {import_name}, validation: {validation_status}", { - "stage": - import_stage.name, - "latency": - timer.time(), - "status": - ImportStatus.SUCCESS.name - if validation_status else ImportStatus.FAILURE.name, - }) + AUTO_IMPORT_JOB_STAGE, + "INFO" if input_validation_status else "ERROR", + f"Import: {import_name}, input: {import_prefix}, validation: {input_validation_status}", + validation_metrics) if os.path.exists(validation_output_path): # Upload output to GCS. diff --git a/import-automation/executor/test/import_executor_test.py b/import-automation/executor/test/import_executor_test.py index 5a863f00e0..763b8a929a 100644 --- a/import-automation/executor/test/import_executor_test.py +++ b/import-automation/executor/test/import_executor_test.py @@ -22,7 +22,8 @@ import threading from app.executor import import_executor -from app.executor.import_executor import ImportStatus +from app.executor.import_executor import ImportStatus, ImportStatusSummary +from tools.import_validation.result import ValidationResult, ValidationStatus class ImportExecutorTest(unittest.TestCase): @@ -120,3 +121,62 @@ def test_construct_process_message_no_output(self): '[Subprocess command]: exit 0\n' '[Subprocess return code]: 0') self.assertEqual(expected, message) + + @mock.patch.object(import_executor, 'log_import_status') + @mock.patch.object(import_executor, 'log_metric') + @mock.patch.object(import_executor, 'ValidationRunner') + def test_validation_metrics_include_deleted_percent_for_each_input( + self, mock_validation_runner, mock_log_metric, _): + validation_runners = [] + for input_index, percent in enumerate((10.0, 20.0)): + rule_id = f'deleted_percent_{input_index}' + runner = mock.Mock() + runner.config.rules = [{ + 'rule_id': rule_id, + 'validator': 'DELETED_RECORDS_PERCENT', + }] + runner.run_validations.return_value = (False, [ + ValidationResult(ValidationStatus.FAILED, + rule_id, + details={ + 'percent': percent, + 'deleted_records_count': input_index + 1, + 'previous_obs_count': 10, + }) + ]) + validation_runners.append(runner) + mock_validation_runner.side_effect = validation_runners + + config = mock.Mock(invoke_differ_tool=True, + ignore_validation_status=False, + enable_skip_status=True) + executor = import_executor.ImportExecutor(mock.Mock(), mock.Mock(), + config) + executor._get_latest_version = mock.Mock(return_value='previous') + executor._invoke_differ_summary = mock.Mock(return_value={ + 'obs_diff_count': 1, + 'schema_diff_count': 0, + }) + executor._get_validation_config_file = mock.Mock( + return_value='validation_config.json') + executor._upload_file_helper = mock.Mock() + import_summary = ImportStatusSummary('test_import') + + with tempfile.TemporaryDirectory() as import_dir: + status = executor._invoke_import_validation( + 'repo', 'relative', import_dir, { + 'import_name': 'test_import', + 'import_inputs': [{}, {}], + }, 'version', import_summary) + + self.assertFalse(status) + self.assertEqual(2, mock_log_metric.call_count) + for input_index, call in enumerate(mock_log_metric.call_args_list): + self.assertEqual('ERROR', call.args[1]) + self.assertEqual(f'input{input_index}', + call.args[3]['import_input']) + self.assertEqual((input_index + 1) * 10.0, + call.args[3]['deleted_records_percent']) + self.assertEqual(input_index + 1, + call.args[3]['deleted_records_count']) + self.assertEqual(10, call.args[3]['previous_obs_count']) From 7c97709b47c1e765cece466f1818101a9190ba7f Mon Sep 17 00:00:00 2001 From: rohit kumar Date: Wed, 26 Aug 2026 12:15:08 +0000 Subject: [PATCH 2/5] Harden deletion metric logging --- .../executor/app/executor/import_executor.py | 4 +-- .../executor/test/import_executor_test.py | 27 ++++++++++++++----- 2 files changed, 22 insertions(+), 9 deletions(-) diff --git a/import-automation/executor/app/executor/import_executor.py b/import-automation/executor/app/executor/import_executor.py index 4f2a314aa4..87e9d1d599 100644 --- a/import-automation/executor/app/executor/import_executor.py +++ b/import-automation/executor/app/executor/import_executor.py @@ -691,9 +691,9 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, validation_status = overall_status validator_by_rule_id = { - rule['rule_id']: rule['validator'] + rule['rule_id']: rule.get('validator') for rule in validation.config.rules - if rule.get('enabled', True) + if rule.get('enabled', True) and rule.get('rule_id') } for result in current_results: if (validator_by_rule_id.get( diff --git a/import-automation/executor/test/import_executor_test.py b/import-automation/executor/test/import_executor_test.py index 763b8a929a..03e89df8c5 100644 --- a/import-automation/executor/test/import_executor_test.py +++ b/import-automation/executor/test/import_executor_test.py @@ -128,15 +128,24 @@ def test_construct_process_message_no_output(self): def test_validation_metrics_include_deleted_percent_for_each_input( self, mock_validation_runner, mock_log_metric, _): validation_runners = [] + input_statuses = (False, True) for input_index, percent in enumerate((10.0, 20.0)): rule_id = f'deleted_percent_{input_index}' runner = mock.Mock() - runner.config.rules = [{ - 'rule_id': rule_id, - 'validator': 'DELETED_RECORDS_PERCENT', - }] - runner.run_validations.return_value = (False, [ - ValidationResult(ValidationStatus.FAILED, + runner.config.rules = [ + { + 'validator': 'UNKNOWN_VALIDATOR', + }, + { + 'rule_id': rule_id, + 'validator': 'DELETED_RECORDS_PERCENT', + }, + ] + input_status = input_statuses[input_index] + result_status = (ValidationStatus.PASSED + if input_status else ValidationStatus.FAILED) + runner.run_validations.return_value = (input_status, [ + ValidationResult(result_status, rule_id, details={ 'percent': percent, @@ -172,9 +181,13 @@ def test_validation_metrics_include_deleted_percent_for_each_input( self.assertFalse(status) self.assertEqual(2, mock_log_metric.call_count) for input_index, call in enumerate(mock_log_metric.call_args_list): - self.assertEqual('ERROR', call.args[1]) + expected_level = 'INFO' if input_statuses[input_index] else 'ERROR' + expected_status = ('SUCCESS' + if input_statuses[input_index] else 'FAILURE') + self.assertEqual(expected_level, call.args[1]) self.assertEqual(f'input{input_index}', call.args[3]['import_input']) + self.assertEqual(expected_status, call.args[3]['status']) self.assertEqual((input_index + 1) * 10.0, call.args[3]['deleted_records_percent']) self.assertEqual(input_index + 1, From 6ec5dfa0d27c14c86a547f7f7b42465fabbfd8f9 Mon Sep 17 00:00:00 2001 From: rohit kumar Date: Wed, 26 Aug 2026 12:39:19 +0000 Subject: [PATCH 3/5] Address validation metric review feedback --- .../executor/app/executor/import_executor.py | 4 ++- .../executor/test/import_executor_test.py | 26 ++++++++++++++++--- 2 files changed, 26 insertions(+), 4 deletions(-) diff --git a/import-automation/executor/app/executor/import_executor.py b/import-automation/executor/app/executor/import_executor.py index 87e9d1d599..e414d5d473 100644 --- a/import-automation/executor/app/executor/import_executor.py +++ b/import-automation/executor/app/executor/import_executor.py @@ -693,7 +693,8 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, validator_by_rule_id = { rule['rule_id']: rule.get('validator') for rule in validation.config.rules - if rule.get('enabled', True) and rule.get('rule_id') + if isinstance(rule, dict) and rule.get('enabled', True) and + isinstance(rule.get('rule_id'), str) } for result in current_results: if (validator_by_rule_id.get( @@ -710,6 +711,7 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, except ValueError as e: logging.error('ValidationRunner failed: %s', e) validation_status = False + input_validation_status = False validation_metrics.update({ "latency": timer.time(), diff --git a/import-automation/executor/test/import_executor_test.py b/import-automation/executor/test/import_executor_test.py index 03e89df8c5..487b05a881 100644 --- a/import-automation/executor/test/import_executor_test.py +++ b/import-automation/executor/test/import_executor_test.py @@ -133,9 +133,14 @@ def test_validation_metrics_include_deleted_percent_for_each_input( rule_id = f'deleted_percent_{input_index}' runner = mock.Mock() runner.config.rules = [ + 'INVALID_RULE', { 'validator': 'UNKNOWN_VALIDATOR', }, + { + 'rule_id': ['invalid'], + 'validator': 'UNKNOWN_VALIDATOR', + }, { 'rule_id': rule_id, 'validator': 'DELETED_RECORDS_PERCENT', @@ -154,6 +159,16 @@ def test_validation_metrics_include_deleted_percent_for_each_input( }) ]) validation_runners.append(runner) + + class InvalidValidationResults: + + def __iter__(self): + raise ValueError('invalid validation results') + + error_runner = mock.Mock() + error_runner.run_validations.return_value = (True, + InvalidValidationResults()) + validation_runners.append(error_runner) mock_validation_runner.side_effect = validation_runners config = mock.Mock(invoke_differ_tool=True, @@ -175,12 +190,12 @@ def test_validation_metrics_include_deleted_percent_for_each_input( status = executor._invoke_import_validation( 'repo', 'relative', import_dir, { 'import_name': 'test_import', - 'import_inputs': [{}, {}], + 'import_inputs': [{}, {}, {}], }, 'version', import_summary) self.assertFalse(status) - self.assertEqual(2, mock_log_metric.call_count) - for input_index, call in enumerate(mock_log_metric.call_args_list): + self.assertEqual(3, mock_log_metric.call_count) + for input_index, call in enumerate(mock_log_metric.call_args_list[:2]): expected_level = 'INFO' if input_statuses[input_index] else 'ERROR' expected_status = ('SUCCESS' if input_statuses[input_index] else 'FAILURE') @@ -193,3 +208,8 @@ def test_validation_metrics_include_deleted_percent_for_each_input( self.assertEqual(input_index + 1, call.args[3]['deleted_records_count']) self.assertEqual(10, call.args[3]['previous_obs_count']) + + error_call = mock_log_metric.call_args_list[2] + self.assertEqual('ERROR', error_call.args[1]) + self.assertEqual('input2', error_call.args[3]['import_input']) + self.assertEqual('FAILURE', error_call.args[3]['status']) From a185a61fec0196d57532a932c57f89a68f3b8099 Mon Sep 17 00:00:00 2001 From: rohit kumar Date: Wed, 26 Aug 2026 17:15:54 +0000 Subject: [PATCH 4/5] Log structured import validation results Include per-input rule details with safe serialization and bounded payload size. --- .../executor/app/executor/import_executor.py | 106 +++++---- .../executor/test/import_executor_test.py | 216 ++++++++++++------ 2 files changed, 217 insertions(+), 105 deletions(-) diff --git a/import-automation/executor/app/executor/import_executor.py b/import-automation/executor/app/executor/import_executor.py index e414d5d473..c2846a7861 100644 --- a/import-automation/executor/app/executor/import_executor.py +++ b/import-automation/executor/app/executor/import_executor.py @@ -20,6 +20,8 @@ import glob import json import logging +import math +import numbers import os import shutil import shlex @@ -79,6 +81,7 @@ IMPORT_SUMMARY_FILE = "import_summary.json" STAGING_VERSION_FILE = "staging_version.txt" MAX_LOG_CHUNK_SIZE = 50000 +MAX_VALIDATION_RESULTS_LOG_SIZE_BYTES = 150 * 1024 class ImportStatus(Enum): @@ -99,6 +102,39 @@ class ImportStage(Enum): FINISH = 6 +def _format_validation_log_result(input_prefix: str, + result: ValidationResult) -> dict: + details = [] + if isinstance(result.details, dict): + for field, value in result.details.items(): + if not isinstance(field, str): + continue + + detail = {'field': field} + if isinstance(value, bool): + detail['bool_value'] = value + elif isinstance(value, numbers.Real): + try: + number_value = float(value) + except (OverflowError, TypeError, ValueError): + continue + if not math.isfinite(number_value): + continue + detail['number_value'] = number_value + elif isinstance(value, str) and len(value) <= 256: + detail['string_value'] = value + else: + continue + details.append(detail) + + return { + 'input_prefix': input_prefix, + 'rule_id': str(result.name), + 'status': result.status.name, + 'details': details, + } + + @dataclasses.dataclass class ExecutionResult: """Describes the result of the execution of an import.""" @@ -617,6 +653,7 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, validation_status = True differ_status = False validation_results = [] + validation_log_results = [] import_dir = f'{relative_import_dir}/{import_spec["import_name"]}' latest_version = self._get_latest_version(import_dir) @@ -666,11 +703,6 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, if not differ_status: differ_status = diff_found timer = Timer() - input_validation_status = False - validation_metrics = { - "stage": import_stage.name, - "import_input": import_prefix, - } try: config_file_path = self._get_validation_config_file( repo_dir, absolute_import_dir, import_spec, @@ -685,45 +717,27 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, lint_report=report_json, validation_output=validation_output_file) overall_status, current_results = validation.run_validations() - input_validation_status = overall_status validation_results.extend(current_results) if validation_status: validation_status = overall_status - validator_by_rule_id = { - rule['rule_id']: rule.get('validator') - for rule in validation.config.rules - if isinstance(rule, dict) and rule.get('enabled', True) and - isinstance(rule.get('rule_id'), str) - } - for result in current_results: - if (validator_by_rule_id.get( - result.name) == 'DELETED_RECORDS_PERCENT'): - validation_metrics.update({ - 'deleted_records_percent': - result.details.get('percent'), - 'deleted_records_count': - result.details.get('deleted_records_count'), - 'previous_obs_count': - result.details.get('previous_obs_count'), - }) - break + validation_log_results.extend( + _format_validation_log_result(import_prefix, result) + for result in current_results) except ValueError as e: logging.error('ValidationRunner failed: %s', e) validation_status = False - input_validation_status = False - validation_metrics.update({ - "latency": - timer.time(), - "status": - ImportStatus.SUCCESS.name - if input_validation_status else ImportStatus.FAILURE.name, - }) log_metric( - AUTO_IMPORT_JOB_STAGE, - "INFO" if input_validation_status else "ERROR", - f"Import: {import_name}, input: {import_prefix}, validation: {input_validation_status}", - validation_metrics) + AUTO_IMPORT_JOB_STAGE, "INFO" if validation_status else "ERROR", + f"Import: {import_name}, validation: {validation_status}", { + "stage": + import_stage.name, + "latency": + timer.time(), + "status": + ImportStatus.SUCCESS.name + if validation_status else ImportStatus.FAILURE.name, + }) if os.path.exists(validation_output_path): # Upload output to GCS. @@ -742,11 +756,13 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, import_summary.import_stats['validation_data_size'] = data_size validation_message = self._get_validation_message(validation_results) log_import_status( - import_name, import_stage, + import_name, + import_stage, ImportStatus.SUCCESS if validation_status else ImportStatus.FAILURE, import_summary.import_stats.get('validation_execution_time', 0), - import_summary.import_stats.get('validation_data_size', - 0), validation_message) + import_summary.import_stats.get('validation_data_size', 0), + validation_message, + validation_results=validation_log_results) if not self.config.ignore_validation_status and not validation_status: logging.error( "Marking import as VALIDATION due to validation failure.") @@ -1496,7 +1512,8 @@ def log_import_status(import_name: str, latency_secs: int = 0, data_size: int = 0, message: str = '', - level: str = "INFO") -> None: + level: str = "INFO", + validation_results: list[dict] = None) -> None: """Logs the status of import. """ import_metrics = { @@ -1506,6 +1523,15 @@ def log_import_status(import_name: str, "latency_secs": int(latency_secs), "data_bytes": data_size } + if validation_results is not None: + import_metrics['validation_results'] = validation_results + # Preserve the final status log if validation details would exceed the + # Cloud Logging entry limit, leaving headroom for the log envelope. + if len(json.dumps(import_metrics).encode( + 'utf-8')) > MAX_VALIDATION_RESULTS_LOG_SIZE_BYTES: + import_metrics['validation_results'] = [] + import_metrics['validation_results_omitted_count'] = len( + validation_results) if not message: message = f'Import: {import_name} stage: {import_stage.name} status: {status.name}' log_metric(AUTO_IMPORT_JOB_STATUS, level, message, import_metrics) diff --git a/import-automation/executor/test/import_executor_test.py b/import-automation/executor/test/import_executor_test.py index 487b05a881..5cf9b1f8c3 100644 --- a/import-automation/executor/test/import_executor_test.py +++ b/import-automation/executor/test/import_executor_test.py @@ -125,51 +125,34 @@ def test_construct_process_message_no_output(self): @mock.patch.object(import_executor, 'log_import_status') @mock.patch.object(import_executor, 'log_metric') @mock.patch.object(import_executor, 'ValidationRunner') - def test_validation_metrics_include_deleted_percent_for_each_input( - self, mock_validation_runner, mock_log_metric, _): - validation_runners = [] - input_statuses = (False, True) - for input_index, percent in enumerate((10.0, 20.0)): - rule_id = f'deleted_percent_{input_index}' - runner = mock.Mock() - runner.config.rules = [ - 'INVALID_RULE', - { - 'validator': 'UNKNOWN_VALIDATOR', - }, - { - 'rule_id': ['invalid'], - 'validator': 'UNKNOWN_VALIDATOR', - }, - { - 'rule_id': rule_id, - 'validator': 'DELETED_RECORDS_PERCENT', - }, - ] - input_status = input_statuses[input_index] - result_status = (ValidationStatus.PASSED - if input_status else ValidationStatus.FAILED) - runner.run_validations.return_value = (input_status, [ - ValidationResult(result_status, - rule_id, - details={ - 'percent': percent, - 'deleted_records_count': input_index + 1, - 'previous_obs_count': 10, - }) - ]) - validation_runners.append(runner) - - class InvalidValidationResults: - - def __iter__(self): - raise ValueError('invalid validation results') - - error_runner = mock.Mock() - error_runner.run_validations.return_value = (True, - InvalidValidationResults()) - validation_runners.append(error_runner) - mock_validation_runner.side_effect = validation_runners + def test_validation_results_are_added_to_final_import_status( + self, mock_validation_runner, mock_log_metric, + mock_log_import_status): + first_runner = mock.Mock() + first_runner.run_validations.return_value = (False, [ + ValidationResult(ValidationStatus.FAILED, + 'check_deleted_records_percent', + details={ + 'percent': 10, + 'threshold': 1.5, + 'enabled': True, + 'note': 'short value', + 'long_value': 'x' * 257, + 'missing_goldens': ['dcid:example'], + 'not_a_number': float('nan'), + 'too_large': 10**10000, + }), + ValidationResult(ValidationStatus.PASSED, + 'check_non_scalar_details', + details={'failed_rows': []}), + ]) + second_runner = mock.Mock() + second_runner.run_validations.return_value = (True, [ + ValidationResult(ValidationStatus.PASSED, + 123, + details={'missing_refs_count': 2}) + ]) + mock_validation_runner.side_effect = [first_runner, second_runner] config = mock.Mock(invoke_differ_tool=True, ignore_validation_status=False, @@ -190,26 +173,129 @@ def __iter__(self): status = executor._invoke_import_validation( 'repo', 'relative', import_dir, { 'import_name': 'test_import', - 'import_inputs': [{}, {}, {}], + 'import_inputs': [{}, {}], }, 'version', import_summary) self.assertFalse(status) - self.assertEqual(3, mock_log_metric.call_count) - for input_index, call in enumerate(mock_log_metric.call_args_list[:2]): - expected_level = 'INFO' if input_statuses[input_index] else 'ERROR' - expected_status = ('SUCCESS' - if input_statuses[input_index] else 'FAILURE') - self.assertEqual(expected_level, call.args[1]) - self.assertEqual(f'input{input_index}', - call.args[3]['import_input']) - self.assertEqual(expected_status, call.args[3]['status']) - self.assertEqual((input_index + 1) * 10.0, - call.args[3]['deleted_records_percent']) - self.assertEqual(input_index + 1, - call.args[3]['deleted_records_count']) - self.assertEqual(10, call.args[3]['previous_obs_count']) - - error_call = mock_log_metric.call_args_list[2] - self.assertEqual('ERROR', error_call.args[1]) - self.assertEqual('input2', error_call.args[3]['import_input']) - self.assertEqual('FAILURE', error_call.args[3]['status']) + self.assertEqual(2, mock_log_metric.call_count) + for call in mock_log_metric.call_args_list: + self.assertEqual(import_executor.AUTO_IMPORT_JOB_STAGE, + call.args[0]) + self.assertEqual('ERROR', call.args[1]) + self.assertEqual('Import: test_import, validation: False', + call.args[2]) + self.assertEqual({'stage', 'latency', 'status'}, set(call.args[3])) + self.assertEqual('VALIDATION', call.args[3]['stage']) + self.assertEqual('FAILURE', call.args[3]['status']) + + self.assertEqual(ImportStatus.FAILURE, + mock_log_import_status.call_args.args[2]) + logged_results = mock_log_import_status.call_args.kwargs[ + 'validation_results'] + self.assertEqual([ + ('input0', 'check_deleted_records_percent', 'FAILED'), + ('input0', 'check_non_scalar_details', 'PASSED'), + ('input1', '123', 'PASSED'), + ], [(result['input_prefix'], result['rule_id'], result['status']) + for result in logged_results]) + self.assertEqual([ + { + 'field': 'percent', + 'number_value': 10.0, + }, + { + 'field': 'threshold', + 'number_value': 1.5, + }, + { + 'field': 'enabled', + 'bool_value': True, + }, + { + 'field': 'note', + 'string_value': 'short value', + }, + ], logged_results[0]['details']) + self.assertEqual([], logged_results[1]['details']) + self.assertEqual([{ + 'field': 'missing_refs_count', + 'number_value': 2.0, + }], logged_results[2]['details']) + + @mock.patch.object(import_executor, 'log_import_status') + @mock.patch.object(import_executor, 'log_metric') + @mock.patch.object(import_executor, 'ValidationRunner') + def test_validation_with_no_rules_logs_empty_results( + self, mock_validation_runner, _, mock_log_import_status): + mock_validation_runner.return_value.run_validations.return_value = ( + True, []) + config = mock.Mock(invoke_differ_tool=False, + ignore_validation_status=False, + enable_skip_status=False) + executor = import_executor.ImportExecutor(mock.Mock(), mock.Mock(), + config) + executor._get_latest_version = mock.Mock(return_value='previous') + executor._get_validation_config_file = mock.Mock( + return_value='validation_config.json') + import_summary = ImportStatusSummary('test_import') + + with tempfile.TemporaryDirectory() as import_dir: + status = executor._invoke_import_validation( + 'repo', 'relative', import_dir, { + 'import_name': 'test_import', + 'import_inputs': [{}], + }, 'version', import_summary) + + self.assertTrue(status) + self.assertEqual(ImportStatus.SUCCESS, + mock_log_import_status.call_args.args[2]) + self.assertEqual([], mock_log_import_status.call_args.kwargs[ + 'validation_results']) + + @mock.patch.object(import_executor, 'log_metric') + def test_log_import_status_includes_validation_results( + self, mock_log_metric): + for validation_results in ([{ + 'input_prefix': 'input0', + 'rule_id': 'check_empty_import', + 'status': 'PASSED', + 'details': [], + }], []): + with self.subTest(validation_results=validation_results): + import_executor.log_import_status( + 'test_import', + import_executor.ImportStage.VALIDATION, + ImportStatus.SUCCESS, + validation_results=validation_results) + + self.assertEqual( + validation_results, + mock_log_metric.call_args.args[3]['validation_results']) + self.assertNotIn('validation_results_omitted_count', + mock_log_metric.call_args.args[3]) + mock_log_metric.reset_mock() + + @mock.patch.object(import_executor, 'log_metric') + def test_log_import_status_omits_oversized_validation_results( + self, mock_log_metric): + validation_results = [{ + 'input_prefix': 'input0', + 'rule_id': 'check_empty_import', + 'status': 'PASSED', + 'details': [], + }] + + with mock.patch.object(import_executor, + 'MAX_VALIDATION_RESULTS_LOG_SIZE_BYTES', 1): + import_executor.log_import_status( + 'test_import', + import_executor.ImportStage.VALIDATION, + ImportStatus.SUCCESS, + validation_results=validation_results) + + metrics = mock_log_metric.call_args.args[3] + self.assertEqual([], metrics['validation_results']) + self.assertEqual(1, metrics['validation_results_omitted_count']) + self.assertEqual('test_import', metrics['import_name']) + self.assertEqual('VALIDATION', metrics['stage_name']) + self.assertEqual('SUCCESS', metrics['status']) From e4a14fe47f7a66aa244526f1c12c3557dcabc4df Mon Sep 17 00:00:00 2001 From: rohit kumar Date: Wed, 26 Aug 2026 17:24:57 +0000 Subject: [PATCH 5/5] Improve validation logging coverage --- .../executor/test/import_executor_test.py | 25 ++++++++++++++----- 1 file changed, 19 insertions(+), 6 deletions(-) diff --git a/import-automation/executor/test/import_executor_test.py b/import-automation/executor/test/import_executor_test.py index 5cf9b1f8c3..900f3dc360 100644 --- a/import-automation/executor/test/import_executor_test.py +++ b/import-automation/executor/test/import_executor_test.py @@ -225,10 +225,16 @@ def test_validation_results_are_added_to_final_import_status( @mock.patch.object(import_executor, 'log_import_status') @mock.patch.object(import_executor, 'log_metric') @mock.patch.object(import_executor, 'ValidationRunner') - def test_validation_with_no_rules_logs_empty_results( + def test_successful_validation_logs_results_when_an_input_has_no_rules( self, mock_validation_runner, _, mock_log_import_status): - mock_validation_runner.return_value.run_validations.return_value = ( - True, []) + mock_validation_runner.return_value.run_validations.side_effect = [ + (True, []), + (True, [ + ValidationResult(ValidationStatus.PASSED, + 'check_missing_refs', + details={'missing_refs_count': 0}) + ]), + ] config = mock.Mock(invoke_differ_tool=False, ignore_validation_status=False, enable_skip_status=False) @@ -243,14 +249,21 @@ def test_validation_with_no_rules_logs_empty_results( status = executor._invoke_import_validation( 'repo', 'relative', import_dir, { 'import_name': 'test_import', - 'import_inputs': [{}], + 'import_inputs': [{}, {}], }, 'version', import_summary) self.assertTrue(status) self.assertEqual(ImportStatus.SUCCESS, mock_log_import_status.call_args.args[2]) - self.assertEqual([], mock_log_import_status.call_args.kwargs[ - 'validation_results']) + self.assertEqual([{ + 'input_prefix': 'input1', + 'rule_id': 'check_missing_refs', + 'status': 'PASSED', + 'details': [{ + 'field': 'missing_refs_count', + 'number_value': 0.0, + }], + }], mock_log_import_status.call_args.kwargs['validation_results']) @mock.patch.object(import_executor, 'log_metric') def test_log_import_status_includes_validation_results(