Spaces:
Sleeping
Sleeping
| """Tests for Pinecone RAG backend.""" | |
| import time | |
| from unittest.mock import MagicMock, patch | |
| import pytest | |
| from langchain_core.documents import Document | |
| from app.core.config import settings | |
| from app.core.schemas import TeamRole | |
| def _make_pinecone_service( | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| index_upsert=None, | |
| index_delete=None, | |
| index_query=None, | |
| ): | |
| """Helper to create a fully initialized PineconeRAGService. | |
| When index_upsert, index_delete, or index_query are provided, they are | |
| set as ``side_effect`` on the MagicMock so call_count and other mock | |
| attributes remain accessible. | |
| """ | |
| mock_embeddings = MagicMock() | |
| mock_embeddings.embed_query.return_value = [0.1, 0.2, 0.3] | |
| mock_embeddings.embed_documents.return_value = [[0.1, 0.2, 0.3]] | |
| mock_get_embeddings.return_value = mock_embeddings | |
| mock_index = MagicMock() | |
| if index_upsert is not None: | |
| mock_index.upsert.side_effect = index_upsert | |
| if index_delete is not None: | |
| mock_index.delete.side_effect = index_delete | |
| if index_query is not None: | |
| mock_index.query.side_effect = index_query | |
| mock_pc = MagicMock() | |
| mock_pc.has_index.return_value = True | |
| mock_pc.Index.return_value = mock_index | |
| mock_pc.describe_index.return_value = {"dimension": 3} | |
| mock_pinecone_cls.return_value = mock_pc | |
| with ( | |
| patch.object(settings, "pinecone_api_key", "test-key"), | |
| patch.object(settings, "pinecone_index", "multi-agent-index"), | |
| ): | |
| from app.core.pinecone_rag import PineconeRAGService | |
| return PineconeRAGService(), mock_index | |
| class TestPineconeRAGContracts: | |
| """Core behavior tests for PineconeRAGService.""" | |
| def test_namespace_map_contains_role_mappings(self) -> None: | |
| """Role mapping should match expected Pinecone namespace names.""" | |
| from app.core.pinecone_rag import ROLE_NAMESPACE_MAP | |
| assert ROLE_NAMESPACE_MAP[TeamRole.PRODUCT_OWNER] == "rag_product_owner" | |
| assert ROLE_NAMESPACE_MAP[TeamRole.BUSINESS_ANALYST] == "rag_business_analyst" | |
| assert ROLE_NAMESPACE_MAP[TeamRole.SPEC_COORDINATOR] is None | |
| def test_namespace_map_covers_all_roles(self) -> None: | |
| """Every TeamRole must have an entry in ROLE_NAMESPACE_MAP (DS-08).""" | |
| from app.core.pinecone_rag import ROLE_NAMESPACE_MAP | |
| for role in TeamRole: | |
| # Must not raise KeyError | |
| ns = ROLE_NAMESPACE_MAP[role] | |
| # Must be either a string or None (both acceptable) | |
| assert ns is None or isinstance(ns, str) | |
| def test_retrieve_returns_documents_from_role_namespace( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """retrieve() should query namespace and convert matches to Documents.""" | |
| from app.core.pinecone_rag import PineconeRAGService | |
| mock_embeddings = MagicMock() | |
| mock_embeddings.embed_query.return_value = [0.1, 0.2, 0.3] | |
| mock_get_embeddings.return_value = mock_embeddings | |
| mock_index = MagicMock() | |
| mock_index.query.return_value = { | |
| "matches": [ | |
| { | |
| "id": "doc-1", | |
| "score": 0.93, | |
| "metadata": { | |
| "content": "Product owner context", | |
| "source": "role_playbook.txt", | |
| "role": "product_owner", | |
| }, | |
| } | |
| ] | |
| } | |
| mock_pc = MagicMock() | |
| mock_pc.has_index.return_value = True | |
| mock_pc.Index.return_value = mock_index | |
| mock_pc.describe_index.return_value = {"dimension": 3} | |
| mock_pinecone_cls.return_value = mock_pc | |
| with patch.dict( | |
| "os.environ", | |
| { | |
| "PINECONE_API_KEY": "test-key", | |
| "PINECONE_INDEX": "multi-agent-index", | |
| }, | |
| clear=False, | |
| ): | |
| service = PineconeRAGService() | |
| docs = service.retrieve("need PRD guidance", TeamRole.PRODUCT_OWNER, k=1) | |
| assert len(docs) == 1 | |
| assert docs[0].page_content == "Product owner context" | |
| assert docs[0].metadata["source"] == "role_playbook.txt" | |
| mock_index.query.assert_called_once() | |
| def test_health_check_reports_not_configured_without_env( | |
| self, mock_get_embeddings | |
| ) -> None: | |
| """health_check() should report not_configured when env vars missing.""" | |
| from app.core.pinecone_rag import PineconeRAGService | |
| mock_get_embeddings.return_value = MagicMock() | |
| with ( | |
| patch.object(settings, "pinecone_api_key", ""), | |
| patch.object(settings, "pinecone_index", ""), | |
| ): | |
| service = PineconeRAGService() | |
| health = service.health_check() | |
| assert health["status"] == "not_configured" | |
| class TestPineconeInitRetry: | |
| """DS-04: Pinecone init should retry on transient failures.""" | |
| def test_init_retries_and_succeeds_on_second_attempt( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """Init should retry when Pinecone raises transiently, then succeed.""" | |
| from app.core.pinecone_rag import PineconeRAGService | |
| call_count = 0 | |
| def pinecone_side_effect(*args, **kwargs): | |
| nonlocal call_count | |
| call_count += 1 | |
| if call_count == 1: | |
| raise ConnectionError("Transient failure on first attempt") | |
| mock_pc = MagicMock() | |
| mock_pc.has_index.return_value = True | |
| mock_index = MagicMock() | |
| mock_pc.Index.return_value = mock_index | |
| mock_pc.describe_index.return_value = {"dimension": 3} | |
| return mock_pc | |
| mock_pinecone_cls.side_effect = pinecone_side_effect | |
| mock_embeddings = MagicMock() | |
| mock_embeddings.embed_query.return_value = [0.1, 0.2, 0.3] | |
| mock_get_embeddings.return_value = mock_embeddings | |
| with ( | |
| patch.object(settings, "pinecone_api_key", "test-key"), | |
| patch.object(settings, "pinecone_index", "multi-agent-index"), | |
| ): | |
| service = PineconeRAGService() | |
| assert service.is_available() | |
| assert call_count >= 2 | |
| def test_init_always_fails_graceful_fallback( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """Init should not crash after all retries exhausted - graceful fallback.""" | |
| from app.core.pinecone_rag import PineconeRAGService | |
| mock_pinecone_cls.side_effect = ConnectionError("Pinecone permanently down") | |
| mock_get_embeddings.return_value = MagicMock() | |
| with ( | |
| patch.object(settings, "pinecone_api_key", "test-key"), | |
| patch.object(settings, "pinecone_index", "multi-agent-index"), | |
| ): | |
| service = PineconeRAGService() | |
| assert not service.is_available() | |
| health = service.health_check() | |
| assert health["status"] == "disconnected" | |
| def test_init_time_is_bounded_by_retries( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """After permanent failure, init should return in bounded time (retries done).""" | |
| from app.core.pinecone_rag import PineconeRAGService | |
| mock_pinecone_cls.side_effect = ConnectionError("Pinecone permanently down") | |
| mock_get_embeddings.return_value = MagicMock() | |
| start = time.time() | |
| with ( | |
| patch.object(settings, "pinecone_api_key", "test-key"), | |
| patch.object(settings, "pinecone_index", "multi-agent-index"), | |
| ): | |
| service = PineconeRAGService() | |
| elapsed = time.time() - start | |
| assert not service.is_available() | |
| assert elapsed < 30 # Sanity check: shouldn't hang | |
| class TestPineconeWritePathRetry: | |
| """DS-05: Pinecone write path should retry on transient failures.""" | |
| async def test_upsert_retries_and_succeeds_on_second_attempt( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """_add_documents_sync should retry upsert on transient failure, then succeed.""" | |
| call_count = 0 | |
| def upsert_side_effect(*args, **kwargs): | |
| nonlocal call_count | |
| call_count += 1 | |
| if call_count == 1: | |
| raise ConnectionError("Transient upsert failure") | |
| return {"upserted_count": 1} | |
| service, mock_index = _make_pinecone_service( | |
| mock_pinecone_cls, mock_get_embeddings, | |
| index_upsert=upsert_side_effect, | |
| ) | |
| docs = [Document(page_content="Test content", metadata={"source": "test"})] | |
| result = await service.add_documents(docs, TeamRole.PRODUCT_OWNER) | |
| assert len(result) == 1 | |
| assert call_count >= 2 | |
| assert mock_index.upsert.call_count >= 2 | |
| async def test_upsert_always_fails_raises_error( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """_add_documents_sync should raise after all retries exhausted.""" | |
| service, mock_index = _make_pinecone_service( | |
| mock_pinecone_cls, mock_get_embeddings, | |
| index_upsert=MagicMock(side_effect=ConnectionError("Pinecone write failure")), | |
| ) | |
| docs = [Document(page_content="Test content", metadata={"source": "test"})] | |
| with pytest.raises(ConnectionError, match="Pinecone write failure"): | |
| await service.add_documents(docs, TeamRole.PRODUCT_OWNER) | |
| async def test_delete_retries_and_succeeds_on_second_attempt( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """_delete_documents_sync should retry delete on transient failure, then succeed.""" | |
| call_count = 0 | |
| def delete_side_effect(*args, **kwargs): | |
| nonlocal call_count | |
| call_count += 1 | |
| if call_count == 1: | |
| raise ConnectionError("Transient delete failure") | |
| return True | |
| service, mock_index = _make_pinecone_service( | |
| mock_pinecone_cls, mock_get_embeddings, | |
| index_delete=delete_side_effect, | |
| ) | |
| result = await service.delete_documents(["doc-1"], TeamRole.PRODUCT_OWNER) | |
| assert result is True | |
| assert call_count >= 2 | |
| class TestPineconeCircuitBreaker: | |
| """DS-06: SyncCircuitBreaker protects pinecone operations.""" | |
| def test_circuit_breaker_opens_after_threshold_failures(self) -> None: | |
| """After failure_threshold failures, the circuit opens.""" | |
| from app.core.resilience import CircuitState, SyncCircuitBreaker | |
| cb = SyncCircuitBreaker("test", failure_threshold=3, recovery_timeout=30.0) | |
| assert cb.state is CircuitState.CLOSED | |
| for _ in range(3): | |
| assert cb.can_execute() is True | |
| cb.record_failure(Exception("fail")) | |
| assert cb.state is CircuitState.OPEN | |
| assert cb.can_execute() is False | |
| def test_circuit_breaker_recovers_after_timeout(self) -> None: | |
| """After recovery timeout, the circuit transitions to half-open.""" | |
| from app.core.resilience import CircuitState, SyncCircuitBreaker | |
| cb = SyncCircuitBreaker("test", failure_threshold=2, recovery_timeout=0.05) | |
| assert cb.can_execute() is True | |
| cb.record_failure(Exception("fail")) | |
| assert cb.can_execute() is True | |
| cb.record_failure(Exception("fail")) | |
| assert cb.state is CircuitState.OPEN | |
| assert cb.can_execute() is False # Still open | |
| time.sleep(0.06) # Past recovery_timeout | |
| assert cb.can_execute() is True # Now half-open | |
| assert cb.state is CircuitState.HALF_OPEN | |
| def test_circuit_breaker_closes_on_success_in_half_open(self) -> None: | |
| """A success in half-open closes the circuit.""" | |
| from app.core.resilience import CircuitState, SyncCircuitBreaker | |
| cb = SyncCircuitBreaker("test", failure_threshold=2, recovery_timeout=0.05) | |
| # Trip the circuit | |
| for _ in range(2): | |
| cb.record_failure(Exception("fail")) | |
| assert cb.state is CircuitState.OPEN | |
| time.sleep(0.06) | |
| # Half-open: success should close it | |
| assert cb.can_execute() is True | |
| cb.record_success() | |
| assert cb.state is CircuitState.CLOSED | |
| async def test_upsert_fails_fast_when_circuit_open( | |
| self, | |
| mock_pinecone_cls, | |
| mock_get_embeddings, | |
| ) -> None: | |
| """When circuit is open, add_documents should fail fast without calling Pinecone.""" | |
| from app.core.resilience import CircuitOpenError | |
| service, mock_index = _make_pinecone_service( | |
| mock_pinecone_cls, mock_get_embeddings, | |
| index_upsert=MagicMock(side_effect=ConnectionError("fail")), | |
| ) | |
| # Open the circuit by exhausting retries | |
| docs = [Document(page_content="Test", metadata={"source": "test"})] | |
| # The first call will fail and record failures through retry+circuit | |
| with pytest.raises(ConnectionError): | |
| await service.add_documents(docs, TeamRole.PRODUCT_OWNER) | |
| # Manually open the circuit so the next call fails fast | |
| if hasattr(service, '_pinecone_cb'): | |
| service._pinecone_cb.state = "open" | |
| service._pinecone_cb.last_failure_time = time.time() | |
| with pytest.raises(CircuitOpenError): | |
| await service.add_documents(docs, TeamRole.PRODUCT_OWNER) | |