diff --git a/tests/unit/test_storage/conftest.py b/tests/unit/test_storage/conftest.py new file mode 100644 index 00000000..594e2b2f --- /dev/null +++ b/tests/unit/test_storage/conftest.py @@ -0,0 +1,162 @@ +""" +Shared fixtures for storage tests +""" + +import pytest +from unittest.mock import AsyncMock, MagicMock + + +@pytest.fixture +def base_storage_config(): + """Base configuration for storage processors""" + return { + 'taskgroup': AsyncMock(), + 'id': 'test-storage-processor' + } + + +@pytest.fixture +def qdrant_storage_config(base_storage_config): + """Configuration for Qdrant storage processors""" + return base_storage_config | { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key' + } + + +@pytest.fixture +def mock_qdrant_client(): + """Mock Qdrant client""" + mock_client = MagicMock() + mock_client.collection_exists.return_value = True + mock_client.create_collection.return_value = None + mock_client.upsert.return_value = None + return mock_client + + +@pytest.fixture +def mock_uuid(): + """Mock UUID generation""" + mock_uuid = MagicMock() + mock_uuid.uuid4.return_value = MagicMock() + mock_uuid.uuid4.return_value.__str__ = MagicMock(return_value='test-uuid-123') + return mock_uuid + + +# Document embeddings fixtures +@pytest.fixture +def mock_document_embeddings_message(): + """Mock document embeddings message""" + mock_message = MagicMock() + mock_message.metadata.user = 'test_user' + mock_message.metadata.collection = 'test_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'test document chunk' + mock_chunk.vectors = [[0.1, 0.2, 0.3]] + + mock_message.chunks = [mock_chunk] + return mock_message + + +@pytest.fixture +def mock_document_embeddings_multiple_chunks(): + """Mock document embeddings message with multiple chunks""" + mock_message = MagicMock() + mock_message.metadata.user = 'multi_user' + mock_message.metadata.collection = 'multi_collection' + + mock_chunk1 = MagicMock() + mock_chunk1.chunk.decode.return_value = 'first document chunk' + mock_chunk1.vectors = [[0.1, 0.2]] + + mock_chunk2 = MagicMock() + mock_chunk2.chunk.decode.return_value = 'second document chunk' + mock_chunk2.vectors = [[0.3, 0.4]] + + mock_message.chunks = [mock_chunk1, mock_chunk2] + return mock_message + + +@pytest.fixture +def mock_document_embeddings_multiple_vectors(): + """Mock document embeddings message with multiple vectors per chunk""" + mock_message = MagicMock() + mock_message.metadata.user = 'vector_user' + mock_message.metadata.collection = 'vector_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'multi-vector document chunk' + mock_chunk.vectors = [ + [0.1, 0.2, 0.3], + [0.4, 0.5, 0.6], + [0.7, 0.8, 0.9] + ] + + mock_message.chunks = [mock_chunk] + return mock_message + + +@pytest.fixture +def mock_document_embeddings_empty_chunk(): + """Mock document embeddings message with empty chunk""" + mock_message = MagicMock() + mock_message.metadata.user = 'empty_user' + mock_message.metadata.collection = 'empty_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = "" # Empty string + mock_chunk.vectors = [[0.1, 0.2]] + + mock_message.chunks = [mock_chunk] + return mock_message + + +# Graph embeddings fixtures +@pytest.fixture +def mock_graph_embeddings_message(): + """Mock graph embeddings message""" + mock_message = MagicMock() + mock_message.metadata.user = 'test_user' + mock_message.metadata.collection = 'test_collection' + + mock_entity = MagicMock() + mock_entity.entity.value = 'test_entity' + mock_entity.vectors = [[0.1, 0.2, 0.3]] + + mock_message.entities = [mock_entity] + return mock_message + + +@pytest.fixture +def mock_graph_embeddings_multiple_entities(): + """Mock graph embeddings message with multiple entities""" + mock_message = MagicMock() + mock_message.metadata.user = 'multi_user' + mock_message.metadata.collection = 'multi_collection' + + mock_entity1 = MagicMock() + mock_entity1.entity.value = 'entity_one' + mock_entity1.vectors = [[0.1, 0.2]] + + mock_entity2 = MagicMock() + mock_entity2.entity.value = 'entity_two' + mock_entity2.vectors = [[0.3, 0.4]] + + mock_message.entities = [mock_entity1, mock_entity2] + return mock_message + + +@pytest.fixture +def mock_graph_embeddings_empty_entity(): + """Mock graph embeddings message with empty entity""" + mock_message = MagicMock() + mock_message.metadata.user = 'empty_user' + mock_message.metadata.collection = 'empty_collection' + + mock_entity = MagicMock() + mock_entity.entity.value = "" # Empty string + mock_entity.vectors = [[0.1, 0.2]] + + mock_message.entities = [mock_entity] + return mock_message \ No newline at end of file diff --git a/tests/unit/test_storage/test_doc_embeddings_qdrant_processor.py b/tests/unit/test_storage/test_doc_embeddings_qdrant_processor.py new file mode 100644 index 00000000..016a356e --- /dev/null +++ b/tests/unit/test_storage/test_doc_embeddings_qdrant_processor.py @@ -0,0 +1,569 @@ +""" +Unit tests for trustgraph.storage.doc_embeddings.qdrant.write +Testing document embeddings storage functionality +""" + +import pytest +from unittest.mock import AsyncMock, MagicMock, patch +from unittest import IsolatedAsyncioTestCase + +# Import the service under test +from trustgraph.storage.doc_embeddings.qdrant.write import Processor + + +class TestQdrantDocEmbeddingsProcessor(IsolatedAsyncioTestCase): + """Test Qdrant document embeddings storage functionality""" + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_processor_initialization_basic(self, mock_base_init, mock_qdrant_client): + """Test basic Qdrant processor initialization""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + # Act + processor = Processor(**config) + + # Assert + # Verify base class initialization was called + mock_base_init.assert_called_once() + + # Verify QdrantClient was created with correct parameters + mock_qdrant_client.assert_called_once_with(url='http://localhost:6333', api_key='test-api-key') + + # Verify processor attributes + assert hasattr(processor, 'qdrant') + assert processor.qdrant == mock_qdrant_instance + assert hasattr(processor, 'last_collection') + assert processor.last_collection is None + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_processor_initialization_with_defaults(self, mock_base_init, mock_qdrant_client): + """Test processor initialization with default values""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + # No store_uri or api_key provided - should use defaults + } + + # Act + processor = Processor(**config) + + # Assert + # Verify QdrantClient was created with default URI and None API key + mock_qdrant_client.assert_called_once_with(url='http://localhost:6333', api_key=None) + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.storage.doc_embeddings.qdrant.write.uuid') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_store_document_embeddings_basic(self, mock_base_init, mock_uuid, mock_qdrant_client): + """Test storing document embeddings with basic message""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True # Collection already exists + mock_qdrant_client.return_value = mock_qdrant_instance + mock_uuid.uuid4.return_value = MagicMock() + mock_uuid.uuid4.return_value.__str__ = MagicMock(return_value='test-uuid-123') + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with chunks and vectors + mock_message = MagicMock() + mock_message.metadata.user = 'test_user' + mock_message.metadata.collection = 'test_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'test document chunk' + mock_chunk.vectors = [[0.1, 0.2, 0.3]] # Single vector with 3 dimensions + + mock_message.chunks = [mock_chunk] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Verify collection existence was checked + expected_collection = 'd_test_user_test_collection_3' + mock_qdrant_instance.collection_exists.assert_called_once_with(expected_collection) + + # Verify upsert was called + mock_qdrant_instance.upsert.assert_called_once() + + # Verify upsert parameters + upsert_call_args = mock_qdrant_instance.upsert.call_args + assert upsert_call_args[1]['collection_name'] == expected_collection + assert len(upsert_call_args[1]['points']) == 1 + + point = upsert_call_args[1]['points'][0] + assert point.vector == [0.1, 0.2, 0.3] + assert point.payload['doc'] == 'test document chunk' + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.storage.doc_embeddings.qdrant.write.uuid') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_store_document_embeddings_multiple_chunks(self, mock_base_init, mock_uuid, mock_qdrant_client): + """Test storing document embeddings with multiple chunks""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True + mock_qdrant_client.return_value = mock_qdrant_instance + mock_uuid.uuid4.return_value = MagicMock() + mock_uuid.uuid4.return_value.__str__ = MagicMock(return_value='test-uuid') + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with multiple chunks + mock_message = MagicMock() + mock_message.metadata.user = 'multi_user' + mock_message.metadata.collection = 'multi_collection' + + mock_chunk1 = MagicMock() + mock_chunk1.chunk.decode.return_value = 'first document chunk' + mock_chunk1.vectors = [[0.1, 0.2]] + + mock_chunk2 = MagicMock() + mock_chunk2.chunk.decode.return_value = 'second document chunk' + mock_chunk2.vectors = [[0.3, 0.4]] + + mock_message.chunks = [mock_chunk1, mock_chunk2] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Should be called twice (once per chunk) + assert mock_qdrant_instance.upsert.call_count == 2 + + # Verify both chunks were processed + upsert_calls = mock_qdrant_instance.upsert.call_args_list + + # First chunk + first_call = upsert_calls[0] + first_point = first_call[1]['points'][0] + assert first_point.vector == [0.1, 0.2] + assert first_point.payload['doc'] == 'first document chunk' + + # Second chunk + second_call = upsert_calls[1] + second_point = second_call[1]['points'][0] + assert second_point.vector == [0.3, 0.4] + assert second_point.payload['doc'] == 'second document chunk' + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.storage.doc_embeddings.qdrant.write.uuid') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_store_document_embeddings_multiple_vectors_per_chunk(self, mock_base_init, mock_uuid, mock_qdrant_client): + """Test storing document embeddings with multiple vectors per chunk""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True + mock_qdrant_client.return_value = mock_qdrant_instance + mock_uuid.uuid4.return_value = MagicMock() + mock_uuid.uuid4.return_value.__str__ = MagicMock(return_value='test-uuid') + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with chunk having multiple vectors + mock_message = MagicMock() + mock_message.metadata.user = 'vector_user' + mock_message.metadata.collection = 'vector_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'multi-vector document chunk' + mock_chunk.vectors = [ + [0.1, 0.2, 0.3], + [0.4, 0.5, 0.6], + [0.7, 0.8, 0.9] + ] + + mock_message.chunks = [mock_chunk] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Should be called 3 times (once per vector) + assert mock_qdrant_instance.upsert.call_count == 3 + + # Verify all vectors were processed + upsert_calls = mock_qdrant_instance.upsert.call_args_list + + expected_vectors = [ + [0.1, 0.2, 0.3], + [0.4, 0.5, 0.6], + [0.7, 0.8, 0.9] + ] + + for i, call in enumerate(upsert_calls): + point = call[1]['points'][0] + assert point.vector == expected_vectors[i] + assert point.payload['doc'] == 'multi-vector document chunk' + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_store_document_embeddings_empty_chunk(self, mock_base_init, mock_qdrant_client): + """Test storing document embeddings skips empty chunks""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with empty chunk + mock_message = MagicMock() + mock_message.metadata.user = 'empty_user' + mock_message.metadata.collection = 'empty_collection' + + mock_chunk_empty = MagicMock() + mock_chunk_empty.chunk.decode.return_value = "" # Empty string + mock_chunk_empty.vectors = [[0.1, 0.2]] + + mock_message.chunks = [mock_chunk_empty] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Should not call upsert for empty chunks + mock_qdrant_instance.upsert.assert_not_called() + mock_qdrant_instance.collection_exists.assert_not_called() + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_collection_creation_when_not_exists(self, mock_base_init, mock_qdrant_client): + """Test collection creation when it doesn't exist""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = False # Collection doesn't exist + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'new_user' + mock_message.metadata.collection = 'new_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'test chunk' + mock_chunk.vectors = [[0.1, 0.2, 0.3, 0.4, 0.5]] # 5 dimensions + + mock_message.chunks = [mock_chunk] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + expected_collection = 'd_new_user_new_collection_5' + + # Verify collection existence check and creation + mock_qdrant_instance.collection_exists.assert_called_once_with(expected_collection) + mock_qdrant_instance.create_collection.assert_called_once() + + # Verify create_collection was called with correct parameters + create_call_args = mock_qdrant_instance.create_collection.call_args + assert create_call_args[1]['collection_name'] == expected_collection + + # Verify upsert was still called after collection creation + mock_qdrant_instance.upsert.assert_called_once() + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_collection_creation_exception(self, mock_base_init, mock_qdrant_client): + """Test collection creation handles exceptions""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = False + mock_qdrant_instance.create_collection.side_effect = Exception("Qdrant connection failed") + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'error_user' + mock_message.metadata.collection = 'error_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'test chunk' + mock_chunk.vectors = [[0.1, 0.2]] + + mock_message.chunks = [mock_chunk] + + # Act & Assert + with pytest.raises(Exception, match="Qdrant connection failed"): + await processor.store_document_embeddings(mock_message) + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_collection_caching_behavior(self, mock_base_init, mock_qdrant_client): + """Test collection caching with last_collection""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create first mock message + mock_message1 = MagicMock() + mock_message1.metadata.user = 'cache_user' + mock_message1.metadata.collection = 'cache_collection' + + mock_chunk1 = MagicMock() + mock_chunk1.chunk.decode.return_value = 'first chunk' + mock_chunk1.vectors = [[0.1, 0.2, 0.3]] + + mock_message1.chunks = [mock_chunk1] + + # First call + await processor.store_document_embeddings(mock_message1) + + # Reset mock to track second call + mock_qdrant_instance.reset_mock() + + # Create second mock message with same dimensions + mock_message2 = MagicMock() + mock_message2.metadata.user = 'cache_user' + mock_message2.metadata.collection = 'cache_collection' + + mock_chunk2 = MagicMock() + mock_chunk2.chunk.decode.return_value = 'second chunk' + mock_chunk2.vectors = [[0.4, 0.5, 0.6]] # Same dimension (3) + + mock_message2.chunks = [mock_chunk2] + + # Act - Second call with same collection + await processor.store_document_embeddings(mock_message2) + + # Assert + expected_collection = 'd_cache_user_cache_collection_3' + assert processor.last_collection == expected_collection + + # Verify second call skipped existence check (cached) + mock_qdrant_instance.collection_exists.assert_not_called() + mock_qdrant_instance.create_collection.assert_not_called() + + # But upsert should still be called + mock_qdrant_instance.upsert.assert_called_once() + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_different_dimensions_different_collections(self, mock_base_init, mock_qdrant_client): + """Test that different vector dimensions create different collections""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with different dimension vectors + mock_message = MagicMock() + mock_message.metadata.user = 'dim_user' + mock_message.metadata.collection = 'dim_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'dimension test chunk' + mock_chunk.vectors = [ + [0.1, 0.2], # 2 dimensions + [0.3, 0.4, 0.5] # 3 dimensions + ] + + mock_message.chunks = [mock_chunk] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Should check existence of both collections + expected_collections = ['d_dim_user_dim_collection_2', 'd_dim_user_dim_collection_3'] + actual_calls = [call.args[0] for call in mock_qdrant_instance.collection_exists.call_args_list] + assert actual_calls == expected_collections + + # Should upsert to both collections + assert mock_qdrant_instance.upsert.call_count == 2 + + upsert_calls = mock_qdrant_instance.upsert.call_args_list + assert upsert_calls[0][1]['collection_name'] == 'd_dim_user_dim_collection_2' + assert upsert_calls[1][1]['collection_name'] == 'd_dim_user_dim_collection_3' + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_add_args_calls_parent(self, mock_base_init, mock_qdrant_client): + """Test that add_args() calls parent add_args method""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_client.return_value = MagicMock() + mock_parser = MagicMock() + + # Act + with patch('trustgraph.base.DocumentEmbeddingsStoreService.add_args') as mock_parent_add_args: + Processor.add_args(mock_parser) + + # Assert + mock_parent_add_args.assert_called_once_with(mock_parser) + + # Verify processor-specific arguments were added + assert mock_parser.add_argument.call_count >= 2 # At least store-uri and api-key + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.storage.doc_embeddings.qdrant.write.uuid') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_utf8_decoding_handling(self, mock_base_init, mock_uuid, mock_qdrant_client): + """Test proper UTF-8 decoding of chunk text""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_instance.collection_exists.return_value = True + mock_qdrant_client.return_value = mock_qdrant_instance + mock_uuid.uuid4.return_value = MagicMock() + mock_uuid.uuid4.return_value.__str__ = MagicMock(return_value='test-uuid') + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with UTF-8 encoded text + mock_message = MagicMock() + mock_message.metadata.user = 'utf8_user' + mock_message.metadata.collection = 'utf8_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.return_value = 'UTF-8 text with special chars: café, naïve, résumé' + mock_chunk.vectors = [[0.1, 0.2]] + + mock_message.chunks = [mock_chunk] + + # Act + await processor.store_document_embeddings(mock_message) + + # Assert + # Verify chunk.decode was called with 'utf-8' + mock_chunk.chunk.decode.assert_called_with('utf-8') + + # Verify the decoded text was stored in payload + upsert_call_args = mock_qdrant_instance.upsert.call_args + point = upsert_call_args[1]['points'][0] + assert point.payload['doc'] == 'UTF-8 text with special chars: café, naïve, résumé' + + @patch('trustgraph.storage.doc_embeddings.qdrant.write.QdrantClient') + @patch('trustgraph.base.DocumentEmbeddingsStoreService.__init__') + async def test_chunk_decode_exception_handling(self, mock_base_init, mock_qdrant_client): + """Test handling of chunk decode exceptions""" + # Arrange + mock_base_init.return_value = None + mock_qdrant_instance = MagicMock() + mock_qdrant_client.return_value = mock_qdrant_instance + + config = { + 'store_uri': 'http://localhost:6333', + 'api_key': 'test-api-key', + 'taskgroup': AsyncMock(), + 'id': 'test-doc-qdrant-processor' + } + + processor = Processor(**config) + + # Create mock message with decode error + mock_message = MagicMock() + mock_message.metadata.user = 'decode_user' + mock_message.metadata.collection = 'decode_collection' + + mock_chunk = MagicMock() + mock_chunk.chunk.decode.side_effect = UnicodeDecodeError('utf-8', b'', 0, 1, 'invalid start byte') + mock_chunk.vectors = [[0.1, 0.2]] + + mock_message.chunks = [mock_chunk] + + # Act & Assert + with pytest.raises(UnicodeDecodeError): + await processor.store_document_embeddings(mock_message) + + +if __name__ == '__main__': + pytest.main([__file__]) \ No newline at end of file