Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 2 additions & 6 deletions src/dve/core_engine/backends/base/rules.py
Original file line number Diff line number Diff line change
Expand Up @@ -411,7 +411,7 @@
_, no_orphs = self.identify_orphans(
entities=entities,
config=OrphanIdentification(
id=list(node.join_fields.values())[0],

Check warning on line 414 in src/dve/core_engine/backends/base/rules.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Replace "list(...)[0]" with "next(iter(...))" to avoid materializing the entire iterable.

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaDpziBZDSmfIj2RlwO2&open=AaDpziBZDSmfIj2RlwO2&pullRequest=168
entity_name=node.entity_name,
target_name=node.parent_entity,
join_condition=join_expr,
Expand All @@ -421,7 +421,7 @@
if no_orphs > 0:
self.logger.info(f"Removing records with missing parent from {node.entity_name}")
issues_found = True
location = list(node.join_fields.values())[0]

Check warning on line 424 in src/dve/core_engine/backends/base/rules.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Replace "list(...)[0]" with "next(iter(...))" to avoid materializing the entire iterable.

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaDpziBZDSmfIj2RlwO3&open=AaDpziBZDSmfIj2RlwO3&pullRequest=168
with BackgroundMessageWriter(
working_directory=working_directory,
dve_stage=self.__stage_name__,
Expand All @@ -448,8 +448,7 @@
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",
Expand Down Expand Up @@ -522,10 +521,7 @@
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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")))
Expand Down
6 changes: 4 additions & 2 deletions src/dve/core_engine/message.py
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@
class UserMessage:
"""The structure of the message that is used to populate the error report."""

Entity: Optional[str]
ReportingEntity: Optional[str]

Check warning on line 93 in src/dve/core_engine/message.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Rename this field "ReportingEntity" to match the regular expression ^[_a-z][_a-z0-9]*$.

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaDpziE5DSmfIj2RlwO4&open=AaDpziE5DSmfIj2RlwO4&pullRequest=168
"""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"
Expand Down Expand Up @@ -176,7 +176,8 @@
"""The category of the error."""

HEADER: ClassVar[list[str]] = [
"Entity",
"ReportingEntity",
"OriginalEntity",
"Key",
"FailureType",
"Status",
Expand Down Expand Up @@ -307,6 +308,7 @@

return (
self.entity,
self.original_entity,
key,
self.failure_type,
"informational" if self.is_informational else "error",
Expand Down
1 change: 1 addition & 0 deletions src/dve/core_engine/type_hints.py
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,7 @@
"""The record index that the error relates to (if applicable)"""

MessageTuple = tuple[
Optional[EntityName],
Optional[EntityName],
Key,
FailureType,
Expand Down
2 changes: 1 addition & 1 deletion src/dve/pipeline/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -545,7 +545,7 @@

return processed_files, failed_processing

def apply_business_rules( # pylint: disable=R0914,R0915

Check failure on line 548 in src/dve/pipeline/pipeline.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 17 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaDpziIMDSmfIj2RlwO5&open=AaDpziIMDSmfIj2RlwO5&pullRequest=168
self, submission_info: SubmissionInfo, submission_status: Optional[SubmissionStatus] = None
) -> tuple[SubmissionInfo, SubmissionStatus]:
"""Apply the business rules to a given submission, the submission may have failed at the
Expand Down Expand Up @@ -667,12 +667,12 @@

entity_issues: dict[EntityName, bool] = {
entity: any(
val
for val in (
orph_issues_1.get(entity, False),
grp_issues_1.get(entity, False),
orph_issues_2.get(entity, False),
)

Check warning on line 675 in src/dve/pipeline/pipeline.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Replace this comprehension with passing the iterable to the collection constructor call

See more on https://sonarcloud.io/project/issues?id=NHSDigital_data-validation-engine&issues=AaDpziIMDSmfIj2RlwO6&open=AaDpziIMDSmfIj2RlwO6&pullRequest=168
)
for entity in orph_issues_1.keys()
}
Expand Down Expand Up @@ -865,7 +865,7 @@
.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
Expand Down
2 changes: 1 addition & 1 deletion src/dve/reporting/error_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 5 additions & 5 deletions tests/features/movies.feature
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion tests/features/steps/steps_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion tests/features/steps/utilities.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
from dve.parser.type_hints import URI

ERROR_DF_FIELDS: List[str] = [
"Entity",
"ReportingEntity",
"Key",
"ErrorCode",
"FailureType",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"
assert messages[1].ReportingEntity == "test_rename"
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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"


Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
4 changes: 2 additions & 2 deletions tests/test_pipeline/pipeline_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand Down
10 changes: 5 additions & 5 deletions tests/test_pipeline/test_spark_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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",
Expand All @@ -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",
Expand Down Expand Up @@ -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",
Expand All @@ -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",
Expand Down
Loading