Skip to content

Commit 8f5e1ae

Browse files
authored
fix: ensure that original entity detials persisted in feedback messages for further processing (#168)
1 parent b3e03b4 commit 8f5e1ae

15 files changed

Lines changed: 33 additions & 28 deletions

File tree

‎src/dve/core_engine/backends/implementations/duckdb/duckdb_helpers.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -284,11 +284,11 @@ def _ddb_filter_contract_errors(
284284
"RecordIndex": "INTEGER",
285285
"FailureType": "STRING",
286286
"Status": "STRING",
287-
"Entity": "STRING",
287+
"OriginalEntity": "STRING",
288288
},
289289
)
290290
.filter(
291-
f"FailureType == 'record' AND Status != 'informational' AND Entity = '{entity_name}'"
291+
f"FailureType == 'record' AND Status != 'informational' AND OriginalEntity = '{entity_name}'" # pylint: disable=C0301
292292
) # pylint: disable=C0301
293293
.select("RecordIndex")
294294
.distinct()

‎src/dve/core_engine/backends/implementations/spark/spark_helpers.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -415,14 +415,14 @@ def _spark_filter_contract_errors(
415415
st.StructField("RecordIndex", st.IntegerType()),
416416
st.StructField("FailureType", st.StringType()),
417417
st.StructField("Status", st.StringType()),
418-
st.StructField("Entity", st.StringType()),
418+
st.StructField("OriginalEntity", st.StringType()),
419419
]
420420
),
421421
)
422422
.filter(
423423
(sf.col("FailureType") == sf.lit("record"))
424424
& (sf.col("Status") != sf.lit("informational"))
425-
& (sf.col("Entity") == sf.lit(entity_name))
425+
& (sf.col("OriginalEntity") == sf.lit(entity_name))
426426
)
427427
.distinct()
428428
.orderBy(sf.asc(sf.col("RecordIndex")))

‎src/dve/core_engine/message.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -90,7 +90,7 @@ def extract_error_value(records, error_location):
9090
class UserMessage:
9191
"""The structure of the message that is used to populate the error report."""
9292

93-
Entity: Optional[str]
93+
ReportingEntity: Optional[str]
9494
"""The entity that the message pertains to (if applicable)."""
9595
Key: Optional[str]
9696
"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
176176
"""The category of the error."""
177177

178178
HEADER: ClassVar[list[str]] = [
179-
"Entity",
179+
"ReportingEntity",
180+
"OriginalEntity",
180181
"Key",
181182
"FailureType",
182183
"Status",
@@ -307,6 +308,7 @@ def to_row(
307308

308309
return (
309310
self.entity,
311+
self.original_entity,
310312
key,
311313
self.failure_type,
312314
"informational" if self.is_informational else "error",

‎src/dve/core_engine/type_hints.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -147,6 +147,7 @@
147147
"""The record index that the error relates to (if applicable)"""
148148

149149
MessageTuple = tuple[
150+
Optional[EntityName],
150151
Optional[EntityName],
151152
Key,
152153
FailureType,

‎src/dve/pipeline/pipeline.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -908,7 +908,7 @@ def _get_error_dataframes(self, submission_id: str):
908908
.alias("error_type") # type: ignore
909909
)
910910
df = df.select(
911-
pl.col("Entity").alias("Table"), # type: ignore
911+
pl.col("ReportingEntity").alias("Table"), # type: ignore
912912
pl.col("error_type").alias("Type"), # type: ignore
913913
pl.col("ErrorCode").alias("Error_Code"), # type: ignore
914914
pl.col("ReportingField").alias("Data_Item"), # type: ignore

‎src/dve/reporting/error_report.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,7 @@ def create_error_dataframe(errors: deque[FeedbackMessage], key_fields):
9191
.alias("error_type")
9292
)
9393
df = df.select( # type: ignore
94-
col("Entity").alias("Table"), # type: ignore
94+
col("ReportingEntity").alias("Table"), # type: ignore
9595
col("error_type").alias("Type"), # type: ignore
9696
col("ErrorCode").alias("Error_Code"), # type: ignore
9797
col("ReportingField").alias("Data_Item"), # type: ignore

‎tests/features/movies.feature‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ Feature: Pipeline tests using the movies dataset
2222
Then there is 1 submission rejection from the data_contract phase
2323
And there are 3 record rejections from the data_contract phase
2424
And there are errors with the following details and associated error_count from the data_contract phase
25-
| Entity | ErrorCode | ErrorMessage | RecordIndex | error_count |
25+
| ReportingEntity | ErrorCode | ErrorMessage | RecordIndex | error_count |
2626
| movies | BLANKYEAR | year not provided | 2 | 1 |
2727
| movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 |
2828
| movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 |
@@ -61,10 +61,10 @@ Feature: Pipeline tests using the movies dataset
6161
Then there is 1 submission rejection from the data_contract phase
6262
And there are 3 record rejections from the data_contract phase
6363
And there are errors with the following details and associated error_count from the data_contract phase
64-
| Entity | ErrorCode | ErrorMessage | RecordIndex | error_count |
65-
| movies | BLANKYEAR | year not provided | 2 | 1 |
66-
| movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 |
67-
| movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 |
64+
| ReportingEntity | ErrorCode | ErrorMessage | RecordIndex | error_count |
65+
| movies | BLANKYEAR | year not provided | 2 | 1 |
66+
| movies_rename_test | DODGYYEAR | year value (NOT_A_NUMBER) is invalid | 1 | 1 |
67+
| movies | DODGYDATE | date_joined value is not valid: daft_date | 1 | 1 |
6868
| movies | BLANKTITLE | title should not be blank | 4 | 1 |
6969
And the movies entity is stored as a parquet after the data_contract phase
7070
And the latest audit record for the submission is marked with processing status business_rules

‎tests/features/steps/steps_pipeline.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -284,7 +284,7 @@ def check_rows_removed_with_error_code(context: Context, entity_name: str, error
284284
err_df = get_all_errors_df(context)
285285

286286
recs_with_err_code = err_df.filter(
287-
(pl.col("Entity").eq(entity_name)) & (pl.col("ErrorCode").eq(error_code))
287+
(pl.col("ReportingEntity").eq(entity_name)) & (pl.col("ErrorCode").eq(error_code))
288288
).shape[0]
289289
assert recs_with_err_code >= 1
290290

‎tests/features/steps/utilities.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
from dve.parser.type_hints import URI
1616

1717
ERROR_DF_FIELDS: List[str] = [
18-
"Entity",
18+
"ReportingEntity",
1919
"Key",
2020
"ErrorCode",
2121
"FailureType",

‎tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_data_contract.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -381,4 +381,4 @@ def test_duckdb_data_contract_custom_error_details(nested_all_string_parquet_w_e
381381
assert messages[0].ErrorMessage == "subfield id is invalid: subfield.id - WRONG"
382382
assert messages[1].ErrorCode == "TESTIDBAD"
383383
assert messages[1].ErrorMessage == "id is invalid: id - WRONG"
384-
assert messages[1].Entity == "test_rename"
384+
assert messages[1].ReportingEntity == "test_rename"

0 commit comments

Comments
 (0)