Skip to content

Commit 434ca4a

Browse files
Merge branch 'release_v010' of https://github.com/NHSDigital/data-validation-engine into feature/gr-ndsp-619-add_group_level_rejections
2 parents 953f1a3 + a6b6154 commit 434ca4a

16 files changed

Lines changed: 1245 additions & 82 deletions

File tree

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

Lines changed: 17 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -396,33 +396,35 @@ def identify_and_remove_orphans(
396396

397397
def process_node(
398398
node: HierarchyNode,
399-
parent_entity_name: Optional[EntityName],
399+
orph_messages: Messages | None = None,
400400
processed: bool = False,
401-
) -> bool:
402-
"""Recursive helper to process a node and its children."""
403-
current_entity_name = node.entity_name
401+
):
402+
"""Identify orphans and remove in a given node"""
404403

405-
if parent_entity_name is not None:
406-
self.logger.info(f"Identifying orphans in {current_entity_name}")
404+
if orph_messages is None:
405+
orph_messages = []
406+
407+
if node.parent_entity is not None:
408+
self.logger.info(f"Identifying orphans in {node.entity_name}")
407409

408410
join_expr = " AND ".join(
409-
f"{parent_entity_name}.{k} = {current_entity_name}.{v}"
411+
f"{node.parent_entity}.{k} = {node.entity_name}.{v}"
410412
for k, v in node.join_fields.items()
411413
)
412414

413415
_, no_orphs = self.identify_orphans(
414416
entities=entities,
415417
config=OrphanIdentification(
416418
id=list(node.join_fields.values())[0],
417-
entity_name=current_entity_name,
418-
target_name=parent_entity_name,
419+
entity_name=node.entity_name,
420+
target_name=node.parent_entity,
419421
join_condition=join_expr,
420422
),
421423
)
422424

