diff --git a/import-automation/executor/app/executor/import_executor.py b/import-automation/executor/app/executor/import_executor.py index 47ec5f79e1..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) @@ -683,6 +720,10 @@ def _invoke_import_validation(self, repo_dir: str, relative_import_dir: str, validation_results.extend(current_results) if validation_status: validation_status = overall_status + + 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 @@ -715,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.") @@ -1469,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 = { @@ -1479,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 5a863f00e0..900f3dc360 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,194 @@ 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_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, + 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 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_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.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) + 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([{ + '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( + 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'])