diff --git a/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_csv_reader_args.schema.json b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_csv_reader_args.schema.json index 56e2acc6..ea049216 100644 --- a/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_csv_reader_args.schema.json +++ b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_csv_reader_args.schema.json @@ -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", @@ -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." diff --git a/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_reader_args.schema.json b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_reader_args.schema.json new file mode 100644 index 00000000..3da16c23 --- /dev/null +++ b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_reader_args.schema.json @@ -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 + } + } +} \ No newline at end of file diff --git a/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_xml_reader_args.schema.json b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_xml_reader_args.schema.json index d82ce720..40f893c9 100644 --- a/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_xml_reader_args.schema.json +++ b/docs/advanced_guidance/json_schemas/contract/components/reader_constraints/global_xml_reader_args.schema.json @@ -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", @@ -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." diff --git a/src/dve/core_engine/backends/base/reader.py b/src/dve/core_engine/backends/base/reader.py index ff3586c9..93b97071 100644 --- a/src/dve/core_engine/backends/base/reader.py +++ b/src/dve/core_engine/backends/base/reader.py @@ -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, @@ -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__` @@ -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, + ), ) diff --git a/src/dve/core_engine/backends/exceptions.py b/src/dve/core_engine/backends/exceptions.py index bb585168..f8d93079 100644 --- a/src/dve/core_engine/backends/exceptions.py +++ b/src/dve/core_engine/backends/exceptions.py @@ -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", + ) ) diff --git a/src/dve/core_engine/backends/implementations/duckdb/readers/csv.py b/src/dve/core_engine/backends/implementations/duckdb/readers/csv.py index 012673ac..dedd66bd 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/readers/csv.py +++ b/src/dve/core_engine/backends/implementations/duckdb/readers/csv.py @@ -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 @@ -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 @@ -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, **_, ): @@ -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( @@ -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: @@ -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: @@ -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 @@ -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." diff --git a/src/dve/core_engine/backends/implementations/duckdb/readers/xml.py b/src/dve/core_engine/backends/implementations/duckdb/readers/xml.py index ac111694..7e591e58 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/readers/xml.py +++ b/src/dve/core_engine/backends/implementations/duckdb/readers/xml.py @@ -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, @@ -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 diff --git a/src/dve/core_engine/backends/implementations/spark/readers/csv.py b/src/dve/core_engine/backends/implementations/spark/readers/csv.py index 5cd2f568..6b79bb35 100644 --- a/src/dve/core_engine/backends/implementations/spark/readers/csv.py +++ b/src/dve/core_engine/backends/implementations/spark/readers/csv.py @@ -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: @@ -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( diff --git a/src/dve/core_engine/backends/implementations/spark/readers/xml.py b/src/dve/core_engine/backends/implementations/spark/readers/xml.py index 4d6df6a7..275182e8 100644 --- a/src/dve/core_engine/backends/implementations/spark/readers/xml.py +++ b/src/dve/core_engine/backends/implementations/spark/readers/xml.py @@ -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: @@ -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, ) diff --git a/src/dve/core_engine/backends/readers/csv.py b/src/dve/core_engine/backends/readers/csv.py index cfa2dcde..619a3d19 100644 --- a/src/dve/core_engine/backends/readers/csv.py +++ b/src/dve/core_engine/backends/readers/csv.py @@ -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. @@ -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] = { @@ -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, ) diff --git a/src/dve/core_engine/backends/readers/utilities.py b/src/dve/core_engine/backends/readers/utilities.py index 3948d709..99ca6abc 100644 --- a/src/dve/core_engine/backends/readers/utilities.py +++ b/src/dve/core_engine/backends/readers/utilities.py @@ -5,7 +5,8 @@ from pydantic import BaseModel -from dve.core_engine.backends.exceptions import MessageBearingError +from dve.core_engine.backends.exceptions import CriticalMessageBearingError +from dve.core_engine.constants import PRE_VALIDATION_ENTITY from dve.core_engine.message import FeedbackMessage from dve.core_engine.type_hints import URI, EntityName from dve.parser.file_handling.service import open_stream @@ -55,19 +56,17 @@ def raise_message_bearing_error_on_header_differences( record_details_additional = ( f"additional fields: {', '.join(sorted(additional))};" if additional else "" ) # pylint: disable=C0301 - raise MessageBearingError( + raise CriticalMessageBearingError( "The CSV header doesn't match what is expected", - messages=[ - FeedbackMessage( - entity="Pre-validation", - record={entity_name: f"{record_details_missing}{record_details_additional}"}, - failure_type="submission", - error_location=entity_name, - reporting_field="csv_header", - error_code=field_check_error_code, - error_message=field_check_error_message, - ) - ], + message=FeedbackMessage( + entity=PRE_VALIDATION_ENTITY, + record={entity_name: f"{record_details_missing}{record_details_additional}"}, + failure_type="submission", + error_location=entity_name, + reporting_field="csv_header", + error_code=field_check_error_code, + error_message=field_check_error_message, + ) ) diff --git a/src/dve/core_engine/backends/readers/xml.py b/src/dve/core_engine/backends/readers/xml.py index a3fd437c..05167e18 100644 --- a/src/dve/core_engine/backends/readers/xml.py +++ b/src/dve/core_engine/backends/readers/xml.py @@ -132,8 +132,8 @@ def __init__( encoding: str = "utf-8-sig", n_records_to_read: Optional[int] = None, 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, **_, ): @@ -174,10 +174,8 @@ def __init__( else: self.xsd_location = xsd_location # type: ignore """The URI of the xsd file if wishing to perform xsd validation.""" - self.xsd_error_code = xsd_error_code - """The error code to be reported if xsd validation fails (if xsd)""" - self.xsd_error_message = xsd_error_message - """The error message to be reported if xsd validation fails""" + self.ft_error_code = ft_error_code or "MalformedXMLFile" + self.ft_error_message = ft_error_message super().__init__() self._logger = get_logger(__name__) @@ -296,15 +294,11 @@ def _run_xmllint(self, file_uri: URI) -> FeedbackMessage | None: onto the system to run succesfully.""" if self.xsd_location is None: raise AttributeError("Trying to run XML lint with no `xsd_location` provided.") - if self.xsd_error_code is None: - raise AttributeError("Trying to run XML with no `xsd_error_code` provided.") - if self.xsd_error_message is None: - raise AttributeError("Trying to run XML with no `xsd_error_message` provided.") return run_xmllint( file_uri=file_uri, schema_uri=self.xsd_location, - error_code=self.xsd_error_code, - error_message=self.xsd_error_message, + error_code=self.ft_error_code or "XMLFailedXSDCheck", + error_message=self.ft_error_message or "XML Submission has failed XSD check", ) def read_to_py_iterator( diff --git a/src/dve/core_engine/backends/readers/xml_linting.py b/src/dve/core_engine/backends/readers/xml_linting.py index 529d8ee8..910bb65c 100644 --- a/src/dve/core_engine/backends/readers/xml_linting.py +++ b/src/dve/core_engine/backends/readers/xml_linting.py @@ -10,6 +10,7 @@ from typing import Union from uuid import uuid4 +from dve.core_engine.constants import PRE_VALIDATION_ENTITY from dve.core_engine.message import FeedbackMessage from dve.parser.file_handling import copy_resource, get_file_name, get_resource_exists, open_stream from dve.parser.file_handling.implementations.file import file_uri_to_local_path @@ -131,12 +132,12 @@ def run_xmllint( return None return FeedbackMessage( - entity="xsd_validation", + entity=PRE_VALIDATION_ENTITY, record={}, failure_type="submission", is_informational=False, error_type="xsd check", - error_location="Whole File", + error_location="XSD Validation", error_message=error_message, error_code=error_code, ) diff --git a/src/dve/core_engine/constants.py b/src/dve/core_engine/constants.py index a2a4a655..9b67e471 100644 --- a/src/dve/core_engine/constants.py +++ b/src/dve/core_engine/constants.py @@ -6,3 +6,8 @@ CONTRACT_ERROR_VALUE_FIELD_NAME: str = "__error_value" """The name of the field that can be used to extract the field value that caused a pydantic validation error""" + +PRE_VALIDATION_ENTITY: str = "Pre-validation" +""" +Consistent name for the entity/group where errors are raised during file transformation +""" diff --git a/src/dve/pipeline/pipeline.py b/src/dve/pipeline/pipeline.py index 60e266a1..1737ebf8 100644 --- a/src/dve/pipeline/pipeline.py +++ b/src/dve/pipeline/pipeline.py @@ -30,7 +30,7 @@ from dve.core_engine.backends.base.core import EntityManager from dve.core_engine.backends.base.reference_data import BaseRefDataLoader, ReferenceConfig from dve.core_engine.backends.base.rules import BaseStepImplementations -from dve.core_engine.backends.exceptions import MessageBearingError +from dve.core_engine.backends.exceptions import CriticalMessageBearingError, MessageBearingError from dve.core_engine.backends.readers import BaseFileReader from dve.core_engine.backends.readers.utilities import get_all_model_fields from dve.core_engine.backends.types import EntityType @@ -347,6 +347,9 @@ def file_transformation( except MessageBearingError as exc: self._logger.exception("Unexpected file transformation error:") errors.extend(exc.messages) + except CriticalMessageBearingError as exc: + self._logger.exception("Unexpected critical file transformation error:") + errors.append(exc.message) if errors: dump_feedback_errors( diff --git a/tests/features/books.feature b/tests/features/books.feature index 8551a6fc..a36c9e9f 100644 --- a/tests/features/books.feature +++ b/tests/features/books.feature @@ -37,8 +37,9 @@ Feature: Pipeline tests using the books dataset And I add initial audit entries for the submission Then the latest audit record for the submission is marked with processing status file_transformation When I run the file transformation phase - Then the latest audit record for the submission is marked with processing status failed - # TODO - handle above within the stream xml reader - specific + Then the latest audit record for the submission is marked with processing status error_report + When I run the error report phase + Then An error report is produced Scenario: Handle a file that fails XSD validation (duckdb) Given I submit the books file books_xsd_fail.xml for processing diff --git a/tests/test_core_engine/test_backends/test_readers/test_csv.py b/tests/test_core_engine/test_backends/test_readers/test_csv.py index f2e2c8df..413b6145 100644 --- a/tests/test_core_engine/test_backends/test_readers/test_csv.py +++ b/tests/test_core_engine/test_backends/test_readers/test_csv.py @@ -12,6 +12,7 @@ from pydantic import BaseModel from dve.core_engine.backends.exceptions import ( + CriticalMessageBearingError, EmptyFileError, FieldCountMismatch, MessageBearingError, @@ -287,7 +288,7 @@ def test_base_csv_reader_with_additional_fields( ): """Test that message bearing error raised when additional fields provided""" reader = CSVFileReader(field_check=True) - with pytest.raises(MessageBearingError) as exc_info: + with pytest.raises(CriticalMessageBearingError) as exc_info: list(reader.read_to_py_iterator( planet_additional_field_location, "test", @@ -295,7 +296,7 @@ def test_base_csv_reader_with_additional_fields( get_all_model_fields([Planets]) )) - error_msg = exc_info.value.messages[0] + error_msg = exc_info.value.message assert error_msg.record["test"] == "additional fields: add_field1, add_field2;" assert "missing_fields" not in error_msg.record["test"] @@ -306,7 +307,7 @@ def test_base_csv_reader_with_missing_fields( """Test that message bearing error raised when fields are missing from the expected schema""" reader = CSVFileReader(field_check=True) - with pytest.raises(MessageBearingError) as exc_info: + with pytest.raises(CriticalMessageBearingError) as exc_info: list(reader.read_to_py_iterator( planet_location, "test", @@ -314,6 +315,6 @@ def test_base_csv_reader_with_missing_fields( get_all_model_fields([PlanetsWithExtra]) )) - error_msg = exc_info.value.messages[0] + error_msg = exc_info.value.message assert "additional_fields" not in error_msg.record["test"] assert error_msg.record["test"] == "missing fields: random_null;" diff --git a/tests/test_pipeline/test_foundry_ddb_pipeline.py b/tests/test_pipeline/test_foundry_ddb_pipeline.py index f84073ac..6edc06bc 100644 --- a/tests/test_pipeline/test_foundry_ddb_pipeline.py +++ b/tests/test_pipeline/test_foundry_ddb_pipeline.py @@ -158,7 +158,7 @@ def test_foundry_runner_error(planet_test_files, temp_ddb_conn): .select(pl.col("step_name"), pl.col("error_location"), pl.col("error_message")) ) actual_error_df = ( - pl.read_json(perror_path, schema=perror_schema) + pl.read_ndjson(perror_path, schema=perror_schema) .select(pl.col("step_name"), pl.col("error_location"), pl.col("error_message")) ) assert actual_error_df.equals(expected_error_df) diff --git a/tests/testdata/books/nested_books.dischema.json b/tests/testdata/books/nested_books.dischema.json index 52add54f..9b2e8afb 100644 --- a/tests/testdata/books/nested_books.dischema.json +++ b/tests/testdata/books/nested_books.dischema.json @@ -34,8 +34,8 @@ "record_tag": "bookstore", "n_records_to_read": 1, "xsd_location": "nested_books.xsd", - "xsd_error_code": "TESTXSDERROR", - "xsd_error_message": "the xml is poorly structured" + "ft_error_code": "TESTXSDERROR", + "ft_error_message": "the xml is poorly structured" } } } diff --git a/tests/testdata/books/nested_books_ddb.dischema.json b/tests/testdata/books/nested_books_ddb.dischema.json index d53c4165..f697ffa8 100644 --- a/tests/testdata/books/nested_books_ddb.dischema.json +++ b/tests/testdata/books/nested_books_ddb.dischema.json @@ -34,8 +34,8 @@ "record_tag": "bookstore", "n_records_to_read": 1, "xsd_location": "nested_books.xsd", - "xsd_error_code": "TESTXSDERROR", - "xsd_error_message": "the xml is poorly structured" + "ft_error_code": "TESTXSDERROR", + "ft_error_message": "the xml is poorly structured" } } }