From c6cc14d21365c20eebce1f1b884b7ad360f26634 Mon Sep 17 00:00:00 2001 From: georgeRobertson <50412379+georgeRobertson@users.noreply.github.com> Date: Wed, 19 Aug 2026 15:36:56 +0100 Subject: [PATCH] refactor: Refactor SubmissionStatus for record rejection includes record count and record rejections per entity --- src/dve/core_engine/backends/base/auditing.py | 16 +++- .../implementations/duckdb/auditing.py | 2 + .../implementations/spark/auditing.py | 2 + src/dve/core_engine/models.py | 6 +- src/dve/pipeline/pipeline.py | 44 +++++----- src/dve/pipeline/utils.py | 69 +++++++++++++-- src/dve/reporting/excel_report.py | 8 +- tests/features/books.feature | 20 ++--- tests/features/steps/steps_pipeline.py | 18 ++-- .../test_duckdb/test_audit_ddb.py | 35 ++++++-- .../test_spark/test_audit_spark.py | 30 +++++-- tests/test_core_engine/test_engine.py | 4 +- tests/test_pipeline/pipeline_helpers.py | 6 +- tests/test_pipeline/test_duckdb_pipeline.py | 44 +++++++--- .../test_foundry_ddb_pipeline.py | 14 ++-- tests/test_pipeline/test_pipeline.py | 1 - tests/test_pipeline/test_spark_pipeline.py | 83 +++++++++++++++---- tests/test_reporting/test_excel_report.py | 4 + .../testdata/books/nested_books.dischema.json | 14 ++-- .../books/nested_books_ddb.dischema.json | 14 ++-- 20 files changed, 317 insertions(+), 117 deletions(-) diff --git a/src/dve/core_engine/backends/base/auditing.py b/src/dve/core_engine/backends/base/auditing.py index 24cfd3b..6b45643 100644 --- a/src/dve/core_engine/backends/base/auditing.py +++ b/src/dve/core_engine/backends/base/auditing.py @@ -32,7 +32,7 @@ QueueType, SubmissionResult, ) -from dve.pipeline.utils import SubmissionStatus +from dve.pipeline.utils import EntityStatistics, SubmissionStatus AuditReturnType = TypeVar("AuditReturnType") # pylint: disable=invalid-name @@ -168,12 +168,14 @@ def __init__( submission_statistics: AuditorType, transfers: AuditorType, pool: Optional[ExecutorType] = None, + dataset_id: Optional[str] = None, ): """Audit manager to handle writing of audit information to auditors.""" self._processing_status = processing_status self._submission_info = submission_info self._submission_statistics = submission_statistics self._transfers = transfers + self._dataset_id = dataset_id self.pool = pool if self.pool is not None: thread = isinstance(self.pool, ThreadPoolExecutor) @@ -521,8 +523,16 @@ def get_submission_status(self, submission_id: str) -> Optional[SubmissionStatus sub_status.processing_failed = True if processing_rec.submission_result == "validation_failed": sub_status.validation_failed = True - if sub_stats_rec: - sub_status.number_of_records = sub_stats_rec.record_count + if sub_stats_rec and sub_stats_rec.record_count: + if not self._dataset_id: + raise AttributeError( + f"Unable to find dataset id in {type(self).__name__}. Please ensure that " \ + +f"dataset id is defined in the setup of the {type(self).__name__} " \ + +"before using get_submission_status." + ) + sub_status.entity_stats[self._dataset_id] = EntityStatistics( + no_records=sub_stats_rec.record_count + ) return sub_status diff --git a/src/dve/core_engine/backends/implementations/duckdb/auditing.py b/src/dve/core_engine/backends/implementations/duckdb/auditing.py index 803b696..6255fef 100644 --- a/src/dve/core_engine/backends/implementations/duckdb/auditing.py +++ b/src/dve/core_engine/backends/implementations/duckdb/auditing.py @@ -169,6 +169,7 @@ def __init__( database_uri: URI, pool: Optional[ExecutorType] = None, connection: Optional[DuckDBPyRelation] = None, + dataset_id: Optional[str] = None, ): self._database_uri = database_uri self._connection = ( @@ -209,6 +210,7 @@ def __init__( name="transfers", connection=self._connection, # type: ignore ), + dataset_id=dataset_id, pool=self._pool, ) diff --git a/src/dve/core_engine/backends/implementations/spark/auditing.py b/src/dve/core_engine/backends/implementations/spark/auditing.py index 3f50721..3ede58b 100644 --- a/src/dve/core_engine/backends/implementations/spark/auditing.py +++ b/src/dve/core_engine/backends/implementations/spark/auditing.py @@ -172,6 +172,7 @@ def __init__( pool: Optional[ExecutorType] = None, spark: Optional[SparkSession] = None, table_format: Optional[SparkTableFormat] = "delta", + dataset_id: Optional[str] = None, ): self._database = database self._spark = spark if spark else SparkSession.builder.getOrCreate() @@ -209,6 +210,7 @@ def __init__( spark=self._spark, ), pool=self._pool, + dataset_id=dataset_id, ) def combine_auditor_information( diff --git a/src/dve/core_engine/models.py b/src/dve/core_engine/models.py index bba2986..430ef45 100644 --- a/src/dve/core_engine/models.py +++ b/src/dve/core_engine/models.py @@ -116,10 +116,12 @@ class SubmissionStatisticsRecord(AuditRecord): record_count: Optional[int] """Count of records in the submitted file""" + total_number_of_records_rejected: Optional[int] + """Total number of records rejected in a submitted file""" number_submission_rejections: Optional[int] - """Number of submission rejections raised following validation""" + """Number of submission rejection errors raised following validation""" number_record_rejections: Optional[int] - """Number of record rejections raised following validation""" + """Number of record rejection errors raised following validation""" number_warnings: Optional[int] """Number of warnings raised following validation""" diff --git a/src/dve/pipeline/pipeline.py b/src/dve/pipeline/pipeline.py index a9be3ff..34cd10c 100644 --- a/src/dve/pipeline/pipeline.py +++ b/src/dve/pipeline/pipeline.py @@ -42,7 +42,7 @@ from dve.parser import file_handling as fh from dve.parser.file_handling.implementations.file import LocalFilesystemImplementation from dve.parser.file_handling.service import _get_implementation -from dve.pipeline.utils import SubmissionStatus, deadletter_file, load_config, load_reader +from dve.pipeline.utils import EntityStatistics, SubmissionStatus, deadletter_file, load_config, load_reader from dve.reporting.constants import ErrorReportCategories from dve.reporting.error_report import ERROR_SCHEMA, calculate_aggregates @@ -450,10 +450,13 @@ def apply_data_contract( entity_locations = {} for path, _ in fh.iter_prefix(read_from): - entity_locations[fh.get_file_name(path)] = path - entities[fh.get_file_name(path)] = self.data_contract.add_record_index( + entity_name = fh.get_file_name(path) + entity_locations[entity_name] = path + entity = self.data_contract.add_record_index( self.data_contract.read_parquet(path) ) + entities[entity_name] = entity + submission_status.create_new_entity_stat(entity_name, self.get_entity_count(entity)) key_fields = {model: conf.reporting_fields for model, conf in model_config.items()} @@ -620,8 +623,20 @@ def apply_business_rules( # pylint: disable=R0914 entity, entity_name, ) + entity_ct = self.get_entity_count(filtered_entity) + try: + submission_status.entity_stats[entity_name].number_of_record_rejections = ( + submission_status.entity_stats[entity_name].number_of_records + - entity_ct + ) + except KeyError: + # Handling derived entities + submission_status.create_new_entity_stat(entity_name, entity_ct) else: self._logger.info(f"Skipping {entity_name}. Marked original.") + submission_status.entity_stats[entity_name] = EntityStatistics( + no_records=self.get_entity_count(entity) + ) filtered_entity = entity projected = self._step_implementations.write_parquet( # type: ignore filtered_entity, @@ -636,20 +651,6 @@ def apply_business_rules( # pylint: disable=R0914 projected ) - submission_status.number_of_records = self.get_entity_count( - entity=entity_manager.entities[f"""Original{rules.global_variables.get( - 'entity', - submission_info.dataset_id)}"""] - ) - submission_status.number_of_records_rejected = ( - submission_status.number_of_records - - self.get_entity_count( - entity_manager.entities[ - rules.global_variables.get("entity", submission_info.dataset_id) - ] - ) - ) - return submission_info, submission_status def business_rule_step( @@ -823,7 +824,8 @@ def error_report( self._logger.info("Reading error dataframes") errors_df, aggregates = self._get_error_dataframes(submission_info.submission_id) - if not submission_status.number_of_records: + no_records = submission_status.number_of_records(submission_info.dataset_id) + if not no_records: sub_stats = None else: err_types = { @@ -834,7 +836,10 @@ def error_report( } sub_stats = SubmissionStatisticsRecord( submission_id=submission_info.submission_id, - record_count=submission_status.number_of_records, + record_count=no_records, + total_number_of_records_rejected=submission_status.number_of_record_rejections( + submission_info.dataset_id + ), number_submission_rejections=err_types.get( ErrorReportCategories.FILE_REJECTION.reporting_name, 0 ), @@ -904,7 +909,6 @@ def error_report_step( futures.append((info, status, pool.submit(self.error_report, info, status))) for info_dict, status in failed_file_transformation: - status.number_of_records = 0 futures.append((info_dict, status, pool.submit(self.error_report, info_dict, status))) for sub_info, status, future in futures: diff --git a/src/dve/pipeline/utils.py b/src/dve/pipeline/utils.py index e6122c2..52ab3da 100644 --- a/src/dve/pipeline/utils.py +++ b/src/dve/pipeline/utils.py @@ -14,7 +14,7 @@ from dve.core_engine.backends.readers import _READER_REGISTRY from dve.core_engine.configuration.v1 import SchemaName, V1EngineConfig, _ModelConfig from dve.core_engine.loggers import get_logger -from dve.core_engine.type_hints import URI, SubmissionResult +from dve.core_engine.type_hints import URI, EntityName, SubmissionResult from dve.metadata_parser.model_generator import JSONtoPyd Dataset = dict[SchemaName, _ModelConfig] @@ -79,19 +79,45 @@ def deadletter_file(source_uri: URI) -> None: return None +class EntityStatistics: + """Statistics for a given entity""" + + def __init__(self, no_records: int, no_record_rej: Optional[int] = None): + self._number_of_records = no_records + self._number_of_record_rejections = no_record_rej + + @property + def number_of_records(self) -> int: + """Get the number of record for the entity""" + return self._number_of_records + + @number_of_records.setter + def number_of_records(self, no_records: int): + """Set the number of records for the entity""" + self._number_of_records = no_records + + @property + def number_of_record_rejections(self) -> int: + """Get the number of record rejections for the entity""" + return self._number_of_record_rejections if self._number_of_record_rejections else 0 + + @number_of_record_rejections.setter + def number_of_record_rejections(self, no_record_rej: int): + """Set the number of record rejections for the entity""" + self._number_of_record_rejections = no_record_rej + + class SubmissionStatus: """Submission status for a given submission.""" def __init__( self, validation_failed: bool = False, - number_of_records: Optional[int] = None, - number_of_records_rejected: Optional[int] = None, + entity_stats: Optional[dict[EntityName, EntityStatistics]] = None, processing_failed: bool = False, ): self.validation_failed = validation_failed - self.number_of_records = number_of_records - self.number_of_records_rejected = number_of_records_rejected + self.entity_stats = entity_stats if entity_stats else {} self.processing_failed = processing_failed @property @@ -103,3 +129,36 @@ def submission_result(self) -> SubmissionResult: if self.validation_failed: return "validation_failed" return "success" + + def number_of_records(self, record_entity_name: EntityName) -> int: + """The total number of records across entities for a given submission.""" + if not self.entity_stats: + return 0 + + return self.entity_stats[record_entity_name].number_of_records + + def number_of_record_rejections(self, record_entity_name: EntityName) -> int: + """The total number of record rejections across entities for a given submission.""" + if not self.entity_stats: + return 0 + + return self.entity_stats[record_entity_name].number_of_record_rejections + + def create_new_entity_stat(self, entity_name: str, record_count: int): + """Create a new EntityStatistics object for a given entity.""" + if self.entity_stats.get(entity_name): + raise LookupError("Record count is already set for {entity_name}." \ + +"Use update_number_of_records method instead") + self.entity_stats[entity_name] = EntityStatistics(no_records=record_count) + + def update_number_of_records(self, entity_name: str, record_count: int): + """Update the number of records for a given entity""" + _new_record = self.entity_stats[entity_name] + _new_record.number_of_records = record_count + self.entity_stats[entity_name] = _new_record + + def update_number_of_record_rejections(self, entity_name: str, record_rej_count: int): + """Update the number of record rejections for a given entity""" + _new_record = self.entity_stats[entity_name] + _new_record.number_of_record_rejections = record_rej_count + self.entity_stats[entity_name] = _new_record diff --git a/src/dve/reporting/excel_report.py b/src/dve/reporting/excel_report.py index 5876cdd..8a1276e 100644 --- a/src/dve/reporting/excel_report.py +++ b/src/dve/reporting/excel_report.py @@ -147,8 +147,8 @@ def _add_submission_info(self, status: str, summary: Worksheet): "", "Total Number of Records Processed", ( - self.submission_status.number_of_records - if self.submission_status.number_of_records + self.submission_status.number_of_records(self.summary_dict["Dataset Id"]) + if self.submission_status.number_of_records(self.summary_dict["Dataset Id"]) else 0 ), # pylint: disable=C0301 ] @@ -161,7 +161,9 @@ def _add_submission_info(self, status: str, summary: Worksheet): [ "", "Total Number of Records Rejected", - self.submission_status.number_of_records_rejected, + self.submission_status.number_of_record_rejections( + self.summary_dict["Dataset Id"] + ), ] ) summary.append(["", ""]) diff --git a/tests/features/books.feature b/tests/features/books.feature index 8551a6f..d53ce01 100644 --- a/tests/features/books.feature +++ b/tests/features/books.feature @@ -11,17 +11,17 @@ Feature: Pipeline tests using the books dataset Then the latest audit record for the submission is marked with processing status file_transformation When I run the file transformation phase Then the header entity is stored as a parquet after the file_transformation phase - And the nested_books entity is stored as a parquet after the file_transformation phase + And the books entity is stored as a parquet after the file_transformation phase And the latest audit record for the submission is marked with processing status data_contract When I run the data contract phase Then there is 1 record rejection from the data_contract phase And the header entity is stored as a parquet after the data_contract phase - And the nested_books entity is stored as a parquet after the data_contract phase + And the books 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 When I run the business rules phase - Then The rules restrict "nested_books" to 3 qualifying records - And The entity "nested_books" contains an entry for "17.85" in column "total_value_of_books" - And the nested_books entity is stored as a parquet after the business_rules phase + Then The rules restrict "books" to 3 qualifying records + And The entity "books" contains an entry for "17.85" in column "total_value_of_books" + And the books entity is stored as a parquet after the business_rules phase And 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 @@ -57,17 +57,17 @@ Feature: Pipeline tests using the books dataset Then the latest audit record for the submission is marked with processing status file_transformation When I run the file transformation phase Then the header entity is stored as a parquet after the file_transformation phase - And the nested_books entity is stored as a parquet after the file_transformation phase + And the books entity is stored as a parquet after the file_transformation phase And the latest audit record for the submission is marked with processing status data_contract When I run the data contract phase Then there is 1 record rejection from the data_contract phase And the header entity is stored as a parquet after the data_contract phase - And the nested_books entity is stored as a parquet after the data_contract phase + And the books 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 When I run the business rules phase - Then The rules restrict "nested_books" to 3 qualifying records - And The entity "nested_books" contains an entry for "17.85" in column "total_value_of_books" - And the nested_books entity is stored as a parquet after the business_rules phase + Then The rules restrict "books" to 3 qualifying records + And The entity "books" contains an entry for "17.85" in column "total_value_of_books" + And the books entity is stored as a parquet after the business_rules phase And 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 diff --git a/tests/features/steps/steps_pipeline.py b/tests/features/steps/steps_pipeline.py index c72c873..629d0d8 100644 --- a/tests/features/steps/steps_pipeline.py +++ b/tests/features/steps/steps_pipeline.py @@ -54,6 +54,7 @@ def setup_spark_pipeline( audit_tables=SparkAuditingManager( database="dve", spark=spark, + dataset_id=dataset_id, ), job_run_id=12345, rules_path=rules_path, @@ -80,6 +81,7 @@ def setup_duckdb_pipeline( database_uri=db_file.as_posix(), # pool=ThreadPoolExecutor(1), connection=connection, + dataset_id=dataset_id, ), job_run_id=12345, connection=connection, @@ -92,12 +94,15 @@ def setup_duckdb_pipeline( def run_file_transformation_step(context: Context): """Apply the file transformation stage""" pipeline = ctxt.get_pipeline(context) - _success, failed = pipeline.file_transformation_step( + success, failed = pipeline.file_transformation_step( pool=ThreadPoolExecutor(1), submissions_to_process=[ctxt.get_submission_info(context)] ) + if success: + ctxt.set_submission_status(context, success[0][1]) + if failed: ctxt.set_failed_file_transformation(context, failed[0][0]) - + ctxt.set_submission_status(context, failed[0][1]) @when("I run the data contract phase") @@ -105,10 +110,11 @@ def apply_data_contract_with_error(context: Context): """Apply the data contract stage""" pipeline = ctxt.get_pipeline(context) sub_info = ctxt.get_submission_info(context) - sub_status = pipeline._audit_tables.get_submission_status(sub_info.submission_id) - pipeline.data_contract_step( + sub_status = ctxt.get_submission_status(context) + passed_contract, _failed_contract = pipeline.data_contract_step( pool=ThreadPoolExecutor(1), file_transform_results=[(sub_info, sub_status)] ) + ctxt.set_submission_status(context, passed_contract[0][1]) @when("I run the business rules phase") @@ -117,7 +123,7 @@ def apply_business_rules(context: Context): pipeline = ctxt.get_pipeline(context) sub_info = ctxt.get_submission_info(context) - sub_status = pipeline._audit_tables.get_submission_status(sub_info.submission_id) + sub_status = ctxt.get_submission_status(context) success, failed, _ = pipeline.business_rule_step( pool=ThreadPoolExecutor(1), files=[(sub_info, sub_status)] ) @@ -133,7 +139,7 @@ def create_error_report(context: Context): try: failed_ft = ctxt.get_failed_file_transformation(context) - sub_status = pipeline._audit_tables.get_submission_status(failed_ft.submission_id) + sub_status = ctxt.get_submission_status(context) pipeline.error_report_step( pool=ThreadPoolExecutor(1), diff --git a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_audit_ddb.py b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_audit_ddb.py index 2ae7a65..04e9326 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_audit_ddb.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_duckdb/test_audit_ddb.py @@ -11,7 +11,6 @@ from dve.core_engine.backends.implementations.duckdb.auditing import DDBAuditingManager from dve.core_engine.models import ProcessingStatusRecord, SubmissionInfo, SubmissionStatisticsRecord -from dve.pipeline.utils import SubmissionStatus @pytest.fixture(scope="function") @@ -21,7 +20,11 @@ def ddb_audit_manager() -> Iterator[DDBAuditingManager]: db_file = Path(tmp, db + ".duckdb") conn = connect(database=db_file, read_only=False) - yield DDBAuditingManager(database_uri=db_file.as_uri(), connection=conn) + yield DDBAuditingManager( + database_uri=db_file.as_uri(), + connection=conn, + dataset_id="TEST_DATASET", + ) @pytest.fixture(scope="function") @@ -32,7 +35,12 @@ def ddb_audit_manager_threaded() -> Iterator[DDBAuditingManager]: conn = connect(database=db_file, read_only=False) with ThreadPoolExecutor(1) as pool: - yield DDBAuditingManager(database_uri=db_file.as_uri(), pool=pool, connection=conn) + yield DDBAuditingManager( + database_uri=db_file.as_uri(), + pool=pool, + connection=conn, + dataset_id="TEST_DATASET", + ) @pytest.fixture @@ -420,15 +428,28 @@ def test_get_submission_status(ddb_audit_manager_threaded: DDBAuditingManager): ] ) aud.add_submission_statistics_records([ - SubmissionStatisticsRecord(submission_id=sub_1.submission_id, record_count=5, number_submission_rejections=0, number_record_rejections=2, number_warnings=3), - SubmissionStatisticsRecord(submission_id=sub_4.submission_id, record_count=20, number_submission_rejections=0, number_record_rejections=0, number_warnings=1) + SubmissionStatisticsRecord( + submission_id=sub_1.submission_id, record_count=5, + total_number_of_records_rejected=1, + number_submission_rejections=0, + number_record_rejections=2, + number_warnings=3 + ), + SubmissionStatisticsRecord( + submission_id=sub_4.submission_id, + record_count=20, + total_number_of_records_rejected=0, + number_submission_rejections=0, + number_record_rejections=0, + number_warnings=1 + ) ]) sub_stats_1 = ddb_audit_manager_threaded.get_submission_status(sub_1.submission_id) assert sub_stats_1.submission_result == "validation_failed" assert sub_stats_1.validation_failed assert not sub_stats_1.processing_failed - assert sub_stats_1.number_of_records == 5 + assert sub_stats_1.number_of_records("TEST_DATASET") == 5 sub_stats_2 = ddb_audit_manager_threaded.get_submission_status(sub_2.submission_id) assert sub_stats_2.submission_result == "processing_failed" assert not sub_stats_2.validation_failed @@ -440,5 +461,5 @@ def test_get_submission_status(ddb_audit_manager_threaded: DDBAuditingManager): assert sub_stats_4.submission_result == "success" assert not sub_stats_4.validation_failed assert not sub_stats_4.processing_failed - assert sub_stats_4.number_of_records == 20 + assert sub_stats_4.number_of_records("TEST_DATASET") == 20 assert not ddb_audit_manager_threaded.get_submission_status("5") diff --git a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_audit_spark.py b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_audit_spark.py index 9d47dfd..8d04cb6 100644 --- a/tests/test_core_engine/test_backends/test_implementations/test_spark/test_audit_spark.py +++ b/tests/test_core_engine/test_backends/test_implementations/test_spark/test_audit_spark.py @@ -11,8 +11,6 @@ from dve.core_engine.backends.implementations.spark.auditing import SparkAuditingManager from dve.core_engine.models import ProcessingStatusRecord, SubmissionInfo, SubmissionStatisticsRecord -from dve.pipeline.utils import SubmissionStatus, unpersist_all_rdds -from dve.core_engine.backends.implementations.spark.spark_helpers import PYTHON_TYPE_TO_SPARK_TYPE from .....conftest import get_test_file_path from .....fixtures import spark, spark_test_database @@ -22,7 +20,8 @@ def spark_audit_manager(spark, spark_test_database) -> Iterator[SparkAuditingManager]: yield SparkAuditingManager(database=spark_test_database, table_format="delta", - spark=spark + spark=spark, + dataset_id="TEST_DATASET", ) @@ -32,7 +31,8 @@ def spark_audit_manager_threaded(spark, spark_test_database) -> Iterator[SparkAu yield SparkAuditingManager(database=spark_test_database, table_format="delta", spark = spark, - pool=pool + pool=pool, + dataset_id="TEST_DATASET", ) @@ -435,15 +435,29 @@ def test_get_submission_status(spark_audit_manager: SparkAuditingManager): ] ) aud.add_submission_statistics_records([ - SubmissionStatisticsRecord(submission_id=sub_1.submission_id, record_count=5, number_submission_rejections=0, number_record_rejections=2, number_warnings=3), - SubmissionStatisticsRecord(submission_id=sub_4.submission_id, record_count=20, number_submission_rejections=0, number_record_rejections=0, number_warnings=1) + SubmissionStatisticsRecord( + submission_id=sub_1.submission_id, + record_count=5, + total_number_of_records_rejected=1, + number_submission_rejections=0, + number_record_rejections=2, + number_warnings=3 + ), + SubmissionStatisticsRecord( + submission_id=sub_4.submission_id, + record_count=20, + total_number_of_records_rejected=0, + number_submission_rejections=0, + number_record_rejections=0, + number_warnings=1 + ) ]) sub_stats_1 = aud.get_submission_status(sub_1.submission_id) assert sub_stats_1.submission_result == "validation_failed" assert sub_stats_1.validation_failed assert not sub_stats_1.processing_failed - assert sub_stats_1.number_of_records == 5 + assert sub_stats_1.number_of_records("TEST_DATASET") == 5 sub_stats_2 = aud.get_submission_status(sub_2.submission_id) assert sub_stats_2.submission_result == "processing_failed" assert not sub_stats_2.validation_failed @@ -455,5 +469,5 @@ def test_get_submission_status(spark_audit_manager: SparkAuditingManager): assert sub_stats_4.submission_result == "success" assert not sub_stats_4.validation_failed assert not sub_stats_4.processing_failed - assert sub_stats_4.number_of_records == 20 + assert sub_stats_4.number_of_records("TEST_DATASET") == 20 assert not aud.get_submission_status("5") diff --git a/tests/test_core_engine/test_engine.py b/tests/test_core_engine/test_engine.py index 7e0fd6e..73a6230 100644 --- a/tests/test_core_engine/test_engine.py +++ b/tests/test_core_engine/test_engine.py @@ -99,7 +99,7 @@ def test_dummy_books_run(self, spark, temp_dir: str): _, errors_uri = test_instance.run_pipeline( entity_locations={ "header": get_test_file_path("books/nested_books.XML").as_posix(), - "nested_books": get_test_file_path("books/nested_books.XML").as_posix(), + "books": get_test_file_path("books/nested_books.XML").as_posix(), } ) @@ -115,4 +115,4 @@ def test_dummy_books_run(self, spark, temp_dir: str): if dir_item.startswith("part-0000") and dir_item.endswith("parquet"): if path.name not in check_dirs: check_dirs.append(path.name) - assert sorted(check_dirs) == sorted(["nested_books"]) + assert sorted(check_dirs) == sorted(["books"]) diff --git a/tests/test_pipeline/pipeline_helpers.py b/tests/test_pipeline/pipeline_helpers.py index b13bef3..6c58d76 100644 --- a/tests/test_pipeline/pipeline_helpers.py +++ b/tests/test_pipeline/pipeline_helpers.py @@ -15,7 +15,7 @@ import pytest from dve.core_engine.models import SubmissionInfo -from dve.pipeline.utils import SubmissionStatus +from dve.pipeline.utils import EntityStatistics, SubmissionStatus import dve.pipeline.utils @@ -347,11 +347,13 @@ def planets_data_after_business_rules() -> Iterator[Tuple[SubmissionInfo, str, S }, }, } + entity_stats = {} for entity_name, data_schema in data_post_br.items(): planet_contract_df = pl.DataFrame(data_schema["data"], data_schema["schema"]) planet_contract_df.write_parquet(Path(output_path, f"{entity_name}.parquet")) + entity_stats[entity_name] = EntityStatistics(no_records=planet_contract_df.shape[0]) - submission_status = SubmissionStatus(False, 1) + submission_status = SubmissionStatus(False, entity_stats) yield submitted_file_info, tdir, submission_status diff --git a/tests/test_pipeline/test_duckdb_pipeline.py b/tests/test_pipeline/test_duckdb_pipeline.py index 4e41af0..1558ee9 100644 --- a/tests/test_pipeline/test_duckdb_pipeline.py +++ b/tests/test_pipeline/test_duckdb_pipeline.py @@ -20,7 +20,7 @@ from dve.core_engine.models import ProcessingStatusRecord, SubmissionInfo, SubmissionStatisticsRecord import dve.parser.file_handling as fh from dve.pipeline.duckdb_pipeline import DDBDVEPipeline -from dve.pipeline.utils import SubmissionStatus +from dve.pipeline.utils import EntityStatistics, SubmissionStatus from ..conftest import get_test_file_path from ..fixtures import temp_ddb_conn # pylint: disable=unused-import @@ -38,7 +38,9 @@ def test_audit_received_step( planet_test_files: str, temp_ddb_conn: Tuple[Path, DuckDBPyConnection] ): # pylint: disable=redefined-outer-name db_file, conn = temp_ddb_conn - with DDBAuditingManager(db_file.as_uri(), ThreadPoolExecutor(1), conn) as audit_manager: + with DDBAuditingManager( + db_file.as_uri(), ThreadPoolExecutor(1), conn, "planets" + ) as audit_manager: dve_pipeline = DDBDVEPipeline( processed_files_path=planet_test_files, audit_tables=audit_manager, @@ -78,7 +80,9 @@ def test_file_transformation_step( planet_test_files: str, temp_ddb_conn: Tuple[Path, DuckDBPyConnection] ): # pylint: disable=redefined-outer-name db_file, conn = temp_ddb_conn - with DDBAuditingManager(db_file.as_uri(), ThreadPoolExecutor(1), conn) as audit_manager: + with DDBAuditingManager( + db_file.as_uri(), ThreadPoolExecutor(1), conn, "planets" + ) as audit_manager: dve_pipeline = DDBDVEPipeline( processed_files_path=planet_test_files, audit_tables=audit_manager, @@ -115,7 +119,9 @@ def test_data_contract_step( ): # pylint: disable=redefined-outer-name db_file, conn = temp_ddb_conn sub_info, processed_file_path = planet_data_after_file_transformation - with DDBAuditingManager(db_file.as_uri(), ThreadPoolExecutor(1), conn) as audit_manager: + with DDBAuditingManager( + db_file.as_uri(), ThreadPoolExecutor(1), conn, "planets" + ) as audit_manager: dve_pipeline = DDBDVEPipeline( processed_files_path=processed_file_path, audit_tables=audit_manager, @@ -148,7 +154,9 @@ def test_business_rule_step( db_file, conn = temp_ddb_conn sub_info, processed_files_path = planets_data_after_data_contract - with DDBAuditingManager(db_file.as_uri(), ThreadPoolExecutor(1), conn) as audit_manager: + with DDBAuditingManager( + db_file.as_uri(), ThreadPoolExecutor(1), conn, "planets" + ) as audit_manager: dve_pipeline = DDBDVEPipeline( processed_files_path=processed_files_path, audit_tables=audit_manager, @@ -160,7 +168,14 @@ def test_business_rule_step( audit_manager.add_new_submissions([sub_info], job_run_id=1) successful_files, unsuccessful_files, failed_processing = dve_pipeline.business_rule_step( - pool=ThreadPoolExecutor(2), files=[(sub_info, SubmissionStatus())] + pool=ThreadPoolExecutor(2), + files=[( + sub_info, + SubmissionStatus(entity_stats={ + "planets": EntityStatistics(no_records=1), + "largest_satellites": EntityStatistics(no_records=1) + }) + )] ) assert len(successful_files) == 1 @@ -183,7 +198,9 @@ def test_error_report_step( db_file, conn = temp_ddb_conn submitted_file_info, processed_files_path, status = planets_data_after_business_rules - with DDBAuditingManager(db_file.as_uri(), ThreadPoolExecutor(1), conn) as audit_manager: + with DDBAuditingManager( + db_file.as_uri(), ThreadPoolExecutor(1), conn, "planets" + ) as audit_manager: dve_pipeline = DDBDVEPipeline( processed_files_path=processed_files_path, audit_tables=audit_manager, @@ -206,7 +223,7 @@ def test_error_report_step( def test_get_submission_status(temp_ddb_conn): db_file, conn = temp_ddb_conn - with DDBAuditingManager(db_file.as_uri(), connection = conn) as aud: + with DDBAuditingManager(db_file.as_uri(), connection = conn, dataset_id="planets") as aud: dve_pipeline = DDBDVEPipeline( processed_files_path="fake_path", audit_tables=aud, @@ -250,14 +267,21 @@ def test_get_submission_status(temp_ddb_conn): ] ) aud.add_submission_statistics_records([ - SubmissionStatisticsRecord(submission_id=sub_one.submission_id, record_count=5, number_submission_rejections=0, number_record_rejections=2, number_warnings=3), + SubmissionStatisticsRecord( + submission_id=sub_one.submission_id, + record_count=5, + total_number_of_records_rejected=1, + number_submission_rejections=0, + number_record_rejections=2, + number_warnings=3 + ), ]) sub_stats_one = dve_pipeline.get_submission_status("test", sub_one.submission_id) assert sub_stats_one.submission_result == "validation_failed" assert sub_stats_one.validation_failed assert not sub_stats_one.processing_failed - assert sub_stats_one.number_of_records == 5 + assert sub_stats_one.number_of_records("planets") == 5 sub_stats_two = dve_pipeline.get_submission_status("test", sub_two.submission_id) assert sub_stats_two.submission_result == "processing_failed" assert not sub_stats_two.validation_failed diff --git a/tests/test_pipeline/test_foundry_ddb_pipeline.py b/tests/test_pipeline/test_foundry_ddb_pipeline.py index 9b7b60d..7946129 100644 --- a/tests/test_pipeline/test_foundry_ddb_pipeline.py +++ b/tests/test_pipeline/test_foundry_ddb_pipeline.py @@ -41,7 +41,9 @@ def prep_multithreading_test(): conn.read_parquet( get_test_file_path("movies/refdata/movies_sequels.parquet").as_posix() ).to_table("movies_refdata.sequels") - sub_details[f"submission_{idx}"] = (conn, tmp_dir, DDBAuditingManager(None, None, conn)) + sub_details[f"submission_{idx}"] = ( + conn, tmp_dir, DDBAuditingManager(None, None, conn, "movies") + ) yield sub_details for con, db_dir, aud in sub_details.values(): @@ -60,7 +62,7 @@ def test_foundry_runner_validation_fail(planet_test_files, temp_ddb_conn): shutil.copytree(planet_test_files, sub_folder) - with DDBAuditingManager(db_file.as_uri(), None, conn) as audit_manager: + with DDBAuditingManager(db_file.as_uri(), None, conn, "planets") as audit_manager: dve_pipeline = FoundryDDBPipeline( processed_files_path=processing_folder, audit_tables=audit_manager, @@ -92,7 +94,7 @@ def test_foundry_runner_validation_success(movies_test_files, temp_ddb_conn): shutil.copytree(movies_test_files, sub_folder) - with DDBAuditingManager(db_file.as_uri(), None, conn) as audit_manager: + with DDBAuditingManager(db_file.as_uri(), None, conn, "movies") as audit_manager: dve_pipeline = FoundryDDBPipeline( processed_files_path=processing_folder, audit_tables=audit_manager, @@ -116,7 +118,7 @@ def test_foundry_runner_error(planet_test_files, temp_ddb_conn): shutil.copytree(planet_test_files, sub_folder) - with DDBAuditingManager(db_file.as_uri(), None, conn) as audit_manager: + with DDBAuditingManager(db_file.as_uri(), None, conn, "planets") as audit_manager: dve_pipeline = FoundryDDBPipeline( processed_files_path=processing_folder, audit_tables=audit_manager, @@ -185,7 +187,7 @@ def test_foundry_runner_with_submitted_files_path(movies_test_files, temp_ddb_co datetime_received=datetime(2025,11,5) ) - with DDBAuditingManager(db_file.as_uri(), None, conn) as audit_manager: + with DDBAuditingManager(db_file.as_uri(), None, conn, "movies") as audit_manager: dve_pipeline = FoundryDDBPipeline( processed_files_path=processing_folder, audit_tables=audit_manager, @@ -216,7 +218,7 @@ def test_foundry_runner_error_at_bi_rules(movies_test_files, temp_ddb_conn): datetime_received=datetime(2025,11,5) ) - with DDBAuditingManager(db_file.as_uri(), None, conn) as audit_manager: + with DDBAuditingManager(db_file.as_uri(), None, conn, "movies") as audit_manager: dve_pipeline = FoundryDDBPipeline( processed_files_path=processing_folder, audit_tables=audit_manager, diff --git a/tests/test_pipeline/test_pipeline.py b/tests/test_pipeline/test_pipeline.py index a8f59c7..470c3db 100644 --- a/tests/test_pipeline/test_pipeline.py +++ b/tests/test_pipeline/test_pipeline.py @@ -9,7 +9,6 @@ from pyspark.sql import DataFrame from uuid import uuid4 -from dve.core_engine.backends.implementations.duckdb.auditing import DDBAuditingManager from dve.core_engine.models import SubmissionInfo from dve.pipeline.pipeline import BaseDVEPipeline diff --git a/tests/test_pipeline/test_spark_pipeline.py b/tests/test_pipeline/test_spark_pipeline.py index dd28e26..9f82f01 100644 --- a/tests/test_pipeline/test_spark_pipeline.py +++ b/tests/test_pipeline/test_spark_pipeline.py @@ -25,7 +25,7 @@ from dve.core_engine.models import ProcessingStatusRecord, SubmissionInfo, SubmissionStatisticsRecord import dve.parser.file_handling as fh from dve.pipeline.spark_pipeline import SparkDVEPipeline -from dve.pipeline.utils import SubmissionStatus +from dve.pipeline.utils import EntityStatistics, SubmissionStatus from ..conftest import get_test_file_path from ..fixtures import spark, spark_test_database # pylint: disable=unused-import @@ -42,7 +42,9 @@ def test_audit_received_step(planet_test_files, spark, spark_test_database): - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_tables: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_tables: dve_pipeline = SparkDVEPipeline( processed_files_path=planet_test_files, audit_tables=audit_tables, @@ -83,7 +85,9 @@ def test_file_transformation_step( spark_test_database: str, planet_test_files: str, ): # pylint: disable=redefined-outer-name - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=planet_test_files, audit_tables=audit_manager, @@ -217,7 +221,9 @@ def test_data_contract_step( ): # pylint: disable=redefined-outer-name sub_info, processed_file_path = planet_data_after_file_transformation sub_status = SubmissionStatus() - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=processed_file_path, audit_tables=audit_manager, @@ -247,7 +253,9 @@ def test_apply_business_rules_success( ): # pylint: disable=redefined-outer-name sub_info, processed_file_path = planets_data_after_data_contract - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=processed_file_path, audit_tables=audit_manager, @@ -257,10 +265,17 @@ def test_apply_business_rules_success( spark=spark, ) - _, status = dve_pipeline.apply_business_rules(sub_info, SubmissionStatus()) + _, status = dve_pipeline.apply_business_rules( + sub_info, + SubmissionStatus(entity_stats={ + "planets": EntityStatistics(no_records=1), + "largest_satellites": EntityStatistics(no_records=1), + }) + ) + assert not status.validation_failed - assert status.number_of_records == 1 + assert status.number_of_records("planets") == 1 planets_entity_path = Path( Path(processed_file_path), sub_info.submission_id, "business_rules", "planets" @@ -288,7 +303,9 @@ def test_apply_business_rules_with_data_errors( # pylint: disable=redefined-out ): sub_info, processed_file_path = planets_data_after_data_contract_that_break_business_rules - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=processed_file_path, audit_tables=audit_manager, @@ -298,10 +315,16 @@ def test_apply_business_rules_with_data_errors( # pylint: disable=redefined-out spark=spark, ) - _, status = dve_pipeline.apply_business_rules(sub_info, SubmissionStatus()) + _, status = dve_pipeline.apply_business_rules( + sub_info, + SubmissionStatus(entity_stats={ + "planets": EntityStatistics(no_records=1), + "largest_satellites": EntityStatistics(no_records=1), + }) + ) assert status.validation_failed - assert status.number_of_records == 1 + assert status.number_of_records("planets") == 1 br_path = Path( Path(processed_file_path), @@ -367,7 +390,9 @@ def test_business_rule_step( ): # pylint: disable=redefined-outer-name sub_info, processed_files_path = planets_data_after_data_contract - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=processed_files_path, audit_tables=audit_manager, @@ -379,7 +404,13 @@ def test_business_rule_step( audit_manager.add_new_submissions([sub_info], job_run_id=1) successful_files, unsuccessful_files, failed_processing = dve_pipeline.business_rule_step( - pool=ThreadPoolExecutor(2), files=[(sub_info, SubmissionStatus())] + pool=ThreadPoolExecutor(2), files=[( + sub_info, + SubmissionStatus(entity_stats={ + "planets": EntityStatistics(no_records=1), + "largest_satellites": EntityStatistics(no_records=1), + }) + )] ) assert len(successful_files) == 1 @@ -408,8 +439,11 @@ def test_error_report_where_report_is_expected( # pylint: disable=redefined-out spark=spark, ) + mock_stats = { + "planets": EntityStatistics(no_records=9, no_record_rej=2) + } submission_info, status, stats, report_uri = dve_pipeline.error_report( - sub_info, SubmissionStatus(True, 9, 2) + sub_info, SubmissionStatus(True, mock_stats) ) assert status.validation_failed @@ -521,7 +555,9 @@ def test_error_report_step( ): # pylint: disable=redefined-outer-name submitted_file_info, processed_files_path, status = planets_data_after_business_rules - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path=processed_files_path, audit_tables=audit_manager, @@ -547,7 +583,9 @@ def test_cluster_pipeline_run( spark: SparkSession, planet_test_files: str, spark_test_database ): # pylint: disable=redefined-outer-name - audit_manager = SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark) + audit_manager = SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark, dataset_id="planets" + ) dve_pipeline = SparkDVEPipeline( processed_files_path=planet_test_files, @@ -568,7 +606,9 @@ def test_cluster_pipeline_run( assert Path(report_uri).exists() def test_get_submission_status(spark, spark_test_database): - with SparkAuditingManager(spark_test_database, ThreadPoolExecutor(1), spark=spark) as audit_manager: + with SparkAuditingManager( + spark_test_database, ThreadPoolExecutor(1), spark=spark, dataset_id="planets" + ) as audit_manager: dve_pipeline = SparkDVEPipeline( processed_files_path="a_path", audit_tables=audit_manager, @@ -612,14 +652,21 @@ def test_get_submission_status(spark, spark_test_database): ] ) audit_manager.add_submission_statistics_records([ - SubmissionStatisticsRecord(submission_id=sub_one.submission_id, record_count=5, number_submission_rejections=0, number_record_rejections=2, number_warnings=3), + SubmissionStatisticsRecord( + submission_id=sub_one.submission_id, + record_count=5, + total_number_of_records_rejected=1, + number_submission_rejections=0, + number_record_rejections=2, + number_warnings=3 + ), ]) sub_stats_one = dve_pipeline.get_submission_status("test", sub_one.submission_id) assert sub_stats_one.submission_result == "validation_failed" assert sub_stats_one.validation_failed assert not sub_stats_one.processing_failed - assert sub_stats_one.number_of_records == 5 + assert sub_stats_one.number_of_records("planets") == 5 sub_stats_two = dve_pipeline.get_submission_status("test", sub_two.submission_id) assert sub_stats_two.submission_result == "processing_failed" assert not sub_stats_two.validation_failed diff --git a/tests/test_reporting/test_excel_report.py b/tests/test_reporting/test_excel_report.py index 1a998ac..d38447c 100644 --- a/tests/test_reporting/test_excel_report.py +++ b/tests/test_reporting/test_excel_report.py @@ -104,6 +104,7 @@ def test_excel_report(report_dfs): "Sender": "X26", "Datetime_sent": datetime.datetime.now(), "Datetime_processed": datetime.datetime.now(), + "Dataset Id": "TEST", }, row_headings=["Submission Failure", "Warning"], table_columns=["Planet", "Derived"], @@ -147,6 +148,7 @@ def test_excel_report_overflow(big_report_dfs): "Sender": "X26", "Datetime_sent": datetime.datetime.now(), "Datetime_processed": datetime.datetime.now(), + "Dataset Id": "TEST", }, row_headings=["Submission Failure", "Warning"], table_columns=["Planet", "Derived"], @@ -171,6 +173,7 @@ def test_excel_report_empty_dfs(): "Sender": "X26", "Datetime_sent": datetime.datetime.now(), "Datetime_processed": datetime.datetime.now(), + "Dataset Id": "TEST", }, row_headings=["Submission Failure", "Warning"], table_columns=["Planet", "Derived"], @@ -191,6 +194,7 @@ def test_sub_status_failed_processing(): "Sender": "X26", "Datetime_sent": datetime.datetime.now(), "Datetime_processed": datetime.datetime.now(), + "Dataset Id": "TEST", }, row_headings=["Submission Failure", "Warning"], table_columns=["Planet", "Derived"], diff --git a/tests/testdata/books/nested_books.dischema.json b/tests/testdata/books/nested_books.dischema.json index 52add54..8fa5036 100644 --- a/tests/testdata/books/nested_books.dischema.json +++ b/tests/testdata/books/nested_books.dischema.json @@ -40,7 +40,7 @@ } } }, - "nested_books": { + "books": { "fields": { "name": "str", "book": { @@ -64,12 +64,12 @@ }, "transformations": { "parameters": { - "entity": "nested_books" + "entity": "books" }, "rules": [ { "operation": "join_header", - "entity": "nested_books", + "entity": "books", "target": "header", "header_column_name": "bookstore" }, @@ -79,7 +79,7 @@ }, { "operation": "select", - "entity": "nested_books", + "entity": "books", "new_entity_name": "exploded_books", "columns": { "explode(book)": "book", @@ -98,9 +98,9 @@ }, { "operation": "join", - "entity": "nested_books", + "entity": "books", "target": "exploded_books", - "join_condition": "nested_books.name == exploded_books.name", + "join_condition": "books.name == exploded_books.name", "new_columns": { "exploded_books.total_value_of_books": "total_value_of_books" } @@ -112,7 +112,7 @@ ], "filters": [ { - "entity": "nested_books", + "entity": "books", "name": "author_has_books", "expression": "book IS NOT NULL AND size(book) >= 1" } diff --git a/tests/testdata/books/nested_books_ddb.dischema.json b/tests/testdata/books/nested_books_ddb.dischema.json index d53c416..4f0e27a 100644 --- a/tests/testdata/books/nested_books_ddb.dischema.json +++ b/tests/testdata/books/nested_books_ddb.dischema.json @@ -40,7 +40,7 @@ } } }, - "nested_books": { + "books": { "fields": { "name": "str", "book": { @@ -64,12 +64,12 @@ }, "transformations": { "parameters": { - "entity": "nested_books" + "entity": "books" }, "rules": [ { "operation": "join_header", - "entity": "nested_books", + "entity": "books", "target": "header", "header_column_name": "bookstore" }, @@ -79,7 +79,7 @@ }, { "operation": "select", - "entity": "nested_books", + "entity": "books", "new_entity_name": "exploded_books", "columns": { "unnest(book)": "book", @@ -98,9 +98,9 @@ }, { "operation": "join", - "entity": "nested_books", + "entity": "books", "target": "exploded_books", - "join_condition": "nested_books.name == exploded_books.name", + "join_condition": "books.name == exploded_books.name", "new_columns": { "exploded_books.total_value_of_books": "total_value_of_books" } @@ -112,7 +112,7 @@ ], "filters": [ { - "entity": "nested_books", + "entity": "books", "name": "author_has_books", "expression": "book IS NOT NULL AND length(book) >= 1" }