Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion cli/eval_demo.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion cli/rag_demo.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion cli/swap_demo.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion demo.ipynb
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
39 changes: 0 additions & 39 deletions rag_integration/feature_groups/columnar.py

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion rag_integration/feature_groups/datasets/image/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
2 changes: 1 addition & 1 deletion rag_integration/feature_groups/datasets/text/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
4 changes: 2 additions & 2 deletions rag_integration/feature_groups/deduplication_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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/)
Expand Down Expand Up @@ -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.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions rag_integration/feature_groups/rag_pipeline/chunking/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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 = []

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
4 changes: 2 additions & 2 deletions rag_integration/feature_groups/rag_pipeline/embedding/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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()

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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()

Expand Down
2 changes: 1 addition & 1 deletion tests/connectors/generate/generate_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
2 changes: 1 addition & 1 deletion tests/connectors/graph_rag/graph_rag_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
10 changes: 5 additions & 5 deletions tests/connectors/graph_rag/test_graph_source_chaining.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()))
2 changes: 1 addition & 1 deletion tests/connectors/graph_rag/test_kg_source.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion tests/connectors/orchestrator/orchestrator_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down
Loading
Loading