From 918830a82181cd4547e72da3e1f08d3a141b0270 Mon Sep 17 00:00:00 2001 From: Karunasagar Mohansundar <52268863+Karunasagar12@users.noreply.github.com> Date: Wed, 12 Aug 2026 05:19:54 +0530 Subject: [PATCH] fix(pipeline): resolve broken import and missing `run()` in `PipelineWithProvenance` (#862) * fix(pipeline): resolve broken import and missing run() in PipelineWithProvenance Fix two bugs in pipeline_provenance.py: 1. Wrong import path: `from .pipeline import Pipeline` fails because `semantica/pipeline/pipeline.py` does not exist. Pipeline lives in `pipeline_builder.py`. Fixed to `from .pipeline_builder import Pipeline`. 2. Pipeline dataclass has no run() method. PipelineWithProvenance.run() now delegates to ExecutionEngine.execute_pipeline(), which is the intended execution path for built pipelines. Additional changes: - Constructor now accepts a built Pipeline instance (breaking the previous unusable API that tried to instantiate a dataclass with **config). - Replace deprecated datetime.utcnow() with datetime.now(timezone.utc). - Add test suite covering import, instantiation, execution, attribute delegation, and provenance graceful degradation. Fixes #858 * test: address Qodo review findings - Remove redundant test_import_succeeds (module-level import already guards against import regression at collection time). - Fix test_provenance_disabled_when_import_fails to deterministically simulate ImportError via sys.modules patch and assert provenance is actually toggled off (runner.provenance is False). * fix(pipeline): update provenance callers for Pipeline API --------- Co-authored-by: Sameer Kadam Co-authored-by: Russell Jurney --- CHANGELOG.md | 6 ++ semantica/pipeline/pipeline_provenance.py | 61 ++++++++++++------ semantica/provenance/provenance_usage.md | 10 ++- tests/pipeline/test_pipeline_provenance.py | 64 +++++++++++++++++++ .../provenance/test_all_provenance_modules.py | 25 +++++++- .../provenance/test_provenance_edge_cases.py | 7 +- .../test_real_module_integration.py | 15 +++-- 7 files changed, 156 insertions(+), 32 deletions(-) create mode 100644 tests/pipeline/test_pipeline_provenance.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 34d400d0..8efdcd18 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -71,6 +71,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- **`PipelineWithProvenance` raised `ModuleNotFoundError` on import and `AttributeError` on `.run()`** (#858, closes #858) by @Karunasagar12 + - `from .pipeline import Pipeline` failed because `semantica/pipeline/pipeline.py` does not exist; corrected to `from .pipeline_builder import Pipeline` + - `.run()` called `self._pipeline.run()` on the `Pipeline` dataclass, which has no such method; replaced with `self._engine.execute_pipeline(self._pipeline, ...)` delegating to `ExecutionEngine` + - Constructor now accepts a built `Pipeline` instance (from `PipelineBuilder.build()`) instead of `**config`; the old `Pipeline(**config)` internal construction was invalid and never functional + - Replaced deprecated `datetime.utcnow()` with `datetime.now(timezone.utc)` in `run()` + - **`VectorStore.search_vectors()` returned inconsistent result shapes across backend implementations** (#853, closes #845) by @Sameer6305, reviewed by @KaifAhmad1 - Every built-in backend (FAISS, Milvus, pgvector, Pinecone, Qdrant, SQLite-vec, Weaviate, in-memory) now returns the same canonical `SearchResult` shape (`id`, `score`, `metadata`, `vector`, `distance`), instead of some backends omitting `vector`/`metadata`/`distance` or, for Weaviate, returning a backend-specific `properties` key instead of `metadata` - Added a `SearchResult` `TypedDict` (`semantica/vector_store/vector_store.py`, exported from `semantica.vector_store`) documenting the contract; `metadata` now always defaults to `{}` rather than being absent, and `id` accepts `Union[str, int]` to accommodate Milvus/Qdrant's native integer IDs without casting diff --git a/semantica/pipeline/pipeline_provenance.py b/semantica/pipeline/pipeline_provenance.py index b5c7c99b..a77d7c9a 100644 --- a/semantica/pipeline/pipeline_provenance.py +++ b/semantica/pipeline/pipeline_provenance.py @@ -6,9 +6,14 @@ capturing all steps, inputs, outputs, and transformations. Usage: from semantica.pipeline.pipeline_provenance import PipelineWithProvenance - - pipeline = PipelineWithProvenance(provenance=True) - result = pipeline.run(data) + from semantica.pipeline import PipelineBuilder + + builder = PipelineBuilder() + builder.add_step("ingest", "file_ingest") + pipeline = builder.build("my_pipeline") + + runner = PipelineWithProvenance(pipeline, provenance=True) + result = runner.run(data) # Tracks all pipeline steps with complete lineage Author: Semantica Contributors @@ -16,26 +21,37 @@ License: MIT """ from typing import Optional, Any, Dict, List -from datetime import datetime +from datetime import datetime, timezone import uuid import time +from .pipeline_builder import Pipeline +from .execution_engine import ExecutionEngine + class PipelineWithProvenance: """Pipeline executor with complete provenance tracking.""" - + def __init__( self, + pipeline: Pipeline, provenance: bool = False, agent_id: Optional[str] = None, is_automated: bool = True, - **config, + **engine_config, ): - """Initialize pipeline with optional provenance.""" - from .pipeline import Pipeline + """Initialize provenance-tracked pipeline runner. + Args: + pipeline: A built Pipeline instance (from PipelineBuilder.build()). + provenance: Whether to record provenance metadata. + agent_id: Identifier for the agent running the pipeline. + is_automated: Whether the execution is automated (vs. human-triggered). + **engine_config: Extra keyword arguments forwarded to ExecutionEngine. + """ + self._pipeline = pipeline + self._engine = ExecutionEngine(**engine_config) self.provenance = provenance - self._pipeline = Pipeline(**config) self._prov_manager = None self._agent_id = agent_id or self.__class__.__name__ self._is_automated = is_automated @@ -47,15 +63,24 @@ class PipelineWithProvenance: except ImportError: self.provenance = False - def run(self, data: Any, source: Optional[str] = None, **kwargs): - """Run pipeline with provenance tracking.""" + def run(self, data: Any = None, source: Optional[str] = None, **kwargs): + """Run pipeline with provenance tracking. + + Args: + data: Input data to feed into the pipeline. + source: Provenance source label (defaults to "pipeline_execution"). + **kwargs: Extra options forwarded to ExecutionEngine.execute_pipeline(). + + Returns: + ExecutionResult from the engine. + """ pipeline_id = f"pipeline_{uuid.uuid4().hex[:8]}" start_time = time.time() - activity_started_at_time = datetime.utcnow().isoformat() + activity_started_at_time = datetime.now(timezone.utc).isoformat() - result = self._pipeline.run(data, **kwargs) + result = self._engine.execute_pipeline(self._pipeline, data=data, **kwargs) elapsed = time.time() - start_time - activity_ended_at_time = datetime.utcnow().isoformat() + activity_ended_at_time = datetime.now(timezone.utc).isoformat() if self.provenance and self._prov_manager: self._prov_manager.track_entity( @@ -69,14 +94,14 @@ class PipelineWithProvenance: activity_started_at_time=activity_started_at_time, activity_ended_at_time=activity_ended_at_time, metadata={ - "steps": len(self._pipeline.steps) if hasattr(self._pipeline, 'steps') else 0, + "steps": len(self._pipeline.steps), "duration_seconds": elapsed, - "status": "completed" + "status": "completed" if result.success else "failed", } ) - + return result - + def __getattr__(self, name): return getattr(self._pipeline, name) diff --git a/semantica/provenance/provenance_usage.md b/semantica/provenance/provenance_usage.md index db667038..7113826f 100644 --- a/semantica/provenance/provenance_usage.md +++ b/semantica/provenance/provenance_usage.md @@ -399,12 +399,16 @@ response = llm.generate("What is artificial intelligence?") ```python from semantica.pipeline.pipeline_provenance import PipelineWithProvenance +from semantica.pipeline import PipelineBuilder -# Create pipeline with provenance -pipeline = PipelineWithProvenance(provenance=True) +builder = PipelineBuilder() +builder.add_step("ingest", "file_ingest") +pipeline = builder.build("my_pipeline") + +runner = PipelineWithProvenance(pipeline, provenance=True) # Run pipeline - all steps tracked -result = pipeline.run( +result = runner.run( data=input_data, source="input_file.json" ) diff --git a/tests/pipeline/test_pipeline_provenance.py b/tests/pipeline/test_pipeline_provenance.py new file mode 100644 index 00000000..1756df08 --- /dev/null +++ b/tests/pipeline/test_pipeline_provenance.py @@ -0,0 +1,64 @@ +"""Tests for PipelineWithProvenance. + +Verifies that: +1. The import path is correct (no ModuleNotFoundError). +2. PipelineWithProvenance accepts a built Pipeline and runs it via ExecutionEngine. +3. Provenance tracking gracefully degrades when the provenance package is absent. +""" + +import sys +from unittest.mock import patch + +import pytest + +from semantica.pipeline import PipelineBuilder +from semantica.pipeline.pipeline_provenance import PipelineWithProvenance +from semantica.pipeline.execution_engine import ExecutionResult + + +class TestPipelineWithProvenance: + """Tests for PipelineWithProvenance.""" + + @pytest.fixture + def simple_pipeline(self): + """Build a minimal two-step pipeline for testing.""" + builder = PipelineBuilder() + builder.add_step("ingest", "file_ingest") + builder.add_step("parse", "document_parse") + return builder.build("test_provenance_pipeline") + + def test_instantiation_with_pipeline(self, simple_pipeline): + """Should accept a built Pipeline instance.""" + runner = PipelineWithProvenance(simple_pipeline, provenance=False) + assert runner._pipeline is simple_pipeline + + def test_run_returns_execution_result(self, simple_pipeline): + """run() should delegate to ExecutionEngine and return an ExecutionResult.""" + runner = PipelineWithProvenance(simple_pipeline, provenance=False) + result = runner.run() + assert isinstance(result, ExecutionResult) + assert result.success is True + + def test_getattr_delegates_to_pipeline(self, simple_pipeline): + """Attribute access should fall through to the wrapped Pipeline.""" + runner = PipelineWithProvenance(simple_pipeline, provenance=False) + assert runner.name == "test_provenance_pipeline" + assert len(runner.steps) == 2 + + def test_provenance_disabled_when_import_fails(self, simple_pipeline): + """When semantica.provenance is unavailable, provenance should be disabled.""" + # Force the provenance import to raise ImportError + with patch.dict(sys.modules, {"semantica.provenance": None}): + runner = PipelineWithProvenance(simple_pipeline, provenance=True) + assert runner.provenance is False + assert runner._prov_manager is None + # Should still execute successfully without provenance + result = runner.run() + assert isinstance(result, ExecutionResult) + assert result.success is True + + def test_run_with_data(self, simple_pipeline): + """run() should accept data and kwargs without error.""" + runner = PipelineWithProvenance(simple_pipeline, provenance=False) + result = runner.run(data={"key": "value"}) + assert isinstance(result, ExecutionResult) diff --git a/tests/provenance/test_all_provenance_modules.py b/tests/provenance/test_all_provenance_modules.py index 0d5e5bd4..f0dc25e9 100644 --- a/tests/provenance/test_all_provenance_modules.py +++ b/tests/provenance/test_all_provenance_modules.py @@ -167,7 +167,6 @@ class TestProvenanceEnabledDisabled: """Test all modules work with provenance=False.""" modules_to_test = [ ('semantica.context.context_provenance', 'ContextManagerWithProvenance'), - ('semantica.pipeline.pipeline_provenance', 'PipelineWithProvenance'), ] for module_path, class_name in modules_to_test: @@ -183,7 +182,6 @@ class TestProvenanceEnabledDisabled: """Test all modules work with provenance=True.""" modules_to_test = [ ('semantica.context.context_provenance', 'ContextManagerWithProvenance'), - ('semantica.pipeline.pipeline_provenance', 'PipelineWithProvenance'), ] for module_path, class_name in modules_to_test: @@ -195,6 +193,23 @@ class TestProvenanceEnabledDisabled: except ImportError: pytest.skip(f"{module_path} not available") + def test_pipeline_with_provenance_supports_provenance_flag(self): + """PipelineWithProvenance accepts provenance=True/False. + + PipelineWithProvenance requires a built Pipeline instance (unlike + other *WithProvenance wrappers that own their internal state), so it + cannot participate in the generic no-argument constructor loop above. + """ + try: + from semantica.pipeline.pipeline_builder import Pipeline + from semantica.pipeline.pipeline_provenance import PipelineWithProvenance + + pipeline = Pipeline(name="compat_test") + assert PipelineWithProvenance(pipeline, provenance=False).provenance is False + assert PipelineWithProvenance(pipeline, provenance=True).provenance is True + except ImportError: + pytest.skip("pipeline_provenance not available") + class TestAllModulesEdgeCases: """Test edge cases across all provenance modules.""" @@ -219,10 +234,14 @@ class TestAllModulesEdgeCases: """Test each module has independent provenance manager.""" try: from semantica.context.context_provenance import ContextManagerWithProvenance + from semantica.pipeline.pipeline_builder import Pipeline from semantica.pipeline.pipeline_provenance import PipelineWithProvenance ctx = ContextManagerWithProvenance(provenance=True) - pipe = PipelineWithProvenance(provenance=True) + pipe = PipelineWithProvenance( + Pipeline(name="independence_test"), + provenance=True, + ) # Each should have its own manager assert ctx._prov_manager is not None diff --git a/tests/provenance/test_provenance_edge_cases.py b/tests/provenance/test_provenance_edge_cases.py index 9d0c1f0f..92329482 100644 --- a/tests/provenance/test_provenance_edge_cases.py +++ b/tests/provenance/test_provenance_edge_cases.py @@ -154,10 +154,13 @@ class TestModuleSpecificEdgeCases: def test_pipeline_with_empty_data(self): """Test pipeline with empty data.""" try: + from semantica.pipeline.pipeline_builder import Pipeline from semantica.pipeline.pipeline_provenance import PipelineWithProvenance - pipeline = PipelineWithProvenance(provenance=True) + + pipeline = Pipeline(name="edge_case_test") + runner = PipelineWithProvenance(pipeline, provenance=True) # Should handle empty data - assert pipeline is not None + assert runner is not None except ImportError: pytest.skip("Pipeline not available") diff --git a/tests/provenance/test_real_module_integration.py b/tests/provenance/test_real_module_integration.py index 5dceebed..cd177438 100644 --- a/tests/provenance/test_real_module_integration.py +++ b/tests/provenance/test_real_module_integration.py @@ -29,14 +29,17 @@ class TestRealModuleIntegration: def test_pipeline_real_execution_tracking(self): """Test pipeline tracks execution with provenance.""" try: + from semantica.pipeline.pipeline_builder import Pipeline from semantica.pipeline.pipeline_provenance import PipelineWithProvenance - - pipeline = PipelineWithProvenance(provenance=True) - + + pipeline = Pipeline(name="provenance_tracking_test") + runner = PipelineWithProvenance(pipeline, provenance=True) + # Verify provenance setup - assert pipeline.provenance is True - assert pipeline._prov_manager is not None - + assert runner.provenance is True + assert runner._prov_manager is not None + assert isinstance(runner._prov_manager, ProvenanceManager) + except ImportError: pytest.skip("Pipeline not available")