From 7a3d8d26ed8b4d94ff1c3c48a22f2f0be1d678df Mon Sep 17 00:00:00 2001 From: Tom Kaltofen Date: Mon, 13 Jul 2026 22:28:46 +0000 Subject: [PATCH] refactor: replace hand-rolled columnar helpers with mloda's public ones mloda 0.10.0 ships columnar_to_rows and homogenize_rows in python_dict_utils, so the local rag_integration/feature_groups/columnar.py (added for the 0.9.0 bump when no public inverse of rows_to_columnar existed) is retired. All imports, including demo.ipynb, now point at the mloda helpers. mloda's columnar_to_rows is strict: it raises ValueError on any non-columnar-dict input instead of passing lists through or mapping None/str to []. Call sites fed by the PythonDict framework are unaffected (they always receive a columnar dict). flatten_result keeps the tolerant shape dispatch locally and only hands genuine columnar dicts to the strict helper; a new tests/integration/test_helpers.py pins that dispatch. Unit tests that handed calculate_feature row-wise lists now build the columnar shape the framework actually delivers (via rows_to_columnar + homogenize_rows for heterogeneous rows). calculate_feature overrides that claimed data: List[Dict[str, Any]] are widened to data: Any, matching the runtime contract and the remaining overrides in the repo. Closes #82 --- cli/eval_demo.py | 2 +- cli/rag_demo.py | 2 +- cli/swap_demo.py | 2 +- demo.ipynb | 2 +- rag_integration/feature_groups/columnar.py | 39 ----- .../connectors/graph_rag/base.py | 2 +- .../feature_groups/datasets/image/base.py | 2 +- .../feature_groups/datasets/text/base.py | 2 +- .../feature_groups/deduplication_base.py | 4 +- .../evaluation/faiss_retrieval_evaluator.py | 4 +- .../evaluation/retrieval_evaluator.py | 4 +- .../image_pipeline/embedding/base.py | 4 +- .../image_pipeline/embedding/clip.py | 4 +- .../image_pipeline/image_source/base.py | 2 +- .../image_pipeline/pii_redaction/base.py | 4 +- .../image_pipeline/preprocessing/base.py | 4 +- .../rag_pipeline/chunking/base.py | 4 +- .../rag_pipeline/document_source/base.py | 2 +- .../rag_pipeline/embedding/base.py | 4 +- .../rag_pipeline/pii_redaction/base.py | 4 +- .../rag_pipeline/pii_redaction/pattern.py | 2 +- .../rag_pipeline/vector_store/base.py | 4 +- .../connectors/generate/generate_contract.py | 2 +- .../graph_rag/graph_rag_contract.py | 2 +- .../graph_rag/test_graph_source_chaining.py | 10 +- tests/connectors/graph_rag/test_kg_source.py | 2 +- .../orchestrator/orchestrator_contract.py | 2 +- tests/connectors/rerank/rerank_contract.py | 2 +- .../connectors/retrieve/retrieve_contract.py | 2 +- .../structured/structured_contract.py | 2 +- tests/feature_groups/test_columnar.py | 77 --------- .../feature_groups/test_deduplication_base.py | 20 +-- .../test_faiss_retrieval_evaluator.py | 159 ++++++++++-------- .../test_retrieval_evaluator.py | 49 ++++-- tests/integration/helpers.py | 17 +- tests/integration/test_helpers.py | 38 +++++ 36 files changed, 225 insertions(+), 262 deletions(-) delete mode 100644 rag_integration/feature_groups/columnar.py delete mode 100644 tests/feature_groups/test_columnar.py create mode 100644 tests/integration/test_helpers.py diff --git a/cli/eval_demo.py b/cli/eval_demo.py index 6411cda..b559ec6 100644 --- a/cli/eval_demo.py +++ b/cli/eval_demo.py @@ -44,7 +44,7 @@ from mloda.user import Options -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.datasets.text.scifact import ScifactDatasetSource from rag_integration.feature_groups.datasets.image.flickr30k import Flickr30kDatasetSource from rag_integration.feature_groups.evaluation.metrics import mean_recall_at_k diff --git a/cli/rag_demo.py b/cli/rag_demo.py index 043b1db..1df83c0 100644 --- a/cli/rag_demo.py +++ b/cli/rag_demo.py @@ -17,7 +17,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.rag_pipeline import ( DictDocumentSource, FixedSizeChunker, diff --git a/cli/swap_demo.py b/cli/swap_demo.py index f772285..d5b1fa9 100644 --- a/cli/swap_demo.py +++ b/cli/swap_demo.py @@ -44,7 +44,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.generate import ExtractiveResponder, TemplateResponder from rag_integration.feature_groups.connectors.orchestrator import HaystackOrchestrator from rag_integration.feature_groups.connectors.retrieve import Bm25sRetriever, TfidfRetriever diff --git a/demo.ipynb b/demo.ipynb index 06d0518..c227606 100644 --- a/demo.ipynb +++ b/demo.ipynb @@ -275,7 +275,7 @@ "from rag_integration.feature_groups.rag_pipeline.embedding.sentence_transformer import SentenceTransformerEmbedder\n", "from rag_integration.feature_groups.rag_pipeline.vector_store.faiss_flat import FaissFlatIndexer\n", "from rag_integration.feature_groups.evaluation.faiss_retrieval_evaluator import FaissRetrievalEvaluator\n", - "from rag_integration.feature_groups.columnar import columnar_to_rows\n", + "from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows\n", "\n", "feature_name = \"eval_docs__chunked__deduped__embedded__indexed__evaluated\"\n", "\n", diff --git a/rag_integration/feature_groups/columnar.py b/rag_integration/feature_groups/columnar.py deleted file mode 100644 index 22b9b03..0000000 --- a/rag_integration/feature_groups/columnar.py +++ /dev/null @@ -1,39 +0,0 @@ -"""Columnar helpers for the PythonDict framework. - -mloda 0.9.0 made PythonDict's native representation columnar ``dict[str, list]`` -(one entry per column, values aligned by row index). The feature groups here read -their upstream input row-wise, so they pivot the columnar input back to rows. -""" - -from __future__ import annotations - -from typing import Any, Dict, List - - -def columnar_to_rows(data: Any) -> List[Dict[str, Any]]: - """Pivot columnar ``dict[str, list]`` input to row-wise ``list[dict]``. - - A ``list`` is returned unchanged (already row-wise); anything else, including - the schema-less empty dict, yields an empty list. - """ - if isinstance(data, list): - return data - if not isinstance(data, dict) or not data: - return [] - columns = list(data.keys()) - n_rows = len(data[columns[0]]) - return [{column: data[column][i] for column in columns} for i in range(n_rows)] - - -def homogenize_rows(rows: List[Dict[str, Any]]) -> List[Dict[str, Any]]: - """Give every row the same key set, backfilling missing keys with ``None``. - - mloda 0.9.0's columnar output contract rejects rows with heterogeneous keys. - Dataset sources whose row types carry different fields (e.g. query rows add - ``relevant_doc_ids``) pass through here so the union schema is uniform. - """ - all_keys: Dict[str, None] = {} - for row in rows: - for key in row: - all_keys.setdefault(key, None) - return [{key: row.get(key) for key in all_keys} for row in rows] diff --git a/rag_integration/feature_groups/connectors/graph_rag/base.py b/rag_integration/feature_groups/connectors/graph_rag/base.py index 958b9a6..cde78ec 100644 --- a/rag_integration/feature_groups/connectors/graph_rag/base.py +++ b/rag_integration/feature_groups/connectors/graph_rag/base.py @@ -41,7 +41,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.errors import DuplicateDocIdError, InvalidOptionError from rag_integration.feature_groups.connectors.mixins import ( DocCollectionMixin, diff --git a/rag_integration/feature_groups/datasets/image/base.py b/rag_integration/feature_groups/datasets/image/base.py index 91a5076..29ac1d5 100644 --- a/rag_integration/feature_groups/datasets/image/base.py +++ b/rag_integration/feature_groups/datasets/image/base.py @@ -11,7 +11,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import homogenize_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import homogenize_rows class BaseImageDatasetSource(FeatureGroup): diff --git a/rag_integration/feature_groups/datasets/text/base.py b/rag_integration/feature_groups/datasets/text/base.py index ef5c97f..8b822b4 100644 --- a/rag_integration/feature_groups/datasets/text/base.py +++ b/rag_integration/feature_groups/datasets/text/base.py @@ -11,7 +11,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import homogenize_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import homogenize_rows class BaseTextDatasetSource(FeatureGroup): diff --git a/rag_integration/feature_groups/deduplication_base.py b/rag_integration/feature_groups/deduplication_base.py index 7b6d547..1291121 100644 --- a/rag_integration/feature_groups/deduplication_base.py +++ b/rag_integration/feature_groups/deduplication_base.py @@ -23,7 +23,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseRowDeduplicator(FeatureChainParserMixin, FeatureGroup): @@ -84,7 +84,7 @@ def _item_size(cls, item: Any) -> int: return len(item) @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Deduplicate rows: attach duplicate metadata and filter by keep strategy. Exactly one distinct feature is processed per call. Unlike column-adding feature diff --git a/rag_integration/feature_groups/evaluation/faiss_retrieval_evaluator.py b/rag_integration/feature_groups/evaluation/faiss_retrieval_evaluator.py index 3faa27f..cf094d6 100644 --- a/rag_integration/feature_groups/evaluation/faiss_retrieval_evaluator.py +++ b/rag_integration/feature_groups/evaluation/faiss_retrieval_evaluator.py @@ -31,7 +31,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.evaluation.metrics import mean_recall_at_k @@ -76,7 +76,7 @@ def compute_framework_rule(cls) -> Optional[Set[Type[ComputeFramework]]]: return {PythonDictFramework} @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Build FAISS index from corpus embeddings, search with query embeddings, compute Recall@K.""" try: import faiss diff --git a/rag_integration/feature_groups/evaluation/retrieval_evaluator.py b/rag_integration/feature_groups/evaluation/retrieval_evaluator.py index 7dc675f..32c702d 100644 --- a/rag_integration/feature_groups/evaluation/retrieval_evaluator.py +++ b/rag_integration/feature_groups/evaluation/retrieval_evaluator.py @@ -15,7 +15,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.evaluation.metrics import mean_recall_at_k @@ -64,7 +64,7 @@ def compute_framework_rule(cls) -> Optional[Set[Type[ComputeFramework]]]: return {PythonDictFramework} @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Compute Recall@K over the embedded corpus and query rows.""" try: import numpy as np diff --git a/rag_integration/feature_groups/image_pipeline/embedding/base.py b/rag_integration/feature_groups/image_pipeline/embedding/base.py index 0d9093c..7f3a05f 100644 --- a/rag_integration/feature_groups/image_pipeline/embedding/base.py +++ b/rag_integration/feature_groups/image_pipeline/embedding/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseImageEmbedder(FeatureChainParserMixin, FeatureGroup): @@ -143,7 +143,7 @@ def _embed_image( ... @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Generate embeddings for images, processing row by row for memory efficiency.""" # mloda 0.9.0 passes columnar data; pivot to rows for row-wise reading. rows = columnar_to_rows(data) diff --git a/rag_integration/feature_groups/image_pipeline/embedding/clip.py b/rag_integration/feature_groups/image_pipeline/embedding/clip.py index a4a5f81..89cb9e3 100644 --- a/rag_integration/feature_groups/image_pipeline/embedding/clip.py +++ b/rag_integration/feature_groups/image_pipeline/embedding/clip.py @@ -16,7 +16,7 @@ from mloda.provider import FeatureSet, property_spec from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.image_pipeline.embedding.base import BaseImageEmbedder # Default local model path: looks two levels above the git repo root (mloda/models/) @@ -182,7 +182,7 @@ def _embed_text( return embedding # type: ignore[no-any-return] @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """ Embed each row using CLIP's vision or text encoder based on row content. diff --git a/rag_integration/feature_groups/image_pipeline/image_source/base.py b/rag_integration/feature_groups/image_pipeline/image_source/base.py index 375d67d..44850b0 100644 --- a/rag_integration/feature_groups/image_pipeline/image_source/base.py +++ b/rag_integration/feature_groups/image_pipeline/image_source/base.py @@ -11,7 +11,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import homogenize_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import homogenize_rows class BaseImageSource(FeatureGroup): diff --git a/rag_integration/feature_groups/image_pipeline/pii_redaction/base.py b/rag_integration/feature_groups/image_pipeline/pii_redaction/base.py index d0937e4..463a74e 100644 --- a/rag_integration/feature_groups/image_pipeline/pii_redaction/base.py +++ b/rag_integration/feature_groups/image_pipeline/pii_redaction/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseImagePIIRedactor(FeatureChainParserMixin, FeatureGroup): @@ -159,7 +159,7 @@ def _redact_region_for_feature( return cls._redact_region(image_data, image_format, regions) @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Perform PII redaction on images, processing row by row for memory efficiency.""" # mloda 0.9.0 passes columnar data; pivot to rows for row-wise reading. rows = columnar_to_rows(data) diff --git a/rag_integration/feature_groups/image_pipeline/preprocessing/base.py b/rag_integration/feature_groups/image_pipeline/preprocessing/base.py index b444869..97d746f 100644 --- a/rag_integration/feature_groups/image_pipeline/preprocessing/base.py +++ b/rag_integration/feature_groups/image_pipeline/preprocessing/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseImagePreprocessor(FeatureChainParserMixin, FeatureGroup): @@ -130,7 +130,7 @@ def _preprocess_image( ... @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Preprocess images row by row for memory efficiency.""" # mloda 0.9.0 passes columnar data; pivot to rows for row-wise reading. rows = columnar_to_rows(data) diff --git a/rag_integration/feature_groups/rag_pipeline/chunking/base.py b/rag_integration/feature_groups/rag_pipeline/chunking/base.py index 43c7150..6ee955e 100644 --- a/rag_integration/feature_groups/rag_pipeline/chunking/base.py +++ b/rag_integration/feature_groups/rag_pipeline/chunking/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseChunker(FeatureChainParserMixin, FeatureGroup): @@ -132,7 +132,7 @@ def _chunk_text_for_feature(cls, text: str, feature: Feature) -> List[str]: return cls._chunk_text(text, cls._get_chunk_size(feature), cls._get_chunk_overlap(feature)) @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Perform chunking on the source feature.""" result = [] diff --git a/rag_integration/feature_groups/rag_pipeline/document_source/base.py b/rag_integration/feature_groups/rag_pipeline/document_source/base.py index f05112f..fac4bba 100644 --- a/rag_integration/feature_groups/rag_pipeline/document_source/base.py +++ b/rag_integration/feature_groups/rag_pipeline/document_source/base.py @@ -11,7 +11,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import homogenize_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import homogenize_rows class BaseDocumentSource(FeatureGroup): diff --git a/rag_integration/feature_groups/rag_pipeline/embedding/base.py b/rag_integration/feature_groups/rag_pipeline/embedding/base.py index 3695948..62097d6 100644 --- a/rag_integration/feature_groups/rag_pipeline/embedding/base.py +++ b/rag_integration/feature_groups/rag_pipeline/embedding/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BaseEmbedder(FeatureChainParserMixin, FeatureGroup): @@ -129,7 +129,7 @@ def _embed_texts( ... @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Generate embeddings for the source feature, with optional artifact support.""" artifact_cls = cls.artifact() diff --git a/rag_integration/feature_groups/rag_pipeline/pii_redaction/base.py b/rag_integration/feature_groups/rag_pipeline/pii_redaction/base.py index 461e2e1..a2a97d3 100644 --- a/rag_integration/feature_groups/rag_pipeline/pii_redaction/base.py +++ b/rag_integration/feature_groups/rag_pipeline/pii_redaction/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows class BasePIIRedactor(FeatureChainParserMixin, FeatureGroup): @@ -174,7 +174,7 @@ def _redact_texts_for_feature(cls, texts: List[str], feature: Feature) -> List[s ) @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Perform PII redaction on the source feature.""" # mloda 0.9.0 delivers columnar data; read it row-wise. rows = columnar_to_rows(data) diff --git a/rag_integration/feature_groups/rag_pipeline/pii_redaction/pattern.py b/rag_integration/feature_groups/rag_pipeline/pii_redaction/pattern.py index 369d0ff..8f0a978 100644 --- a/rag_integration/feature_groups/rag_pipeline/pii_redaction/pattern.py +++ b/rag_integration/feature_groups/rag_pipeline/pii_redaction/pattern.py @@ -79,7 +79,7 @@ def _get_patterns(cls, feature: Feature) -> Dict[str, Pattern[str]]: return patterns @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: Any) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: Any) -> List[Dict[str, Any]]: """Extract custom patterns from feature options before redacting.""" for feature in features.features: cls._active_patterns = cls._get_patterns(feature) diff --git a/rag_integration/feature_groups/rag_pipeline/vector_store/base.py b/rag_integration/feature_groups/rag_pipeline/vector_store/base.py index 47c47cc..272b145 100644 --- a/rag_integration/feature_groups/rag_pipeline/vector_store/base.py +++ b/rag_integration/feature_groups/rag_pipeline/vector_store/base.py @@ -13,7 +13,7 @@ ) from mloda.provider import DefaultOptionKeys -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.rag_pipeline.vector_store.vector_store_artifact import VectorStoreArtifact @@ -91,7 +91,7 @@ def _index_type_name(cls) -> str: ... @classmethod - def calculate_feature(cls, data: List[Dict[str, Any]], features: FeatureSet) -> List[Dict[str, Any]]: + def calculate_feature(cls, data: Any, features: FeatureSet) -> List[Dict[str, Any]]: """Build FAISS index from embeddings, save via artifact, attach row metadata.""" artifact_cls = cls.artifact() diff --git a/tests/connectors/generate/generate_contract.py b/tests/connectors/generate/generate_contract.py index d9e8b72..69e7219 100644 --- a/tests/connectors/generate/generate_contract.py +++ b/tests/connectors/generate/generate_contract.py @@ -14,7 +14,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.generate.base import BaseGenerateConnector diff --git a/tests/connectors/graph_rag/graph_rag_contract.py b/tests/connectors/graph_rag/graph_rag_contract.py index 303cdb1..6893242 100644 --- a/tests/connectors/graph_rag/graph_rag_contract.py +++ b/tests/connectors/graph_rag/graph_rag_contract.py @@ -18,7 +18,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.graph_rag.base import BaseGraphRagConnector diff --git a/tests/connectors/graph_rag/test_graph_source_chaining.py b/tests/connectors/graph_rag/test_graph_source_chaining.py index 7f5989f..b5d0859 100644 --- a/tests/connectors/graph_rag/test_graph_source_chaining.py +++ b/tests/connectors/graph_rag/test_graph_source_chaining.py @@ -19,7 +19,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.graph_rag.adjacency_graph_rag import AdjacencyGraphRag from rag_integration.feature_groups.connectors.graph_rag.base import BaseGraphRagConnector from rag_integration.feature_groups.connectors.graph_rag.kg_source import TriplesKnowledgeGraph @@ -145,21 +145,21 @@ def test_chained_top_k_applies() -> None: def test_graph_source_with_inline_nodes_raises() -> None: options = _chained_options(extra={AdjacencyGraphRag.NODES: [{"doc_id": "n", "text": "t"}]}) with pytest.raises(ValueError, match="one graph only"): - AdjacencyGraphRag.calculate_feature([], _feature_set(options)) + AdjacencyGraphRag.calculate_feature({}, _feature_set(options)) def test_graph_source_with_inline_edges_raises() -> None: options = _chained_options(extra={AdjacencyGraphRag.EDGES: [["a", "b"]]}) with pytest.raises(ValueError, match="one graph only"): - AdjacencyGraphRag.calculate_feature([], _feature_set(options)) + AdjacencyGraphRag.calculate_feature({}, _feature_set(options)) def test_graph_source_without_upstream_row_raises() -> None: with pytest.raises(ValueError, match="produced no row"): - AdjacencyGraphRag.calculate_feature([], _feature_set(_chained_options())) + AdjacencyGraphRag.calculate_feature({}, _feature_set(_chained_options())) def test_graph_source_with_malformed_payload_raises() -> None: - data = [{TriplesKnowledgeGraph.ROOT_FEATURE_NAME: ["not", "a", "dict"]}] + data = {TriplesKnowledgeGraph.ROOT_FEATURE_NAME: [["not", "a", "dict"]]} with pytest.raises(ValueError, match="nodes"): AdjacencyGraphRag.calculate_feature(data, _feature_set(_chained_options())) diff --git a/tests/connectors/graph_rag/test_kg_source.py b/tests/connectors/graph_rag/test_kg_source.py index 85d7258..3b2851d 100644 --- a/tests/connectors/graph_rag/test_kg_source.py +++ b/tests/connectors/graph_rag/test_kg_source.py @@ -16,7 +16,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.graph_rag.kg_source import ( BaseKnowledgeGraphSource, TriplesKnowledgeGraph, diff --git a/tests/connectors/orchestrator/orchestrator_contract.py b/tests/connectors/orchestrator/orchestrator_contract.py index 8efe7a9..64d55fa 100644 --- a/tests/connectors/orchestrator/orchestrator_contract.py +++ b/tests/connectors/orchestrator/orchestrator_contract.py @@ -15,7 +15,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.orchestrator.base import BaseOrchestratorConnector diff --git a/tests/connectors/rerank/rerank_contract.py b/tests/connectors/rerank/rerank_contract.py index 67c0e8d..c9cf63d 100644 --- a/tests/connectors/rerank/rerank_contract.py +++ b/tests/connectors/rerank/rerank_contract.py @@ -15,7 +15,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.rerank.base import BaseRerankConnector diff --git a/tests/connectors/retrieve/retrieve_contract.py b/tests/connectors/retrieve/retrieve_contract.py index 32db796..4e6faab 100644 --- a/tests/connectors/retrieve/retrieve_contract.py +++ b/tests/connectors/retrieve/retrieve_contract.py @@ -22,7 +22,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.retrieve.base import BaseRetrieveConnector diff --git a/tests/connectors/structured/structured_contract.py b/tests/connectors/structured/structured_contract.py index 3d80496..32f084c 100644 --- a/tests/connectors/structured/structured_contract.py +++ b/tests/connectors/structured/structured_contract.py @@ -18,7 +18,7 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows from rag_integration.feature_groups.connectors.structured.base import BaseStructuredConnector diff --git a/tests/feature_groups/test_columnar.py b/tests/feature_groups/test_columnar.py deleted file mode 100644 index d11a5d2..0000000 --- a/tests/feature_groups/test_columnar.py +++ /dev/null @@ -1,77 +0,0 @@ -"""Tests for the columnar helpers backing every PythonDict feature group. - -mloda 0.9.0 made PythonDict columnar ``dict[str, list]``; ``columnar_to_rows`` is the -row-wise read path in each ``calculate_feature`` and ``homogenize_rows`` keeps source -rows on the uniform key schema the columnar output contract requires. Both are -hand-rolled here (mloda ships ``rows_to_columnar`` but no public inverse), so this -suite pins their edge cases directly. -""" - -from __future__ import annotations - -from typing import Any - -from rag_integration.feature_groups.columnar import columnar_to_rows, homogenize_rows - - -class TestColumnarToRows: - def test_pivots_columnar_dict_to_rows(self) -> None: - data = {"doc_id": ["a", "b"], "text": ["one", "two"]} - assert columnar_to_rows(data) == [ - {"doc_id": "a", "text": "one"}, - {"doc_id": "b", "text": "two"}, - ] - - def test_preserves_column_order_in_rows(self) -> None: - data = {"z": [1], "a": [2]} - assert list(columnar_to_rows(data)[0].keys()) == ["z", "a"] - - def test_list_passes_through_unchanged(self) -> None: - rows = [{"doc_id": "a"}] - assert columnar_to_rows(rows) is rows - - def test_schemaless_empty_dict_yields_empty_list(self) -> None: - assert columnar_to_rows({}) == [] - - def test_schema_bearing_zero_row_dict_yields_empty_list(self) -> None: - assert columnar_to_rows({"doc_id": []}) == [] - - def test_non_dict_non_list_yields_empty_list(self) -> None: - assert columnar_to_rows(None) == [] - assert columnar_to_rows("text") == [] - - def test_none_cell_values_survive_the_pivot(self) -> None: - data = {"doc_id": ["a", "b"], "author": [None, "jane"]} - assert columnar_to_rows(data) == [ - {"doc_id": "a", "author": None}, - {"doc_id": "b", "author": "jane"}, - ] - - -class TestHomogenizeRows: - def test_backfills_missing_keys_with_none(self) -> None: - rows = [{"doc_id": "a"}, {"doc_id": "b", "author": "jane"}] - assert homogenize_rows(rows) == [ - {"doc_id": "a", "author": None}, - {"doc_id": "b", "author": "jane"}, - ] - - def test_key_order_follows_first_occurrence(self) -> None: - rows = [{"b": 1}, {"a": 2}, {"c": 3}] - assert all(list(row.keys()) == ["b", "a", "c"] for row in homogenize_rows(rows)) - - def test_uniform_rows_are_copied_unchanged(self) -> None: - rows = [{"doc_id": "a", "text": "one"}] - result = homogenize_rows(rows) - assert result == rows - assert result[0] is not rows[0] - - def test_empty_input_yields_empty_list(self) -> None: - assert homogenize_rows([]) == [] - - def test_explicit_none_values_are_kept(self) -> None: - rows: list[dict[str, Any]] = [{"doc_id": "a", "author": None}, {"doc_id": "b"}] - assert homogenize_rows(rows) == [ - {"doc_id": "a", "author": None}, - {"doc_id": "b", "author": None}, - ] diff --git a/tests/feature_groups/test_deduplication_base.py b/tests/feature_groups/test_deduplication_base.py index 34a8135..a81f34d 100644 --- a/tests/feature_groups/test_deduplication_base.py +++ b/tests/feature_groups/test_deduplication_base.py @@ -81,25 +81,25 @@ class TestBaseRowDeduplicatorKeepStrategies: def test_keep_longest_selects_largest_item_per_group(self) -> None: """Strategy "longest" keeps the longest string from each duplicate group.""" - data = [{"item": "aa"}, {"item": "aaaa"}, {"item": "bb"}, {"item": "b"}] + data = {"item": ["aa", "aaaa", "bb", "b"]} result = _LenDeduplicator.calculate_feature(data, _features("longest")) assert [row["items__deduped"] for row in result] == ["aaaa", "bb"] def test_keep_largest_uses_length_for_bytes_and_custom_strategy_value(self) -> None: """A subclass with KEEP_LARGEST_STRATEGY="largest" keeps the largest bytes per group.""" - data = [{"item": b"aa"}, {"item": b"aaaa"}, {"item": b"bb"}] + data = {"item": [b"aa", b"aaaa", b"bb"]} result = _LargestLenDeduplicator.calculate_feature(data, _features("largest")) assert [row["items__deduped"] for row in result] == [b"aaaa", b"bb"] def test_keep_first_keeps_only_non_duplicate_rows(self) -> None: """Strategy "first" keeps the first occurrence of each group and drops the duplicates.""" - data = [{"item": "aa"}, {"item": "aaaa"}, {"item": "bb"}, {"item": "b"}] + data = {"item": ["aa", "aaaa", "bb", "b"]} result = _LenDeduplicator.calculate_feature(data, _features("first")) assert [row["items__deduped"] for row in result] == ["aa", "bb"] def test_all_unique_keeps_all_rows_with_duplicate_metadata(self) -> None: """Strategy "all_unique" keeps every row but still annotates duplicate metadata.""" - data = [{"item": "aa"}, {"item": "aaaa"}, {"item": "bb"}, {"item": "b"}] + data = {"item": ["aa", "aaaa", "bb", "b"]} result = _LenDeduplicator.calculate_feature(data, _features("all_unique")) assert [row["items__deduped"] for row in result] == ["aa", "aaaa", "bb", "b"] assert [(row["is_duplicate"], row["duplicate_of"]) for row in result] == [ @@ -109,21 +109,21 @@ def test_all_unique_keeps_all_rows_with_duplicate_metadata(self) -> None: (True, 2), ] - def test_empty_feature_set_returns_data_unchanged(self) -> None: - """An empty FeatureSet is a no-op: the input rows are returned as-is.""" - data = [{"item": "aa"}, {"item": "bb"}] + def test_empty_feature_set_returns_rows_unfiltered(self) -> None: + """An empty FeatureSet is a no-op: the pivoted rows come back unfiltered.""" + data = {"item": ["aa", "bb"]} result = _LenDeduplicator.calculate_feature(data, _features_named()) - assert result is data + assert result == [{"item": "aa"}, {"item": "bb"}] def test_repeated_same_feature_is_processed_once(self) -> None: """The same feature appearing multiple times (a framework duplicate) collapses by name.""" - data = [{"item": "aa"}, {"item": "aaaa"}, {"item": "bb"}] + data = {"item": ["aa", "aaaa", "bb"]} result = _LenDeduplicator.calculate_feature(data, _features_named("items__deduped", "items__deduped")) # "first" strategy: one row survives per duplicate group, under the single feature name. assert [row["items__deduped"] for row in result] == ["aa", "bb"] def test_distinct_features_raise_instead_of_silently_dropping(self) -> None: """Deduplication is undefined for >1 distinct feature (row-filtering + shared metadata).""" - data = [{"item": "aa"}, {"item": "aaaa"}] + data = {"item": ["aa", "aaaa"]} with pytest.raises(ValueError, match="distinct features"): _LenDeduplicator.calculate_feature(data, _features_named("a__deduped", "b__deduped")) diff --git a/tests/feature_groups_evaluation/test_faiss_retrieval_evaluator.py b/tests/feature_groups_evaluation/test_faiss_retrieval_evaluator.py index 3ee2e3e..5a5ac0f 100644 --- a/tests/feature_groups_evaluation/test_faiss_retrieval_evaluator.py +++ b/tests/feature_groups_evaluation/test_faiss_retrieval_evaluator.py @@ -7,11 +7,22 @@ import pytest +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import ( + homogenize_rows, + rows_to_columnar, +) + from rag_integration.feature_groups.evaluation.faiss_retrieval_evaluator import FaissRetrievalEvaluator pytest.importorskip("numpy") pytest.importorskip("faiss") + +def _columnar(rows: List[Dict[str, Any]]) -> Dict[str, List[Any]]: + """Pivot test rows to the columnar shape the framework delivers.""" + return rows_to_columnar(homogenize_rows(rows)) + + # The embedding feature is one level above __indexed _EMBEDDING_FEATURE = "eval_docs__chunked__deduped__embedded" _INDEXED_FEATURE = f"{_EMBEDDING_FEATURE}__indexed" @@ -39,26 +50,28 @@ def _embed(values: List[float]) -> List[float]: return list(v / np.linalg.norm(v)) -def _make_data(emb_feature: str = _EMBEDDING_FEATURE) -> List[Dict[str, Any]]: +def _make_data(emb_feature: str = _EMBEDDING_FEATURE) -> Dict[str, List[Any]]: """Two corpus docs, two queries; query i matches corpus i exactly.""" - return [ - {"doc_id": "d0", "row_type": "corpus", emb_feature: _embed([1.0, 0.0, 0.0]), _INDEXED_FEATURE: 0}, - {"doc_id": "d1", "row_type": "corpus", emb_feature: _embed([0.0, 1.0, 0.0]), _INDEXED_FEATURE: 1}, - { - "doc_id": "q0", - "row_type": "query", - emb_feature: _embed([1.0, 0.0, 0.0]), - _INDEXED_FEATURE: 2, - "relevant_doc_ids": ["d0"], - }, - { - "doc_id": "q1", - "row_type": "query", - emb_feature: _embed([0.0, 1.0, 0.0]), - _INDEXED_FEATURE: 3, - "relevant_doc_ids": ["d1"], - }, - ] + return _columnar( + [ + {"doc_id": "d0", "row_type": "corpus", emb_feature: _embed([1.0, 0.0, 0.0]), _INDEXED_FEATURE: 0}, + {"doc_id": "d1", "row_type": "corpus", emb_feature: _embed([0.0, 1.0, 0.0]), _INDEXED_FEATURE: 1}, + { + "doc_id": "q0", + "row_type": "query", + emb_feature: _embed([1.0, 0.0, 0.0]), + _INDEXED_FEATURE: 2, + "relevant_doc_ids": ["d0"], + }, + { + "doc_id": "q1", + "row_type": "query", + emb_feature: _embed([0.0, 1.0, 0.0]), + _INDEXED_FEATURE: 3, + "relevant_doc_ids": ["d1"], + }, + ] + ) class TestFaissRetrievalEvaluator: @@ -80,18 +93,20 @@ def test_perfect_recall(self) -> None: def test_zero_recall(self) -> None: """Orthogonal query embeddings never match → Recall@1 = 0.0.""" emb = _EMBEDDING_FEATURE - data = [ - {"doc_id": "d0", "row_type": "corpus", emb: _embed([1.0, 0.0, 0.0]), _INDEXED_FEATURE: 0}, - {"doc_id": "d1", "row_type": "corpus", emb: _embed([0.0, 1.0, 0.0]), _INDEXED_FEATURE: 1}, - # Query points to d0 but its embedding is orthogonal to d0 - { - "doc_id": "q0", - "row_type": "query", - emb: _embed([0.0, 1.0, 0.0]), # matches d1, not d0 - _INDEXED_FEATURE: 2, - "relevant_doc_ids": ["d0"], - }, - ] + data = _columnar( + [ + {"doc_id": "d0", "row_type": "corpus", emb: _embed([1.0, 0.0, 0.0]), _INDEXED_FEATURE: 0}, + {"doc_id": "d1", "row_type": "corpus", emb: _embed([0.0, 1.0, 0.0]), _INDEXED_FEATURE: 1}, + # Query points to d0 but its embedding is orthogonal to d0 + { + "doc_id": "q0", + "row_type": "query", + emb: _embed([0.0, 1.0, 0.0]), # matches d1, not d0 + _INDEXED_FEATURE: 2, + "relevant_doc_ids": ["d0"], + }, + ] + ) features = _make_features() with patch.object(FaissRetrievalEvaluator, "_extract_source_features", return_value=[_INDEXED_FEATURE]): @@ -103,36 +118,38 @@ def test_chunked_doc_recall(self) -> None: """Corpus doc split into 2 chunks; query matches chunk 1 → doc-level Recall@1 = 1.0.""" emb = _EMBEDDING_FEATURE # d0 has two chunks, both share doc_id="d0" - data = [ - { - "doc_id": "d0", - "chunk_id": "d0_chunk_0", - "row_type": "corpus", - emb: _embed([1.0, 0.0, 0.0]), - _INDEXED_FEATURE: 0, - }, - { - "doc_id": "d0", - "chunk_id": "d0_chunk_1", - "row_type": "corpus", - emb: _embed([0.9, 0.1, 0.0]), - _INDEXED_FEATURE: 1, - }, - { - "doc_id": "d1", - "chunk_id": "d1_chunk_0", - "row_type": "corpus", - emb: _embed([0.0, 1.0, 0.0]), - _INDEXED_FEATURE: 2, - }, - { - "doc_id": "q0", - "row_type": "query", - emb: _embed([0.9, 0.1, 0.0]), # closest to d0_chunk_1 - _INDEXED_FEATURE: 3, - "relevant_doc_ids": ["d0"], - }, - ] + data = _columnar( + [ + { + "doc_id": "d0", + "chunk_id": "d0_chunk_0", + "row_type": "corpus", + emb: _embed([1.0, 0.0, 0.0]), + _INDEXED_FEATURE: 0, + }, + { + "doc_id": "d0", + "chunk_id": "d0_chunk_1", + "row_type": "corpus", + emb: _embed([0.9, 0.1, 0.0]), + _INDEXED_FEATURE: 1, + }, + { + "doc_id": "d1", + "chunk_id": "d1_chunk_0", + "row_type": "corpus", + emb: _embed([0.0, 1.0, 0.0]), + _INDEXED_FEATURE: 2, + }, + { + "doc_id": "q0", + "row_type": "query", + emb: _embed([0.9, 0.1, 0.0]), # closest to d0_chunk_1 + _INDEXED_FEATURE: 3, + "relevant_doc_ids": ["d0"], + }, + ] + ) features = _make_features() with patch.object(FaissRetrievalEvaluator, "_extract_source_features", return_value=[_INDEXED_FEATURE]): @@ -144,15 +161,17 @@ def test_chunked_doc_recall(self) -> None: def test_empty_corpus_returns_zero_metrics(self) -> None: """No corpus rows → returns zeros gracefully.""" emb = _EMBEDDING_FEATURE - data = [ - { - "doc_id": "q0", - "row_type": "query", - emb: _embed([1.0, 0.0, 0.0]), - _INDEXED_FEATURE: 0, - "relevant_doc_ids": ["d0"], - } - ] + data = _columnar( + [ + { + "doc_id": "q0", + "row_type": "query", + emb: _embed([1.0, 0.0, 0.0]), + _INDEXED_FEATURE: 0, + "relevant_doc_ids": ["d0"], + } + ] + ) features = _make_features() with patch.object(FaissRetrievalEvaluator, "_extract_source_features", return_value=[_INDEXED_FEATURE]): diff --git a/tests/feature_groups_evaluation/test_retrieval_evaluator.py b/tests/feature_groups_evaluation/test_retrieval_evaluator.py index 964acf6..951b238 100644 --- a/tests/feature_groups_evaluation/test_retrieval_evaluator.py +++ b/tests/feature_groups_evaluation/test_retrieval_evaluator.py @@ -7,11 +7,21 @@ import pytest +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import ( + homogenize_rows, + rows_to_columnar, +) + from rag_integration.feature_groups.evaluation.retrieval_evaluator import RetrievalEvaluator pytest.importorskip("numpy") +def _columnar(rows: List[Dict[str, Any]]) -> Dict[str, List[Any]]: + """Pivot test rows to the columnar shape the framework delivers.""" + return rows_to_columnar(homogenize_rows(rows)) + + def _make_features(source_name: str = "eval_docs__embedded") -> Any: """Build a minimal FeatureSet mock.""" feature = MagicMock() @@ -34,14 +44,16 @@ def _embed(values: List[float]) -> List[float]: return list(v / np.linalg.norm(v)) -def _make_data(source: str = "eval_docs__embedded") -> List[Dict[str, Any]]: +def _make_data(source: str = "eval_docs__embedded") -> Dict[str, List[Any]]: """Two corpus docs, two queries; query 0 matches corpus 0, query 1 matches corpus 1.""" - return [ - {"doc_id": "d0", "row_type": "corpus", source: _embed([1.0, 0.0, 0.0])}, - {"doc_id": "d1", "row_type": "corpus", source: _embed([0.0, 1.0, 0.0])}, - {"doc_id": "q0", "row_type": "query", source: _embed([1.0, 0.0, 0.0]), "relevant_doc_ids": ["d0"]}, - {"doc_id": "q1", "row_type": "query", source: _embed([0.0, 1.0, 0.0]), "relevant_doc_ids": ["d1"]}, - ] + return _columnar( + [ + {"doc_id": "d0", "row_type": "corpus", source: _embed([1.0, 0.0, 0.0])}, + {"doc_id": "d1", "row_type": "corpus", source: _embed([0.0, 1.0, 0.0])}, + {"doc_id": "q0", "row_type": "query", source: _embed([1.0, 0.0, 0.0]), "relevant_doc_ids": ["d0"]}, + {"doc_id": "q1", "row_type": "query", source: _embed([0.0, 1.0, 0.0]), "relevant_doc_ids": ["d1"]}, + ] + ) class TestRetrievalEvaluator: @@ -60,10 +72,17 @@ def test_perfect_recall(self) -> None: def test_zero_recall(self) -> None: source = "eval_docs__embedded" - data = [ - {"doc_id": "d0", "row_type": "corpus", source: _embed([1.0, 0.0, 0.0])}, - {"doc_id": "q0", "row_type": "query", source: _embed([0.0, 1.0, 0.0]), "relevant_doc_ids": ["d_missing"]}, - ] + data = _columnar( + [ + {"doc_id": "d0", "row_type": "corpus", source: _embed([1.0, 0.0, 0.0])}, + { + "doc_id": "q0", + "row_type": "query", + source: _embed([0.0, 1.0, 0.0]), + "relevant_doc_ids": ["d_missing"], + }, + ] + ) features = _make_features(source) with patch.object(RetrievalEvaluator, "_extract_source_features", return_value=[source]): @@ -84,9 +103,11 @@ def test_counts_returned(self) -> None: def test_empty_corpus(self) -> None: source = "eval_docs__embedded" - data = [ - {"doc_id": "q0", "row_type": "query", source: _embed([1.0, 0.0]), "relevant_doc_ids": ["d0"]}, - ] + data = _columnar( + [ + {"doc_id": "q0", "row_type": "query", source: _embed([1.0, 0.0]), "relevant_doc_ids": ["d0"]}, + ] + ) features = _make_features(source) with patch.object(RetrievalEvaluator, "_extract_source_features", return_value=[source]): diff --git a/tests/integration/helpers.py b/tests/integration/helpers.py index e6dd416..4b3c934 100644 --- a/tests/integration/helpers.py +++ b/tests/integration/helpers.py @@ -9,28 +9,29 @@ PythonDictFramework, ) -from rag_integration.feature_groups.columnar import columnar_to_rows +from mloda_plugins.compute_framework.base_implementations.python_dict.python_dict_utils import columnar_to_rows def flatten_result(result: Any) -> List[Dict[str, Any]]: """Unwrap a nested mlodaAPI result to a flat list of row dicts. - PythonDict partitions are columnar ``dict[str, list]`` (mloda 0.9.0); pivot the - first partition back to rows. Accepts a raw result list, a single wrapped - partition, a bare columnar dict, or an already-row-wise list. A leading dict is - treated as a partition only when every value is a list; scalar-valued dicts are - rows and pass through unchanged. + PythonDict partitions are columnar ``dict[str, list]``; pivot the first + partition back to rows. Accepts a raw result list, a single wrapped partition, + a bare columnar dict, or an already-row-wise list. A leading dict is treated as + a partition only when every value is a list; scalar-valued dicts are rows and + pass through unchanged. mloda's ``columnar_to_rows`` raises on non-columnar + input, so the shape dispatch lives here and only genuine columnar dicts reach it. """ if isinstance(result, dict): return columnar_to_rows(result) if result and isinstance(result[0], list): - return columnar_to_rows(result[0]) + return list(result[0]) if result and isinstance(result[0], dict): first = result[0] if not first or all(isinstance(value, list) for value in first.values()): return columnar_to_rows(first) return list(result) - return columnar_to_rows(result) + return list(result) if isinstance(result, list) else [] def get_results_by_feature(raw_result: List[Any], feature_names: List[str]) -> Dict[str, List[Dict[str, Any]]]: diff --git a/tests/integration/test_helpers.py b/tests/integration/test_helpers.py new file mode 100644 index 0000000..ac34645 --- /dev/null +++ b/tests/integration/test_helpers.py @@ -0,0 +1,38 @@ +"""Tests for the ``flatten_result`` shape dispatch. + +mloda 0.10.0's ``columnar_to_rows`` raises on non-columnar input, so the tolerant +shape dispatch (row-wise passthrough, empty fallback) lives in ``flatten_result`` +itself. This suite pins that dispatch. +""" + +from __future__ import annotations + +from tests.integration.helpers import flatten_result + + +class TestFlattenResult: + def test_bare_columnar_dict_pivots_to_rows(self) -> None: + data = {"doc_id": ["a", "b"], "text": ["one", "two"]} + assert flatten_result(data) == [ + {"doc_id": "a", "text": "one"}, + {"doc_id": "b", "text": "two"}, + ] + + def test_wrapped_columnar_partition_pivots_to_rows(self) -> None: + assert flatten_result([{"doc_id": ["a"]}]) == [{"doc_id": "a"}] + + def test_wrapped_row_list_is_kept_row_wise(self) -> None: + rows = [{"doc_id": "a"}, {"doc_id": "b"}] + assert flatten_result([rows]) == rows + + def test_scalar_valued_leading_dict_passes_through_as_rows(self) -> None: + rows = [{"doc_id": "a", "score": 0.5}, {"doc_id": "b", "score": 0.7}] + assert flatten_result(rows) == rows + + def test_empty_leading_dict_yields_empty_list(self) -> None: + assert flatten_result([{}]) == [] + + def test_empty_inputs_yield_empty_list(self) -> None: + assert flatten_result([]) == [] + assert flatten_result({}) == [] + assert flatten_result(None) == []