Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 57 additions & 4 deletions import-automation/executor/app/executor/import_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
import glob
import json
import logging
import math
import numbers
import os
import shutil
import shlex
Expand Down Expand Up @@ -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):
Expand All @@ -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."""
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Comment thread
rohitkumarbhagat marked this conversation as resolved.
Expand Down Expand Up @@ -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.")
Expand Down Expand Up @@ -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 = {
Expand All @@ -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)
194 changes: 193 additions & 1 deletion import-automation/executor/test/import_executor_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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'])
Loading