423425
if no_orphs > 0:
424426
self.logger.info(
425-
f"Removing records with missing parent from {current_entity_name}"
427+
f"Removing records with missing parent from {node.entity_name}"
426428
)
427429
processed = True
428430
location = list(node.join_fields.values())[0]
@@ -435,7 +437,7 @@ def process_node(
435437
_orph_records = self.remove_orphans(
436438
entities=entities,
437439
config=OrphanRemoval(
438-
entity_name=current_entity_name,
440+
entity_name=node.entity_name,
439441
reporting=ReportingConfig(
440442
emit="record_failure",
441443
code=node.missing_parent_id_error_code,
@@ -448,7 +450,7 @@ def process_node(
448450
msg_writer.write_queue.put(
449451
[
450452
FeedbackMessage(
451-
entity=current_entity_name,
453+
entity=node.entity_name,
452454
record=record, # type: ignore
453455
error_location=location,
454456
error_message=node.missing_parent_id_error_message,
@@ -462,16 +464,13 @@ def process_node(
462464
]
463465
)
464466

465-
if node.children:
466-
for child_node in node.children:
467-
processed = process_node(child_node, current_entity_name, processed)
468-
469467
return processed
470468

471469
processed = False
472470

473-
for root_node in entity_hierarchy.entity_trees.values():
474-
processed = process_node(root_node, parent_entity_name=None, processed=processed)
471+
for tree in entity_hierarchy.entity_trees.values():
472+
for node in tree.iterate_root_down():
473+
processed = process_node(node)
475474

476475
_orph_rel = entities.get(ORPHANED_RECORD_ENTITY_NAME)
477476
if _orph_rel is not None:

‎src/dve/core_engine/configuration/v1/hierarchy.py‎

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ class HierarchyNode(BaseModel):
1616
"""Stores entity hierarchy information"""
1717

1818
entity_name: str
19+
parent_entity: Optional[str] = None
1920
children: list["HierarchyNode"] = Field(default_factory=list)
2021
mandatory: bool = False
2122
join_fields: dict[str, str] = Field(default_factory=dict)
@@ -26,14 +27,18 @@ class HierarchyNode(BaseModel):
2627
"Records removed due to no valid parent record"
2728
)
2829

29-
def get_descendents(self) -> list[str]:
30+
def get_descendents(self) -> list["HierarchyNode"]:
3031
"""Recursively list all descendents of the node"""
3132
descendents = []
3233
for node in self.children: # type: ignore
33-
descendents.append(node.entity_name)
34+
descendents.append(node)
3435
descendents.extend(node.get_descendents())
3536
return descendents
3637

38+
def get_descendent_names(self) -> list[str]:
39+
"""Recursively list all names of descendents of the node"""
40+
return [node.entity_name for node in self.get_descendents()]
41+
3742
def get_node(self, entity_name: str) -> Union["HierarchyNode", None]:
3843
"""Recursively search for node and return if found"""
3944
node = None
@@ -65,6 +70,20 @@ def as_dict(self) -> dict[str, dict[str, Any]]:
6570

6671
return {self.entity_name: ret_dict}
6772

73+
def _get_full_tree(self):
74+
"""Get all nodes in tree, including the root"""
75+
desc = self.get_descendents()
76+
desc.insert(0, self)
77+
return desc
78+
79+
def iterate_root_down(self):
80+
"""Iterate through nodes from root to lowest descendent"""
81+
yield from self._get_full_tree()
82+
83+
def iterate_lowest_descendent_up(self):
84+
"""Iterate through nodes from lowest descendent to root"""
85+
yield from self._get_full_tree()[::-1]
86+
6887

6988
class EntityHierarchy:
7089
"""Determines and stores entity hierarchy information from config"""
@@ -83,6 +102,7 @@ def determine_trees(
83102
top_level_parents: dict[EntityName, HierarchyNode] = {
84103
entity_name: HierarchyNode(
85104
entity_name=entity_name,
105+
parent_entity=None,
86106
**config.model_dump(
87107
exclude={
88108
"parent_entity",
@@ -102,6 +122,7 @@ def determine_trees(
102122
for entity_name in default_roots:
103123
top_level_parents[entity_name] = HierarchyNode(
104124
entity_name=entity_name,
125+
parent_entity=None,
105126
missing_parent_id_error_code=None,
106127
missing_parent_id_error_message=None,
107128
)
@@ -110,13 +131,11 @@ def determine_trees(
110131
for main_entity, parent_node in top_level_parents.items():
111132
if (
112133
linkage_detail.parent_entity == main_entity
113-
or linkage_detail.parent_entity in parent_node.get_descendents()
134+
or linkage_detail.parent_entity in parent_node.get_descendent_names()
114135
):
115136
parent_node.add_child_node(
116137
linkage_detail.parent_entity,
117-
HierarchyNode(
118-
entity_name=name, **linkage_detail.model_dump(exclude={"parent_entity"})
119-
),
138+
HierarchyNode(entity_name=name, **linkage_detail.model_dump()),
120139
)
121140
break
122141
else:

‎src/dve/pipeline/pipeline.py‎

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -886,14 +886,18 @@ def error_report(
886886
.agg(pl.col("Count").sum()) # type: ignore
887887
.iter_rows(named=True)
888888
}
889+
submission_rejections = err_types.get(
890+
ErrorReportCategories.FILE_REJECTION.reporting_name, 0
891+
)
889892
sub_stats = SubmissionStatisticsRecord(
890893
submission_id=submission_info.submission_id,
891894
record_count=submission_status.number_of_records,
892-
number_submission_rejections=err_types.get(
893-
ErrorReportCategories.FILE_REJECTION.reporting_name, 0
894-
),
895-
number_record_rejections=err_types.get(
896-
ErrorReportCategories.RECORD_REJECTION.reporting_name, 0
895+
number_submission_rejections=submission_rejections,
896+
number_record_rejections=(
897+
submission_status.number_of_records
898+
if submission_rejections > 0 else err_types.get( # type: ignore
899+
ErrorReportCategories.RECORD_REJECTION.reporting_name, 0
900+
)
897901
),
898902
number_warnings=err_types.get(ErrorReportCategories.WARNING.reporting_name, 0),
899903
)

‎src/dve/pipeline/utils.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,7 @@ def load_reader(
6060
if file_extension:
6161
err_msg = (
6262
f"The supplied file extension `{file_extension}`"
63-
+f" is not a supported file format for {model_name}."
63+
+ f" is not a supported file format for {model_name}."
6464
)
6565
else:
6666
err_msg = "No supplied file extension. Unable to parse file without a file extension."

‎src/dve/reporting/excel_report.py‎

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -153,17 +153,21 @@ def _add_submission_info(self, status: str, summary: Worksheet):
153153
), # pylint: disable=C0301
154154
]
155155
)
156-
if status not in (
156+
if status in (
157157
ErrorReportStatus.PROCESSING_FAILED,
158158
ErrorReportStatus.FILE_REJECTION,
159159
):
160-
summary.append(
161-
[
162-
"",
163-
"Total Number of Records Rejected",
164-
self.submission_status.number_of_records_rejected,
165-
]
166-
)
160+
_records_rejected = self.submission_status.number_of_records
161+
else:
162+
_records_rejected = self.submission_status.number_of_records_rejected
163+
summary.append(
164+
[
165+
"",
166+
"Total Number of Records Rejected",
167+
_records_rejected,
168+
]
169+
)
170+
167171
summary.append(["", ""])
168172

169173

‎tests/features/animals.feature‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -55,5 +55,5 @@ Feature: Pipeline tests using the animal dataset
5555
| parameter | value |
5656
| record_count | 7 |
5757
| number_submission_rejections | 1 |
58-
| number_record_rejections | 2 |
58+
| number_record_rejections | 7 |
5959
| number_warnings | 1 |

‎tests/features/flights.feature‎

Lines changed: 14 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ Feature: Pipeline tests using the flights dataset
1010
When I run the file transformation phase
1111
Then the country entity is stored as a parquet after the file_transformation phase
1212
And the airport entity is stored as a parquet after the file_transformation phase
13+
And the staff entity is stored as a parquet after the file_transformation phase
1314
And the flights entity is stored as a parquet after the file_transformation phase
1415
And the passengers entity is stored as a parquet after the file_transformation phase
1516
And the latest audit record for the submission is marked with processing status data_contract
@@ -36,19 +37,20 @@ Feature: Pipeline tests using the flights dataset
3637
When I run the file transformation phase
3738
Then the country entity is stored as a parquet after the file_transformation phase
3839
And the airport entity is stored as a parquet after the file_transformation phase
40+
And the staff entity is stored as a parquet after the file_transformation phase
3941
And the flights entity is stored as a parquet after the file_transformation phase
4042
And the passengers entity is stored as a parquet after the file_transformation phase
4143
And the latest audit record for the submission is marked with processing status data_contract
4244
When I run the data contract phase
4345
Then there are no file rejections from the data_contract phase
44-
And there are no record rejections from the data_contract phase
46+
And there is 1 record rejection from the data_contract phase
4547
When I run the business rules phase
4648
Then there are errors with the following details and associated error_count from the business_rules phase
4749
| ErrorType | ErrorCode | error_count |
48-
| record | C1 | 1 |
49-
| record | AG1 | 1 |
50-
| record | FG1 | 2 |
51-
| record | PG1 | 4 |
50+
| record | AG1 | 3 |
51+
| record | SG1 | 15 |
52+
| record | FG1 | 10 |
53+
| record | PG1 | 25 |
5254
When I run the error report phase
5355
Then An error report is produced
5456
# TODO - fix the stats calculations as they're currently incorrect for hiearchical datasets
@@ -66,6 +68,7 @@ Feature: Pipeline tests using the flights dataset
6668
When I run the file transformation phase
6769
Then the country entity is stored as a parquet after the file_transformation phase
6870
And the airport entity is stored as a parquet after the file_transformation phase
71+
And the staff entity is stored as a parquet after the file_transformation phase
6972
And the flights entity is stored as a parquet after the file_transformation phase
7073
And the passengers entity is stored as a parquet after the file_transformation phase
7174
And the latest audit record for the submission is marked with processing status data_contract
@@ -76,10 +79,10 @@ Feature: Pipeline tests using the flights dataset
7679
Then there are errors with the following details and associated error_count from the business_rules phase
7780
| ErrorType | ErrorCode | error_count |
7881
| record | F1 | 1 |
79-
| record | PG1 | 2 |
82+
| record | PG1 | 3 |
8083
When I run the error report phase
8184
Then An error report is produced
82-
# TODO - fix the stats calculations as they're currently incorrect for hiearchical datasets
85+
#TODO - fix the stats calculations as they're currently incorrect for hiearchical datasets
8386
# And The statistics entry for the submission shows the following information
8487
# | parameter | value |
8588
# | record_count | 1 |
@@ -95,6 +98,7 @@ Feature: Pipeline tests using the flights dataset
9598
Then the country entity is stored as a parquet after the file_transformation phase
9699
And the airport entity is stored as a parquet after the file_transformation phase
97100
And the flights entity is stored as a parquet after the file_transformation phase
101+
And the staff entity is stored as a parquet after the file_transformation phase
98102
And the passengers entity is stored as a parquet after the file_transformation phase
99103
And the latest audit record for the submission is marked with processing status data_contract
100104
When I run the data contract phase
@@ -104,8 +108,9 @@ Feature: Pipeline tests using the flights dataset
104108
Then there are errors with the following details and associated error_count from the business_rules phase
105109
| ErrorType | ErrorCode | error_count |
106110
| record | F1 | 1 |
107-
| record | PG1 | 2 |
111+
| record | PG1 | 3 |
108112
| record | P1 | 1 |
113+
| record | S1 | 7 |
109114
When I run the error report phase
110115
Then An error report is produced
111116
# TODO - fix the stats calculations as they're currently incorrect for hiearchical datasets
@@ -161,7 +166,7 @@ Feature: Pipeline tests using the flights dataset
161166
| ErrorType | Status | ErrorCode | error_count |
162167
| record | error | F2 | 2 |
163168
| record | error | PG1 | 4 |
164-
| record | informational | A1 | 1 |
169+
# | record | informational | A1 | 1 |
165170
When I run the error report phase
166171
Then An error report is produced
167172
# TODO - fix the stats calculations as they're currently incorrect for hiearchical datasets

‎tests/features/movies.feature‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -42,7 +42,7 @@ Feature: Pipeline tests using the movies dataset
4242
| parameter | value |
4343
| record_count | 5 |
4444
| number_submission_rejections | 1 |
45-
| number_record_rejections | 3 |
45+
| number_record_rejections | 5 |
4646
| number_warnings | 2 |
4747
And the error aggregates are persisted
4848

@@ -81,7 +81,7 @@ Feature: Pipeline tests using the movies dataset
8181
| parameter | value |
8282
| record_count | 5 |
8383
| number_submission_rejections | 1 |
84-
| number_record_rejections | 3 |
84+
| number_record_rejections | 5 |
8585
| number_warnings | 2 |
8686
And the error aggregates are persisted
8787

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -656,9 +656,11 @@ def test_identify_and_remove_orphans(self):
656656
children=[
657657
HierarchyNode(
658658
entity_name="passengers",
659+
parent_entity="flights",
659660
children=[
660661
HierarchyNode(
661662
entity_name="food",
663+
parent_entity="passengers",
662664
children=[],
663665
join_fields={"passenger_id": "passenger_id"},
664666
mandatory=False

0 commit comments

Comments
 (0)