Skip to content

Commit cfbefc8

Browse files
refactor: add CriticalFeedbackMessage which exits the pipeline asap
number of feedback messages to reduce and removal of specific error code+messages for certain readers to become generic ft_error_code/message
1 parent bbaa9fc commit cfbefc8

16 files changed

Lines changed: 118 additions & 103 deletions

File tree

‎docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_csv_reader_args.schema.json‎

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,11 @@
44
"title": "Keyword Arguments used across all CSV readers",
55
"description": "Arguments present in all CSV readers available.",
66
"type": "object",
7+
"anyOf": [
8+
{
9+
"$ref": "global_reader_args.schema.json"
10+
}
11+
],
712
"properties": {
813
"field_check": {
914
"type": "string",
@@ -13,14 +18,6 @@
1318
"false"
1419
]
1520
},
16-
"field_check_error_code": {
17-
"type": "string",
18-
"description": "Allows you to customise the error code generated in the error report when the field check fails."
19-
},
20-
"field_check_error_message": {
21-
"type": "string",
22-
"description": "Allows you to customise the error message generated in the error report when the field check fails."
23-
},
2421
"null_empty_strings": {
2522
"type": "boolean",
2623
"description": "Converts empty string values into 'Null' values."
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
{
2+
"$schema": "https://json-schema.org/draft/2020-12/schema",
3+
"$id": "data-ingest:contract/components/reader_constraints/global_reader_args.schema.json",
4+
"title": "Keyword Arguments used across all readers",
5+
"description": "Arguments present in all readers available.",
6+
"type": "object",
7+
"properties": {
8+
"ft_error_code": {
9+
"type": "string",
10+
"description": "Error code to raise when a reader is unable to read the contents of a file.",
11+
"minLength": 1
12+
},
13+
"ft_error_message": {
14+
"type": "string",
15+
"description": "Error message to raise when a reader is unable to read the contents of a file.",
16+
"minLength": 1
17+
}
18+
}
19+
}

‎docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_xml_reader_args.schema.json‎

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,11 @@
44
"title": "Keyword Arguments used across all XML readers",
55
"description": "Arguments present in all XML readers.",
66
"type": "object",
7+
"anyOf": [
8+
{
9+
"$ref": "global_reader_args.schema.json"
10+
}
11+
],
712
"properties": {
813
"record_tag": {
914
"type": "string",
@@ -40,14 +45,6 @@
4045
"type": "string",
4146
"description": "The Relative path (to the dischema) for the XSD document."
4247
},
43-
"xsd_error_code": {
44-
"type": "string",
45-
"description": "Allows you to customise the error code generated in the error report when the XSD check fails."
46-
},
47-
"xsd_error_message": {
48-
"type": "string",
49-
"description": "Allows you to customise the error message generated in the error report when the XSD check fails."
50-
},
5148
"rules_location": {
5249
"type": "string",
5350
"description": "Allows you to prefix the directory which contains the XSD document."

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

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,10 @@ class BaseFileReader(ABC):
7373
decorated with the '@read_function' decorator, and is used in `read_entity_type`.
7474
7575
"""
76+
ft_error_code: Optional[str] = "MalformedFile"
77+
"""Default error code for when a file/submission cannot be parsed succesfully."""
78+
ft_error_message: Optional[str] = "The resource doesn't seem to be a valid text file"
79+
"""Default error message for when a file/submission cannot be parsed succesfully."""
7680

7781
def __init_subclass__(cls, *_, **__) -> None:
7882
"""When this class is subclassed, create and populate the `__read_methods__`
@@ -212,7 +216,11 @@ def _check_likely_text_file(resource: URI) -> bool:
212216
return False
213217
return True
214218

215-
def raise_if_not_sensible_file(self, resource: URI, entity_name: str):
219+
def raise_if_not_sensible_file(
220+
self,
221+
resource: URI,
222+
entity_name: str,
223+
):
216224
"""Sense check that the file is a text file. Raise error if doesn't
217225
appear to be the case."""
218226
if not self._check_likely_text_file(resource):
@@ -223,7 +231,7 @@ def raise_if_not_sensible_file(self, resource: URI, entity_name: str):
223231
record=None,
224232
failure_type="submission",
225233
error_location="Whole File",
226-
error_code="MalformedFile",
227-
error_message="The resource doesn't seem to be a valid text file",
234+
error_code=self.ft_error_code,
235+
error_message=self.ft_error_message,
228236
),
229237
)

‎src/dve/core_engine/backends/exceptions.py‎

Lines changed: 14 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -44,27 +44,25 @@ def __init__(self, *args: object, message: FeedbackMessage) -> None:
4444
self.message = message
4545
"""The message to be returned as part of the error."""
4646

47-
class UnableToParseCSVError(MessageBearingError):
47+
class UnableToParseCSVError(CriticalMessageBearingError):
4848
"""An error raised when unable to parse a CSV file"""
4949

5050
def __init__(
51-
self, entity_name: str, field_check_error_message: str, field_check_error_code: str
51+
self,
52+
entity_name: Optional[str],
53+
error_message: Optional[str],
54+
error_code: Optional[str],
5255
):
5356
super().__init__(
54-
messages=[
55-
FeedbackMessage(
56-
entity="csv_structure",
57-
record={
58-
entity_name: "Unable to parse file. Please check the structure of the file."
59-
},
60-
failure_type="submission",
61-
is_informational=False,
62-
error_type="csv read",
63-
error_location=entity_name,
64-
error_message=field_check_error_message,
65-
error_code=field_check_error_code,
66-
)
67-
]
57+
message=FeedbackMessage(
58+
entity=entity_name,
59+
record=None,
60+
failure_type="submission",
61+
is_informational=False,
62+
error_type="csv read",
63+
error_message=error_message or "Unable to parse the CSV file. Please check the structure of your CSV.", # pylint: disable=C0301
64+
error_code=error_code or "MalformedCSV",
65+
)
6866
)
6967

7068

‎src/dve/core_engine/backends/implementations/duckdb/readers/csv.py‎

Lines changed: 15 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -26,11 +26,12 @@
2626
duckdb_record_index,
2727
duckdb_write_parquet,
2828
get_duckdb_type_from_annotation,
29+
relation_is_empty,
2930
)
3031
from dve.core_engine.backends.implementations.duckdb.types import SQLType
3132
from dve.core_engine.backends.readers.csv import CSVFileReader
3233
from dve.core_engine.backends.utilities import get_polars_type_from_annotation, polars_record_index
33-
from dve.core_engine.constants import RECORD_INDEX_COLUMN_NAME
34+
from dve.core_engine.constants import PRE_VALIDATION_ENTITY, RECORD_INDEX_COLUMN_NAME
3435
from dve.core_engine.message import FeedbackMessage
3536
from dve.core_engine.type_hints import URI, EntityName
3637
from dve.parser.file_handling import get_content_length
@@ -44,10 +45,7 @@ class DuckDBCSVReader(CSVFileReader):
4445
to the file header, if it exists.
4546
4647
field_check: flag to compare submitted file header to the accompanying pydantic model
47-
field_check_error_code: The error code to provide if the file header doesn't contain
48-
the expected fields
49-
field_check_error_message: The error message to provide if the file header doesn't contain
50-
the expected fields"""
48+
"""
5149

