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
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@
"title": "Keyword Arguments used across all CSV readers",
"description": "Arguments present in all CSV readers available.",
"type": "object",
"anyOf": [
{
"$ref": "global_reader_args.schema.json"
}
],
"properties": {
"field_check": {
"type": "string",
Expand All @@ -13,14 +18,6 @@
"false"
]
},
"field_check_error_code": {
"type": "string",
"description": "Allows you to customise the error code generated in the error report when the field check fails."
},
"field_check_error_message": {
"type": "string",
"description": "Allows you to customise the error message generated in the error report when the field check fails."
},
"null_empty_strings": {
"type": "boolean",
"description": "Converts empty string values into 'Null' values."
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"$id": "data-ingest:contract/components/reader_constraints/global_reader_args.schema.json",
"title": "Keyword Arguments used across all readers",
"description": "Arguments present in all readers available.",
"type": "object",
"properties": {
"ft_error_code": {
"type": "string",
"description": "Error code to raise when a reader is unable to read the contents of a file.",
"minLength": 1
},
"ft_error_message": {
"type": "string",
"description": "Error message to raise when a reader is unable to read the contents of a file.",
"minLength": 1
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@
"title": "Keyword Arguments used across all XML readers",
"description": "Arguments present in all XML readers.",
"type": "object",
"anyOf": [
{
"$ref": "global_reader_args.schema.json"
}
],
"properties": {
"record_tag": {
"type": "string",
Expand Down Expand Up @@ -40,14 +45,6 @@
"type": "string",
"description": "The Relative path (to the dischema) for the XSD document."
},
"xsd_error_code": {
"type": "string",
"description": "Allows you to customise the error code generated in the error report when the XSD check fails."
},
"xsd_error_message": {
"type": "string",
"description": "Allows you to customise the error message generated in the error report when the XSD check fails."
},
"rules_location": {
"type": "string",
"description": "Allows you to prefix the directory which contains the XSD document."
Expand Down
36 changes: 23 additions & 13 deletions src/dve/core_engine/backends/base/reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,11 @@
from pydantic import BaseModel
from typing_extensions import Protocol

from dve.core_engine.backends.exceptions import MessageBearingError, ReaderLacksEntityTypeSupport
from dve.core_engine.backends.exceptions import (
CriticalMessageBearingError,
MessageBearingError,
ReaderLacksEntityTypeSupport
)
from dve.core_engine.backends.types import EntityName, EntityType
from dve.core_engine.configuration.v1 import (
AllowedAdditionalReaderChecks,
Expand Down Expand Up @@ -69,6 +73,10 @@ class BaseFileReader(ABC):
decorated with the '@read_function' decorator, and is used in `read_entity_type`.

"""
ft_error_code: Optional[str] = "MalformedFile"
"""Default error code for when a file/submission cannot be parsed succesfully."""
ft_error_message: Optional[str] = "The resource doesn't seem to be a valid text file"
"""Default error message for when a file/submission cannot be parsed succesfully."""

def __init_subclass__(cls, *_, **__) -> None:
"""When this class is subclassed, create and populate the `__read_methods__`
Expand Down Expand Up @@ -208,20 +216,22 @@ def _check_likely_text_file(resource: URI) -> bool:
return False
return True

def raise_if_not_sensible_file(self, resource: URI, entity_name: str):
def raise_if_not_sensible_file(
self,
resource: URI,
entity_name: str,
):
"""Sense check that the file is a text file. Raise error if doesn't
appear to be the case."""
if not self._check_likely_text_file(resource):
raise MessageBearingError(
raise CriticalMessageBearingError(
"The submitted file doesn't appear to be text",
messages=[
FeedbackMessage(
entity=entity_name,
record=None,
failure_type="submission",
error_location="Whole File",
error_code="MalformedFile",
error_message="The resource doesn't seem to be a valid text file",
)
],
message=FeedbackMessage(
entity=entity_name,
record=None,
failure_type="submission",
error_location="Whole File",
error_code=self.ft_error_code,
error_message=self.ft_error_message,
),
)
41 changes: 25 additions & 16 deletions src/dve/core_engine/backends/exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,28 +32,37 @@ def __init__(self, *args: object, messages: Messages) -> None:
self.messages = messages
"""The messages to be returned as part of the error."""

class CriticalMessageBearingError(BackendError):
"""
A backend error that comes with a pre-created message.
The intention of this exception vs MessageBearingError is that
this should be used to halt the processing of a submission.
"""

def __init__(self, *args: object, message: FeedbackMessage) -> None:
super().__init__(*args)
self.message = message
"""The message to be returned as part of the error."""

class UnableToParseCSVError(MessageBearingError):
class UnableToParseCSVError(CriticalMessageBearingError):
"""An error raised when unable to parse a CSV file"""

def __init__(
self, entity_name: str, field_check_error_message: str, field_check_error_code: str
self,
entity_name: Optional[str],
error_message: Optional[str],
error_code: Optional[str],
):
super().__init__(
messages=[
FeedbackMessage(
entity="csv_structure",
record={
entity_name: "Unable to parse file. Please check the structure of the file."
},
failure_type="submission",
is_informational=False,
error_type="csv read",
error_location=entity_name,
error_message=field_check_error_message,
error_code=field_check_error_code,
)
]
message=FeedbackMessage(
entity=entity_name,
record=None,
failure_type="submission",
is_informational=False,
error_type="csv read",
error_message=error_message or "Unable to parse the CSV file. Please check the structure of your CSV.", # pylint: disable=C0301
error_code=error_code or "MalformedCSV",
)
)


Expand Down
32 changes: 15 additions & 17 deletions src/dve/core_engine/backends/implementations/duckdb/readers/csv.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,12 @@
duckdb_record_index,
duckdb_write_parquet,
get_duckdb_type_from_annotation,
relation_is_empty,
)
from dve.core_engine.backends.implementations.duckdb.types import SQLType
from dve.core_engine.backends.readers.csv import CSVFileReader
from dve.core_engine.backends.utilities import get_polars_type_from_annotation, polars_record_index
from dve.core_engine.constants import RECORD_INDEX_COLUMN_NAME
from dve.core_engine.constants import PRE_VALIDATION_ENTITY, RECORD_INDEX_COLUMN_NAME
from dve.core_engine.message import FeedbackMessage
from dve.core_engine.type_hints import URI, EntityName
from dve.parser.file_handling import get_content_length
Expand All @@ -44,10 +45,7 @@ class DuckDBCSVReader(CSVFileReader):
to the file header, if it exists.

field_check: flag to compare submitted file header to the accompanying pydantic model
field_check_error_code: The error code to provide if the file header doesn't contain
the expected fields
field_check_error_message: The error message to provide if the file header doesn't contain
the expected fields"""
"""

# TODO - the read_to_relation should include the schema and determine whether to
# TODO - stringify or not
Expand All @@ -59,8 +57,8 @@ def __init__(
quotechar: str = '"',
connection: Optional[DuckDBPyConnection] = None,
field_check: bool = False,
field_check_error_code: str = "ExpectedVsActualFieldMismatch",
field_check_error_message: str = "The submitted header is missing fields",
ft_error_code: str = "ExpectedVsActualFieldMismatch",
ft_error_message: str = "The submitted header is missing fields",
null_empty_strings: bool = False,
**_,
):
Expand All @@ -72,8 +70,8 @@ def __init__(
delimiter=delim,
quote_char=quotechar,
field_check=field_check,
field_check_error_code=field_check_error_code,
field_check_error_message=field_check_error_message,
ft_error_code=ft_error_code,
ft_error_message=ft_error_message,
)

def read_to_py_iterator(
Expand Down Expand Up @@ -124,8 +122,8 @@ def read_to_relation( # pylint: disable=unused-argument
except InvalidInputException as exc:
raise UnableToParseCSVError(
entity_name="csv_structure",
field_check_error_message=self.field_check_error_message,
field_check_error_code=self.field_check_error_code,
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
) from exc

if self.null_empty_strings:
Expand Down Expand Up @@ -184,8 +182,8 @@ def read_to_relation( # pylint: disable=unused-argument
except pl.exceptions.PolarsError as exc:
raise UnableToParseCSVError(
entity_name="csv_structure",
field_check_error_message=self.field_check_error_message,
field_check_error_code=self.field_check_error_code,
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Unable to parse CSV file. Structure is likely malformed.", # pylint: disable=C0301
) from exc

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

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

if entity.pl().shape[0] == 0:
if relation_is_empty(entity):
raise UnableToParseCSVError(
entity_name="csv_structure",
field_check_error_message=self.field_check_error_message,
field_check_error_code=self.field_check_error_code,
error_code=self.ft_error_code,
error_message=self.ft_error_message or "Found zero records after loading CSV. File is likely malformed.", # pylint: disable=C0301
)

return entity
Expand Down Expand Up @@ -273,7 +271,7 @@ def read_to_relation( # pylint: disable=unused-argument
messages=[
FeedbackMessage(
record={entity_name: differing_values},
entity="Pre-validation",
entity=PRE_VALIDATION_ENTITY,
failure_type="submission",
error_message=(
f"Found {no_records} distinct combination of header values."
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
from pydantic import BaseModel

from dve.core_engine.backends.base.reader import read_function
from dve.core_engine.backends.exceptions import MessageBearingError
from dve.core_engine.backends.exceptions import CriticalMessageBearingError
from dve.core_engine.backends.implementations.duckdb.duckdb_helpers import (
duckdb_check_entity_empty,
duckdb_write_parquet,
Expand Down Expand Up @@ -45,9 +45,9 @@ def read_to_relation(
if self.xsd_location:
msg = self._run_xmllint(file_uri=resource)
if msg:
raise MessageBearingError(
raise CriticalMessageBearingError(
"Submitted file failed XSD validation.",
messages=[msg],
message=msg,
)

polars_schema: dict[str, pl.DataType] = { # type: ignore
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,8 +40,8 @@ def __init__(
null_empty_strings: bool = False,
spark_session: Optional[SparkSession] = None,
field_check: bool = False,
field_check_error_code: str = "ExpectedVsActualFieldMismatch",
field_check_error_message: str = "The submitted header is missing fields",
ft_error_code: str = "ExpectedVsActualFieldMismatch",
ft_error_message: str = "The submitted header is missing fields",
**_,
) -> None:

Expand All @@ -56,8 +56,8 @@ def __init__(
quote_char=quote_char,
header=header,
field_check=field_check,
field_check_error_code=field_check_error_code,
field_check_error_message=field_check_error_message,
ft_error_code=ft_error_code,
ft_error_message=ft_error_message,
)

def read_to_py_iterator(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,8 @@ def __init__(
namespace=None,
trim_cells=True,
xsd_location: Optional[URI] = None,
xsd_error_code: Optional[str] = None,
xsd_error_message: Optional[str] = None,
ft_error_code: Optional[str] = None,
ft_error_message: Optional[str] = None,
rules_location: Optional[URI] = None,
**_,
) -> None:
Expand All @@ -92,8 +92,8 @@ def __init__(
null_values=null_values,
sanitise_multiline=sanitise_multiline,
xsd_location=xsd_location,
xsd_error_code=xsd_error_code,
xsd_error_message=xsd_error_message,
ft_error_code=ft_error_code,
ft_error_message=ft_error_message,
rules_location=rules_location,
)

Expand Down
14 changes: 6 additions & 8 deletions src/dve/core_engine/backends/readers/csv.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ def __init__(
null_values: Collection[str] = frozenset({"NULL", "null", ""}),
encoding: str = "utf-8-sig",
field_check: bool = False,
field_check_error_code: str = "CSVFieldMismatch",
field_check_error_message: str = "The submitted header is invalid",
ft_error_code: Optional[str] = "MalformedCSVFile",
ft_error_message: Optional[str] = None,
**_,
):
"""Init function for the base CSV reader.
Expand Down Expand Up @@ -89,10 +89,8 @@ def __init__(
"""Encoding of the CSV file."""
self.field_check = field_check
"""Whether to check the fields are correct in the supplied header or not"""
self.field_check_error_code = field_check_error_code
"""Error code to raise when fields are missing or unexpected"""
self.field_check_error_message = field_check_error_message
"""Error message to raise when fields are missing or unexpected"""
self.ft_error_code = ft_error_code
self.ft_error_message = ft_error_message

def _get_reader_args(self) -> dict[str, Any]:
reader_args: dict[str, Any] = {
Expand Down Expand Up @@ -218,8 +216,8 @@ def perform_field_check(
entity_name,
expected_schema,
all_model_fields,
self.field_check_error_code,
self.field_check_error_message,
self.ft_error_code or "CSVFieldMismatch",
self.ft_error_message or "The submitted header is invalid",
self.delimiter,
self.quote_char,
)
Expand Down
Loading
Loading