From b245abd82167cf5afb9c8f093ed9ddbb186b64f0 Mon Sep 17 00:00:00 2001 From: stevenhsd <56357022+stevenhsd@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:54:54 +0100 Subject: [PATCH] fix: ensure that original entity detials persisted in feedback messages for further processing --- src/dve/core_engine/backends/base/rules.py | 8 ++------ .../backends/implementations/duckdb/duckdb_helpers.py | 4 ++-- .../backends/implementations/spark/spark_helpers.py | 4 ++-- src/dve/core_engine/message.py | 6 ++++-- src/dve/core_engine/type_hints.py | 1 + src/dve/pipeline/pipeline.py | 2 +- src/dve/reporting/error_report.py | 2 +- tests/features/movies.feature | 10 +++++----- tests/features/steps/steps_pipeline.py | 2 +- tests/features/steps/utilities.py | 2 +- .../test_duckdb/test_data_contract.py | 2 +- .../test_duckdb/test_duckdb_helpers.py | 6 ++++-- .../test_spark/test_data_contract.py | 2 +- .../test_spark/test_spark_helpers.py | 4 ++-- tests/test_pipeline/pipeline_helpers.py | 4 ++-- tests/test_pipeline/test_spark_pipeline.py | 10 +++++----- 16 files changed, 35 insertions(+), 34 deletions(-) diff --git a/src/dve/core_engine/backends/base/rules.py b/src/dve/core_engine/backends/base/rules.py index ee93106..0b5b8a0 100644 --- a/src/dve/core_engine/backends/base/rules.py +++ b/src/dve/core_engine/backends/base/rules.py @@ -448,8 +448,7 @@ def process_node(node: HierarchyNode): record=record, # type: ignore error_location=location, error_message=template_object( - node.missing_parent_id_error_message, - record + node.missing_parent_id_error_message, record ), failure_type="record", error_type="record", @@ -522,10 +521,7 @@ def process_node(node: HierarchyNode) -> bool: entity=node.parent_entity, record=record, # type: ignore error_location=location, - error_message=template_object( - node.no_valid_records_error_message, - record - ), + error_message=template_object(node.no_valid_records_error_message, record), failure_type="record", error_type="record", error_code=node.no_valid_records_error_code, diff --git a/src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py b/src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py index 588cd7e..d9cd149 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py +++ b/src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py @@ -284,11 +284,11 @@ def _ddb_filter_contract_errors( "RecordIndex": "INTEGER", "FailureType": "STRING", "Status": "STRING", - "Entity": "STRING", + "OriginalEntity": "STRING", }, ) .filter( - f"FailureType == 'record' AND Status != 'informational' AND Entity = '{entity_name}'" + f"FailureType == 'record' AND Status != 'informational' AND OriginalEntity = '{entity_name}'" # pylint: disable=C0301 ) # pylint: disable=C0301 .select("RecordIndex") .distinct() diff --git a/src/dve/core_engine/backends/implementations/spark/spark_helpers.py b/src/dve/core_engine/backends/implementations/spark/spark_helpers.py index 8c14132..0319a26 100644 --- a/src/dve/core_engine/backends/implementations/spark/spark_helpers.py +++ b/src/dve/core_engine/backends/implementations/spark/spark_helpers.py @@ -415,14 +415,14 @@ def _spark_filter_contract_errors( st.StructField("RecordIndex", st.IntegerType()), st.StructField("FailureType", st.StringType()), st.StructField("Status", st.StringType()), - st.StructField("Entity", st.StringType()), + st.StructField("OriginalEntity", st.StringType()), ] ), ) .filter( (sf.col("FailureType") == sf.lit("record")) & (sf.col("Status") != sf.lit("informational")) - & (sf.col("Entity") == sf.lit(entity_name)) + & (sf.col("OriginalEntity") == sf.lit(entity_name)) ) .distinct() .orderBy(sf.asc(sf.col("RecordIndex"))) diff --git a/src/dve/core_engine/message.py b/src/dve/core_engine/message.py index 78024e9..05dbc17 100644 --- a/src/dve/core_engine/message.py +++ b/src/dve/core_engine/message.py @@ -90,7 +90,7 @@ def extract_error_value(records, error_location): class UserMessage: """The structure of the message that is used to populate the error report.""" - Entity: Optional[str] + ReportingEntity: Optional[str] """The entity that the message pertains to (if applicable).""" Key: Optional[str] "The key field(s) in string format to allow users to identify the record" @@ -176,7 +176,8 @@ class FeedbackMessage: # pylint: disable=too-many-instance-attributes """The category of the error.""" HEADER: ClassVar[list[str]] = [ - "Entity", + "ReportingEntity", + "OriginalEntity", "Key", "FailureType", "Status", @@ -307,6 +308,7 @@ def to_row( return ( self.entity, + self.original_entity, key, self.failure_type, "informational" if self.is_informational else "error", diff --git a/src/dve/core_engine/type_hints.py b/src/dve/core_engine/type_hints.py index 48d9eeb..1073886 100644 --- a/src/dve/core_engine/type_hints.py +++ b/src/dve/core_engine/type_hints.py @@ -141,6 +141,7 @@ """The record index that the error relates to (if applicable)""" MessageTuple = tuple[ + Optional[EntityName], Optional[EntityName], Key, FailureType, diff --git a/src/dve/pipeline/pipeline.py b/src/dve/pipeline/pipeline.py index ee5a6bc..8e1379a 100644 --- a/src/dve/pipeline/pipeline.py +++ b/src/dve/pipeline/pipeline.py @@ -865,7 +865,7 @@ def _get_error_dataframes(self, submission_id: str): .alias("error_type") # type: ignore ) df = df.select( - pl.col("Entity").alias("Table"), # type: ignore + pl.col("ReportingEntity").alias("Table"), # type: ignore pl.col("error_type").alias("Type"), # type: ignore pl.col("ErrorCode").alias("Error_Code"), # type: ignore pl.col("ReportingField").alias("Data_Item"), # type: ignore diff --git a/src/dve/reporting/error_report.py b/src/dve/reporting/error_report.py index 95137b5..bf1708a 100644 --- a/src/dve/reporting/error_report.py +++ b/src/dve/reporting/error_report.py @@ -91,7 +91,7 @@ def create_error_dataframe(errors: deque[FeedbackMessage], key_fields): .alias("error_type") ) df = df.select( # type: ignore - col("Entity").alias("Table"), # type: ignore + col("ReportingEntity").alias("Table"), # type: ignore col("error_type").alias("Type"), # type: ignore col("ErrorCode").alias("Error_Code"), # type: ignore col("ReportingField").alias("Data_Item"), # type: ignore diff --git a/tests/features/movies.feature b/tests/features/movies.feature index 750975e..c6d5ad6 100644 --- a/tests/features/movies.feature +++ b/tests/features/movies.feature @@ -22,7 +22,7 @@ Feature: Pipeline tests using the movies dataset Then there is 1 submission rejection from the data_contract phase And there are 3 record rejections from the data_contract phase And there are errors with the following details and associated error_count from the data_contract phase - | Entity | ErrorCode | ErrorMessage | RecordIndex | error_count | + | ReportingEntity | ErrorCode | ErrorMessage | RecordIndex | error_count | | movies | BLANKYEAR | year not provided | 2 | 1 | | movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 | | movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 | @@ -61,10 +61,10 @@ Feature: Pipeline tests using the movies dataset Then there is 1 submission rejection from the data_contract phase And there are 3 record rejections from the data_contract phase And there are errors with the following details and associated error_count from the data_contract phase - | Entity | ErrorCode | ErrorMessage | RecordIndex | error_count | - | movies | BLANKYEAR | year not provided | 2 | 1 | - | movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 | - | movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 | + | ReportingEntity | ErrorCode | ErrorMessage | RecordIndex | error_count | + | movies | BLANKYEAR | year not provided | 2 | 1 | + | movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 | + | movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 | | movies | BLANKTITLE | title should not be blank | 4 | 1 | And the movies entity is stored as a parquet after the data_contract phase And the latest audit record for the submission is marked with processing status business_rules diff --git a/tests/features/steps/steps_pipeline.py b/tests/features/steps/steps_pipeline.py index c71bd45..471bbcc 100644 --- a/tests/features/steps/steps_pipeline.py +++ b/tests/features/steps/steps_pipeline.py @@ -284,7 +284,7 @@ def check_rows_removed_with_error_code(context: Context, entity_name: str, error err_df = get_all_errors_df(context) recs_with_err_code = err_df.filter( - (pl.col("Entity").eq(entity_name)) & (pl.col("ErrorCode").eq(error_code)) + (pl.col("ReportingEntity").eq(entity_name)) & (pl.col("ErrorCode").eq(error_code)) ).shape[0] assert recs_with_err_code >= 1 diff --git a/tests/features/steps/utilities.py b/tests/features/steps/utilities.py index 58edc67..f70dde8 100644 --- a/tests/features/steps/utilities.py +++ b/tests/features/steps/utilities.py @@ -15,7 +15,7 @@ from dve.parser.type_hints import URI ERROR_DF_FIELDS: List[str] = [ - "Entity", + "ReportingEntity", "Key", "ErrorCode", "FailureType", diff --git a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_data_contract.py b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_data_contract.py index 2019a66..2e6bf87 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_data_contract.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_data_contract.py @@ -381,4 +381,4 @@ def test_duckdb_data_contract_custom_error_details(nested_all_string_parquet_w_e assert messages[0].ErrorMessage == "subfield id is invalid: subfield.id - WRONG" assert messages[1].ErrorCode == "TESTIDBAD" assert messages[1].ErrorMessage == "id is invalid: id - WRONG" - assert messages[1].Entity == "test_rename" \ No newline at end of file + assert messages[1].ReportingEntity == "test_rename" \ No newline at end of file diff --git a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_duckdb_helpers.py b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_duckdb_helpers.py index 4eefe05..0e81cf4 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_duckdb_helpers.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_duckdb_helpers.py @@ -76,7 +76,8 @@ def example_data_contract_error_codes(temp_ddb_conn): test_entity = con.sql("SELECT * FROM test_df") error_contract_messages = [ { - "Entity": "test_entity", + "ReportingEntity": "test_entity", + "OriginalEntity": "test_entity", "Key": "", "FailureType": "record", "Status": "error", @@ -90,7 +91,8 @@ def example_data_contract_error_codes(temp_ddb_conn): "Category": "Bad value" }, { - "Entity": "test_entity", + "ReportingEntity": "test_entity", + "OriginalEntity": "test_entity", "Key": "", "FailureType": "record", "Status": "error", diff --git a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_data_contract.py b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_data_contract.py index 70c6b9c..ea5fc46 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_data_contract.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_data_contract.py @@ -242,6 +242,6 @@ def test_spark_data_contract_custom_error_details(nested_all_string_parquet_w_er assert messages[0].ErrorMessage == "subfield id is invalid: subfield.id - WRONG" assert messages[1].ErrorCode == "TESTIDBAD" assert messages[1].ErrorMessage == "id is invalid: id - WRONG" - assert messages[1].Entity == "test_rename" + assert messages[1].ReportingEntity == "test_rename" \ No newline at end of file diff --git a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_spark_helpers.py b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_spark_helpers.py index e7a37eb..41c9e93 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_spark_helpers.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_spark_helpers.py @@ -63,7 +63,7 @@ def example_data_contract_error_codes(spark: SparkSession): ]) error_contract_messages = [ { - "Entity": "test_entity", + "OriginalEntity": "test_entity", "Key": "", "FailureType": "record", "Status": "error", @@ -77,7 +77,7 @@ def example_data_contract_error_codes(spark: SparkSession): "Category": "Bad value" }, { - "Entity": "test_entity", + "OriginalEntity": "test_entity", "Key": "", "FailureType": "record", "Status": "error", diff --git a/tests/test_pipeline/pipeline_helpers.py b/tests/test_pipeline/pipeline_helpers.py index b13bef3..efed6e2 100644 --- a/tests/test_pipeline/pipeline_helpers.py +++ b/tests/test_pipeline/pipeline_helpers.py @@ -372,7 +372,7 @@ def error_data_after_business_rules() -> Iterator[Tuple[SubmissionInfo, str]]: error_data = json.loads( """[ { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", @@ -386,7 +386,7 @@ def error_data_after_business_rules() -> Iterator[Tuple[SubmissionInfo, str]]: "RecordIndex": "1" }, { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", diff --git a/tests/test_pipeline/test_spark_pipeline.py b/tests/test_pipeline/test_spark_pipeline.py index dd28e26..99d2fa4 100644 --- a/tests/test_pipeline/test_spark_pipeline.py +++ b/tests/test_pipeline/test_spark_pipeline.py @@ -162,7 +162,7 @@ def test_apply_data_contract_failed( # pylint: disable=redefined-outer-name expected_errors = [ { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", @@ -176,7 +176,7 @@ def test_apply_data_contract_failed( # pylint: disable=redefined-outer-name "Category": "Bad value", }, { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", @@ -190,7 +190,7 @@ def test_apply_data_contract_failed( # pylint: disable=redefined-outer-name "Category": "Bad value", }, { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", @@ -326,7 +326,7 @@ def test_apply_business_rules_with_data_errors( # pylint: disable=redefined-out expected_errors = [ { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error", @@ -340,7 +340,7 @@ def test_apply_business_rules_with_data_errors( # pylint: disable=redefined-out "RecordIndex": "1" }, { - "Entity": "planets", + "ReportingEntity": "planets", "Key": "", "FailureType": "record", "Status": "error",