Skip to content

Commit b245abd

Browse files
committed
fix: ensure that original entity detials persisted in feedback messages for further processing
1 parent 3eca2f1 commit b245abd

16 files changed

Lines changed: 35 additions & 34 deletions

File tree

‎src/dve/core_engine/backends/base/rules.py‎

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -448,8 +448,7 @@ def process_node(node: HierarchyNode):
448448
record=record, # type: ignore
449449
error_location=location,
450450
error_message=template_object(
451-
node.missing_parent_id_error_message,
452-
record
451+
node.missing_parent_id_error_message, record
453452
),
454453
failure_type="record",
455454
error_type="record",
@@ -522,10 +521,7 @@ def process_node(node: HierarchyNode) -> bool:
522521
entity=node.parent_entity,
523522
record=record, # type: ignore
524523
error_location=location,
525-
error_message=template_object(
526-
node.no_valid_records_error_message,
527-
record
528-
),
524+
error_message=template_object(node.no_valid_records_error_message, record),
529525
failure_type="record",
530526
error_type="record",
531527
error_code=node.no_valid_records_error_code,

‎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
@@ -141,6 +141,7 @@
141141
"""The record index that the error relates to (if applicable)"""
142142

143143
MessageTuple = tuple[
144+
Optional[EntityName],
144145
Optional[EntityName],
145146
Key,
146147
FailureType,

‎src/dve/pipeline/pipeline.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -865,7 +865,7 @@ def _get_error_dataframes(self, submission_id: str):
865865
.alias("error_type") # type: ignore
866866
)
867867
df = df.select(
868-
pl.col("Entity").alias("Table"), # type: ignore
868+
pl.col("ReportingEntity").alias("Table"), # type: ignore
869869
pl.col("error_type").alias("Type"), # type: ignore
870870
pl.col("ErrorCode").alias("Error_Code"), # type: ignore
871871
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",

0 commit comments

Comments
 (0)