5250
# TODO - the read_to_relation should include the schema and determine whether to
5351
# TODO - stringify or not
@@ -59,8 +57,8 @@ def __init__(
5957
quotechar: str = '"',
6058
connection: Optional[DuckDBPyConnection] = None,
6159
field_check: bool = False,
62-
field_check_error_code: str = "ExpectedVsActualFieldMismatch",
63-
field_check_error_message: str = "The submitted header is missing fields",
60+
ft_error_code: str = "ExpectedVsActualFieldMismatch",
61+
ft_error_message: str = "The submitted header is missing fields",
6462
null_empty_strings: bool = False,
6563
**_,
6664
):
@@ -72,8 +70,8 @@ def __init__(
7270
delimiter=delim,
7371
quote_char=quotechar,
7472
field_check=field_check,
75-
field_check_error_code=field_check_error_code,
76-
field_check_error_message=field_check_error_message,
73+
ft_error_code=ft_error_code,
74+
ft_error_message=ft_error_message,
7775
)
7876

7977
def read_to_py_iterator(
@@ -124,8 +122,8 @@ def read_to_relation( # pylint: disable=unused-argument
124122
except InvalidInputException as exc:
125123
raise UnableToParseCSVError(
126124
entity_name="csv_structure",
127-
field_check_error_message=self.field_check_error_message,
128-
field_check_error_code=self.field_check_error_code,
125+
error_code=self.ft_error_code,
126+
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
129127
) from exc
130128

131129
if self.null_empty_strings:
@@ -184,8 +182,8 @@ def read_to_relation( # pylint: disable=unused-argument
184182
except pl.exceptions.PolarsError as exc:
185183
raise UnableToParseCSVError(
186184
entity_name="csv_structure",
187-
field_check_error_message=self.field_check_error_message,
188-
field_check_error_code=self.field_check_error_code,
185+
error_code=self.ft_error_code,
186+
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
189187
) from exc
190188

191189
if self.null_empty_strings:
@@ -198,11 +196,11 @@ def read_to_relation( # pylint: disable=unused-argument
198196

199197
entity = self._connection.sql("SELECT * FROM df")
200198

201-
if entity.pl().shape[0] == 0:
199+
if relation_is_empty(entity):
202200
raise UnableToParseCSVError(
203201
entity_name="csv_structure",
204-
field_check_error_message=self.field_check_error_message,
205-
field_check_error_code=self.field_check_error_code,
202+
error_code=self.ft_error_code,
203+
error_message=self.ft_error_message or "Found zero records after loading CSV. File is likely malformed.", # pylint: disable=C0301
206204
)
207205

208206
return entity
@@ -273,7 +271,7 @@ def read_to_relation( # pylint: disable=unused-argument
273271
messages=[
274272
FeedbackMessage(
275273
record={entity_name: differing_values},
276-
entity="Pre-validation",
274+
entity=PRE_VALIDATION_ENTITY,
277275
failure_type="submission",
278276
error_message=(
279277
f"Found {no_records} distinct combination of header values."

‎src/dve/core_engine/backends/implementations/spark/readers/csv.py‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -40,8 +40,8 @@ def __init__(
4040
null_empty_strings: bool = False,
4141
spark_session: Optional[SparkSession] = None,
4242
field_check: bool = False,
43-
field_check_error_code: str = "ExpectedVsActualFieldMismatch",
44-
field_check_error_message: str = "The submitted header is missing fields",
43+
ft_error_code: str = "ExpectedVsActualFieldMismatch",
44+
ft_error_message: str = "The submitted header is missing fields",
4545
**_,
4646
) -> None:
4747

@@ -56,8 +56,8 @@ def __init__(
5656
quote_char=quote_char,
5757
header=header,
5858
field_check=field_check,
59-
field_check_error_code=field_check_error_code,
60-
field_check_error_message=field_check_error_message,
59+
ft_error_code=ft_error_code,
60+
ft_error_message=ft_error_message,
6161
)
6262

6363
def read_to_py_iterator(

‎src/dve/core_engine/backends/implementations/spark/readers/xml.py‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -79,8 +79,8 @@ def __init__(
7979
namespace=None,
8080
trim_cells=True,
8181
xsd_location: Optional[URI] = None,
82-
xsd_error_code: Optional[str] = None,
83-
xsd_error_message: Optional[str] = None,
82+
ft_error_code: Optional[str] = None,
83+
ft_error_message: Optional[str] = None,
8484
rules_location: Optional[URI] = None,
8585
**_,
8686
) -> None:
@@ -92,8 +92,8 @@ def __init__(
9292
null_values=null_values,
9393
sanitise_multiline=sanitise_multiline,
9494
xsd_location=xsd_location,
95-
xsd_error_code=xsd_error_code,
96-
xsd_error_message=xsd_error_message,
95+
ft_error_code=ft_error_code,
96+
ft_error_message=ft_error_message,
9797
rules_location=rules_location,
9898
)
9999

‎src/dve/core_engine/backends/readers/csv.py‎

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -42,8 +42,8 @@ def __init__(
4242
null_values: Collection[str] = frozenset({"NULL", "null", ""}),
4343
encoding: str = "utf-8-sig",
4444
field_check: bool = False,
45-
field_check_error_code: str = "CSVFieldMismatch",
46-
field_check_error_message: str = "The submitted header is invalid",
45+
ft_error_code: Optional[str] = "MalformedCSVFile",
46+
ft_error_message: Optional[str] = None,
4747
**_,
4848
):
4949
"""Init function for the base CSV reader.
@@ -89,10 +89,8 @@ def __init__(
8989
"""Encoding of the CSV file."""
9090
self.field_check = field_check
9191
"""Whether to check the fields are correct in the supplied header or not"""
92-
self.field_check_error_code = field_check_error_code
93-
"""Error code to raise when fields are missing or unexpected"""
94-
self.field_check_error_message = field_check_error_message
95-
"""Error message to raise when fields are missing or unexpected"""
92+
self.ft_error_code = ft_error_code
93+
self.ft_error_message = ft_error_message
9694

9795
def _get_reader_args(self) -> dict[str, Any]:
9896
reader_args: dict[str, Any] = {
@@ -218,8 +216,8 @@ def perform_field_check(
218216
entity_name,
219217
expected_schema,
220218
all_model_fields,
221-
self.field_check_error_code,
222-
self.field_check_error_message,
219+
self.ft_error_code or "CSVFieldMismatch",
220+
self.ft_error_message or "The submitted header is invalid",
223221
self.delimiter,
224222
self.quote_char,
225223
)

‎src/dve/core_engine/backends/readers/utilities.py‎

Lines changed: 12 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
from pydantic import BaseModel
77

8-
from dve.core_engine.backends.exceptions import MessageBearingError
8+
from dve.core_engine.backends.exceptions import CriticalMessageBearingError
9+
from dve.core_engine.constants import PRE_VALIDATION_ENTITY
910
from dve.core_engine.message import FeedbackMessage
1011
from dve.core_engine.type_hints import URI, EntityName
1112
from dve.parser.file_handling.service import open_stream
@@ -55,19 +56,17 @@ def raise_message_bearing_error_on_header_differences(
5556
record_details_additional = (
5657
f"additional fields: {', '.join(sorted(additional))};" if additional else ""
5758
) # pylint: disable=C0301
58-
raise MessageBearingError(
59+
raise CriticalMessageBearingError(
5960
"The CSV header doesn't match what is expected",
60-
messages=[
61-
FeedbackMessage(
62-
entity="Pre-validation",
63-
record={entity_name: f"{record_details_missing}{record_details_additional}"},
64-
failure_type="submission",
65-
error_location=entity_name,
66-
reporting_field="csv_header",
67-
error_code=field_check_error_code,
68-
error_message=field_check_error_message,
69-
)
70-
],
61+
message=FeedbackMessage(
62+
entity=PRE_VALIDATION_ENTITY,
63+
record={entity_name: f"{record_details_missing}{record_details_additional}"},
64+
failure_type="submission",
65+
error_location=entity_name,
66+
reporting_field="csv_header",
67+
error_code=field_check_error_code,
68+
error_message=field_check_error_message,
69+
)
7170
)
7271

7372

0 commit comments

Comments
 (0)