Skip to content

Commit afba07e

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 4a0288d + 516e5e4 commit afba07e

14 files changed

Lines changed: 208 additions & 91 deletions

File tree

‎docs/advanced_guidance/json_schemas/entity_relationships.schema.json‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,10 @@
2020
"mandatory": {
2121
"type": "boolean"
2222
},
23-
"orphaned_records_error_code": {
23+
"missing_parent_id_error_code": {
2424
"type": "string"
2525
},
26-
"orphaned_records_error_message": {
26+
"missing_parent_id_error_message": {
2727
"type": "string"
2828
}
2929
},

‎src/dve/common/error_utils.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@
55
import logging
66
from collections.abc import Iterable
77
from itertools import chain
8-
from multiprocessing import Queue
8+
from queue import Queue
99
from threading import Thread
1010
from typing import Optional, Union
1111

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

Lines changed: 25 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -46,11 +46,7 @@
4646
TableUnion,
4747
)
4848
from dve.core_engine.backends.types import Entities, EntityType, StageSuccessful
49-
from dve.core_engine.configuration.v1.hierarchy import (
50-
ChildHierarchyNode,
51-
EntityHierarchy,
52-
HierarchyNode,
53-
)
49+
from dve.core_engine.configuration.v1.hierarchy import EntityHierarchy, HierarchyNode
5450
from dve.core_engine.constants import ORPHANED_RECORD_ENTITY_NAME
5551
from dve.core_engine.exceptions import CriticalProcessingError
5652
from dve.core_engine.loggers import get_logger
@@ -399,7 +395,7 @@ def identify_and_remove_orphans(
399395
"""
400396

401397
def process_node(
402-
node: HierarchyNode | ChildHierarchyNode,
398+
node: HierarchyNode,
403399
parent_entity_name: Optional[EntityName],
404400
orph_messages: Messages | None = None,
405401
):
@@ -409,7 +405,7 @@ def process_node(
409405
if orph_messages is None:
410406
orph_messages = []
411407

412-
if isinstance(node, ChildHierarchyNode) and parent_entity_name is not None:
408+
if parent_entity_name is not None:
413409
self.logger.info(f"Identifying orphans in {current_entity_name}")
414410

415411
join_expr = " AND ".join(
@@ -428,7 +424,9 @@ def process_node(
428424
)
429425

430426
if no_orphs > 0:
431-
self.logger.info(f"Removing orphan records from {current_entity_name}")
427+
self.logger.info(
428+
f"Removing records with missing parent from {current_entity_name}"
429+
)
432430
location = list(node.join_fields.values())[0]
433431
with BackgroundMessageWriter(
434432
working_directory=working_directory,
@@ -442,28 +440,29 @@ def process_node(
442440
entity_name=current_entity_name,
443441
reporting=ReportingConfig(
444442
emit="record_failure",
445-
code=node.orphaned_records_error_code,
446-
message=node.orphaned_records_error_message,
443+
code=node.missing_parent_id_error_code,
444+
message=node.missing_parent_id_error_message,
447445
location=location,
448446
),
449447
),
450448
)
451-
for record in _orph_records:
452-
msg_writer.write_queue.put(
453-
[
454-
FeedbackMessage(
455-
entity=current_entity_name,
456-
record=record, # type: ignore
457-
error_location=location,
458-
error_message=node.orphaned_records_error_message,
459-
failure_type="record",
460-
error_type="record",
461-
error_code=node.orphaned_records_error_code,
462-
reporting_field=location,
463-
category="Parent Missing",
464-
)
465-
]
466-
)
449+
# moved to batch the write - risky if large number of
450+
msg_writer.write_queue.put(
451+
[
452+
FeedbackMessage(
453+
entity=current_entity_name,
454+
record=record, # type: ignore
455+
error_location=location,
456+
error_message=node.missing_parent_id_error_message,
457+
failure_type="record",
458+
error_type="record",
459+
error_code=node.missing_parent_id_error_code,
460+
reporting_field=location,
461+
category="Parent Missing",
462+
)
463+
for record in _orph_records
464+
]
465+
)
467466

468467
if node.children:
469468
for child_node in node.children:

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

Lines changed: 32 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,8 @@
33
import json
44
from typing import Any, Optional, Type, Union
55

6-
from pydantic import BaseModel, Field, PrivateAttr, validate_call
6+
from pydantic import BaseModel, Field, PrivateAttr, field_validator, model_validator, validate_call
7+
from pydantic_core.core_schema import FieldValidationInfo
78
from typing_extensions import Literal
89

910
from dve.core_engine.backends.base.reference_data import ReferenceConfig, ReferenceConfigUnion
@@ -93,23 +94,48 @@ class _TypeAliasDefinition(_BaseTypeDefintion):
9394
class _LinkageConfig(BaseModel):
9495
"""Specify how to link entities back to parents if required"""
9596

96-
parent_entity: EntityName
97+
parent_entity: Optional[EntityName] = None
9798
"""The name of the parent entity"""
98-
join_fields: JoinFields
99+
join_fields: JoinFields = Field(default_factory=dict)
99100
"""The fields that can be used to link back to the parent entity"""
100-
mandatory: Optional[bool] = False
101+
is_root_entity: bool = False
102+
"""Whether the entity is the highest level parent in a tree"""
103+
mandatory: bool = False
101104
"""If the entity is a child, is it a mandatory field of the parent"""
102105
no_valid_records_error_code: Optional[ErrorCode] = "NoValidRecords"
103106
"""The error code to emit if the entity has no valid records and is mandatory in the parent entity""" # pylint: disable=C0301
104107
no_valid_records_error_message: Optional[ErrorMessage] = (
105108
"parent record removed as no valid child records"
106109
)
107110
"""The error message to emit if the entity has no valid records and is mandatory in the parent entity""" # pylint: disable=C0301
108-
orphaned_records_error_code: Optional[ErrorCode] = "OrphanedRecords"
111+
missing_parent_id_error_code: Optional[ErrorCode] = "MissingParentRecord"
109112
"""The error code to emit if the entity contains records that are orphaned by parent record rejections""" # pylint: disable=C0301
110-
orphaned_records_error_message: Optional[ErrorMessage] = "Orphaned records removed"
113+
missing_parent_id_error_message: Optional[ErrorMessage] = (
114+
"Records removed due to no valid parent record"
115+
)
111116
"""The error code to emit if the entity contains records that are orphaned by parent record rejections""" # pylint: disable=C0301
112117

118+
@model_validator(mode="after")
119+
def _check_root_no_parent_or_join_keys(self):
120+
if self.is_root_entity and (self.parent_entity or self.join_fields):
121+
raise ValueError(
122+
"If entity is root, neither parent_entity nor join keys should be specified"
123+
)
124+
return self
125+
126+
@model_validator(mode="after")
127+
def _check_root_mandatory(self):
128+
if self.is_root_entity and not self.mandatory:
129+
raise ValueError("If entity is root, it must be labelled mandatory")
130+
return self
131+
132+
@model_validator(mode="after")
133+
def _check_parent_entity_with_join_keys(self):
134+
if self.parent_entity or self.join_fields:
135+
if not (self.parent_entity and self.join_fields):
136+
raise ValueError("Both parent_entity and join_fields must be supplied if one is")
137+
return self
138+
113139

114140
class _SchemaConfig(BaseModel):
115141
"""Configuration for a component schema within a dataset."""

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

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

1818
entity_name: str
19-
children: Optional[list["ChildHierarchyNode"]] = Field(default_factory=list)
19+
children: list["HierarchyNode"] = Field(default_factory=list)
20+
mandatory: bool = False
21+
join_fields: dict[str, str] = Field(default_factory=dict)
22+
no_valid_records_error_code: ErrorCode = "NoValidRecords"
23+
no_valid_records_error_message: ErrorMessage = "parent record removed as no valid child records"
24+
missing_parent_id_error_code: Optional[ErrorCode] = "MissingParentRecord"
25+
missing_parent_id_error_message: Optional[ErrorMessage] = (
26+
"Records removed due to no valid parent record"
27+
)
2028

2129
def get_descendents(self) -> list[str]:
2230
"""Recursively list all descendents of the node"""
@@ -58,19 +66,6 @@ def as_dict(self) -> dict[str, dict[str, Any]]:
5866
return {self.entity_name: ret_dict}
5967

6068

61-
class ChildHierarchyNode(HierarchyNode):
62-
"""Stores child entity hierarchy information"""
63-
64-
join_fields: dict[str, str]
65-
mandatory: Optional[bool] = False
66-
no_valid_records_error_code: Optional[ErrorCode] = "NoValidRecords"
67-
no_valid_records_error_message: Optional[ErrorMessage] = (
68-
"parent record removed as no valid child records"
69-
)
70-
orphaned_records_error_code: Optional[ErrorCode] = "OrphanedRecords"
71-
orphaned_records_error_message: Optional[ErrorMessage] = "Orphaned records removed"
72-
73-
7469
class EntityHierarchy:
7570
"""Determines and stores entity hierarchy information from config"""
7671

@@ -82,12 +77,35 @@ def determine_trees(
8277
all_datasets: Iterable[str], entity_relationships: dict[str, _LinkageConfig]
8378
) -> dict[EntityName, HierarchyNode]:
8479
"""Determine the entity hierarchy trees and store as HierarchyNodes"""
80+
root_entities: dict[str, _LinkageConfig] = dict(
81+
filter(lambda x: x[1].is_root_entity, entity_relationships.items())
82+
)
8583
top_level_parents: dict[EntityName, HierarchyNode] = {
86-
entity_name: HierarchyNode(entity_name=entity_name)
87-
for entity_name in all_datasets
88-
if entity_name not in entity_relationships
84+
entity_name: HierarchyNode(
85+
entity_name=entity_name,
86+
**config.model_dump(
87+
exclude={
88+
"parent_entity",
89+
"missing_parent_id_error_code",
90+
"missing_parent_id_error_message",
91+
}
92+
),
93+
missing_parent_id_error_code=None,
94+
missing_parent_id_error_message=None,
95+
)
96+
for entity_name, config in root_entities.items()
8997
}
9098

99+
if default_roots := [
100+
entity_name for entity_name in all_datasets if entity_name not in entity_relationships
101+
]:
102+
for entity_name in default_roots:
103+
top_level_parents[entity_name] = HierarchyNode(
104+
entity_name=entity_name,
105+
missing_parent_id_error_code=None,
106+
missing_parent_id_error_message=None,
107+
)
108+
91109
for name, linkage_detail in entity_relationships.items():
92110
for main_entity, parent_node in top_level_parents.items():
93111
if (
@@ -96,7 +114,7 @@ def determine_trees(
96114
):
97115
parent_node.add_child_node(
98116
linkage_detail.parent_entity,
99-
ChildHierarchyNode(
117+
HierarchyNode(
100118
entity_name=name, **linkage_detail.model_dump(exclude={"parent_entity"})
101119
),
102120
)

‎src/dve/core_engine/constants.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,5 +7,5 @@
77
"""The name of the field that can be used to extract the field value that caused
88
a pydantic validation error"""
99

10-
ORPHANED_RECORD_ENTITY_NAME: str = "orphaned_records_tracker"
11-
"""Name to keep track of identified orphaned records"""
10+
ORPHANED_RECORD_ENTITY_NAME: str = "orphaned_record_tracker"
11+
"""Name of entity to keep track of records where there is a missing parent record"""

‎src/dve/core_engine/models.py‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,9 @@ def _ensure_just_file_stem(
8282
@property
8383
def file_name_with_ext(self):
8484
"""Return file name with extension."""
85-
return f"{self.file_name}.{self.file_extension}"
85+
if self.file_extension:
86+
return f"{self.file_name}.{self.file_extension}"
87+
return self.file_name
8688

8789
@classmethod
8890
def from_metadata_file(cls, submission_id: str, metadata_uri: Location):

‎src/dve/pipeline/pipeline.py‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -216,10 +216,10 @@ def write_file_to_parquet(
216216

217217
for model_name, model in models.items():
218218
self._logger.info(f"Transforming {model_name} to stringified parquet")
219-
reader: BaseFileReader = load_reader(
220-
dataset, model_name, ext, self.backend_reader_kwargs
221-
)
222219
try:
220+
reader: BaseFileReader = load_reader(
221+
dataset, model_name, ext, self.backend_reader_kwargs
222+
)
223223
if not entity_type:
224224
reader.write_parquet(
225225
reader.read_to_py_iterator(
@@ -242,6 +242,7 @@ def write_file_to_parquet(
242242
f"{out}{model_name}",
243243
)
244244
except MessageBearingError as exc:
245+
self._logger.error(f"Unable to process {model_name}", exc_info=exc)
245246
errors.extend(exc.messages)
246247

247248
return list(dict.fromkeys(errors)) # remove any duplicate errors

‎src/dve/pipeline/utils.py‎

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,11 @@
1111
import dve.core_engine.backends.implementations.duckdb # pylint: disable=unused-import
1212
import dve.core_engine.backends.implementations.spark # pylint: disable=unused-import
1313
import dve.parser.file_handling as fh
14+
from dve.core_engine.backends.exceptions import MessageBearingError
1415
from dve.core_engine.backends.readers import _READER_REGISTRY
1516
from dve.core_engine.configuration.v1 import SchemaName, V1EngineConfig, _ModelConfig
1617
from dve.core_engine.loggers import get_logger
18+
from dve.core_engine.message import FeedbackMessage
1719
from dve.core_engine.type_hints import URI, SubmissionResult
1820
from dve.metadata_parser.model_generator import JSONtoPyd
1921

@@ -52,7 +54,31 @@ def load_reader(
5254
backend_reader_kwargs: Optional[dict[str, Any]] = None,
5355
):
5456
"""Loads the readers for the diven feed, model name and file extension"""
55-
reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"]
57+
try:
58+
reader_config = dataset[model_name].reader_config[f".{file_extension.lower()}"]
59+
except KeyError as exc:
60+
if file_extension:
61+
err_msg = (
62+
f"The supplied file extension `{file_extension}`"
63+
+f" is not a supported file format for {model_name}."
64+
)
65+
else:
66+
err_msg = "No supplied file extension. Unable to parse file without a file extension."
67+
68+
raise MessageBearingError(
69+
f"The file extension provided ({file_extension}) is not supported for this collection.",
70+
messages=[
71+
FeedbackMessage(
72+
entity=model_name,
73+
record=None,
74+
failure_type="submission",
75+
error_location="Whole File",
76+
error_code="InvalidFileExtension",
77+
error_message=err_msg,
78+
)
79+
],
80+
) from exc
81+
5682
reader = _READER_REGISTRY[reader_config.reader](
5783
**reader_config.kwargs_, **backend_reader_kwargs if backend_reader_kwargs else {}
5884
)

‎tests/features/planets.feature‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,9 @@ Feature: Pipeline tests using the planets dataset
4343
And I add initial audit entries for the submission
4444
Then the latest audit record for the submission is marked with processing status file_transformation
4545
When I run the file transformation phase
46-
Then the latest audit record for the submission is marked with processing status failed
46+
Then the latest audit record for the submission is marked with processing status error_report
47+
When I run the error report phase
48+
Then An error report is produced
4749

4850
Scenario: Handle a file with duplicated extension provided (spark)
4951
Given I submit the planets file planets.csv.csv for processing

0 commit comments

Comments
 (0)