From 22e6f86fe013beb457f65294608fdaede4eddb91 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Sat, 12 Jul 2025 10:59:06 +0100 Subject: [PATCH] Cassandra --- .../test_triples_cassandra_query.py | 70 ++++++ .../test_triples_cassandra_storage.py | 212 ++++++++++++++++++ 2 files changed, 282 insertions(+) create mode 100644 tests/unit/test_query/test_triples_cassandra_query.py create mode 100644 tests/unit/test_storage/test_triples_cassandra_storage.py diff --git a/tests/unit/test_query/test_triples_cassandra_query.py b/tests/unit/test_query/test_triples_cassandra_query.py new file mode 100644 index 00000000..eb81c8c8 --- /dev/null +++ b/tests/unit/test_query/test_triples_cassandra_query.py @@ -0,0 +1,70 @@ +""" +Tests for Cassandra triples query service +""" + +import pytest +from unittest.mock import MagicMock + +from trustgraph.query.triples.cassandra.service import Processor +from trustgraph.schema import Value + + +class TestCassandraQueryProcessor: + """Test cases for Cassandra query processor""" + + @pytest.fixture + def processor(self): + """Create a processor instance for testing""" + return Processor( + taskgroup=MagicMock(), + id='test-cassandra-query', + graph_host='localhost' + ) + + def test_create_value_with_http_uri(self, processor): + """Test create_value with HTTP URI""" + result = processor.create_value("http://example.com/resource") + + assert isinstance(result, Value) + assert result.value == "http://example.com/resource" + assert result.is_uri is True + + def test_create_value_with_https_uri(self, processor): + """Test create_value with HTTPS URI""" + result = processor.create_value("https://example.com/resource") + + assert isinstance(result, Value) + assert result.value == "https://example.com/resource" + assert result.is_uri is True + + def test_create_value_with_literal(self, processor): + """Test create_value with literal value""" + result = processor.create_value("just a literal string") + + assert isinstance(result, Value) + assert result.value == "just a literal string" + assert result.is_uri is False + + def test_create_value_with_empty_string(self, processor): + """Test create_value with empty string""" + result = processor.create_value("") + + assert isinstance(result, Value) + assert result.value == "" + assert result.is_uri is False + + def test_create_value_with_partial_uri(self, processor): + """Test create_value with string that looks like URI but isn't complete""" + result = processor.create_value("http") + + assert isinstance(result, Value) + assert result.value == "http" + assert result.is_uri is False + + def test_create_value_with_ftp_uri(self, processor): + """Test create_value with FTP URI (should not be detected as URI)""" + result = processor.create_value("ftp://example.com/file") + + assert isinstance(result, Value) + assert result.value == "ftp://example.com/file" + assert result.is_uri is False \ No newline at end of file diff --git a/tests/unit/test_storage/test_triples_cassandra_storage.py b/tests/unit/test_storage/test_triples_cassandra_storage.py new file mode 100644 index 00000000..1c675a34 --- /dev/null +++ b/tests/unit/test_storage/test_triples_cassandra_storage.py @@ -0,0 +1,212 @@ +""" +Tests for Cassandra triples storage service +""" + +import pytest +from unittest.mock import MagicMock, patch, AsyncMock + +from trustgraph.storage.triples.cassandra.write import Processor +from trustgraph.schema import Value, Triple + + +class TestCassandraStorageProcessor: + """Test cases for Cassandra storage processor""" + + def test_processor_initialization_with_defaults(self): + """Test processor initialization with default parameters""" + taskgroup_mock = MagicMock() + + processor = Processor(taskgroup=taskgroup_mock) + + assert processor.graph_host == ['localhost'] + assert processor.username is None + assert processor.password is None + assert processor.table is None + + def test_processor_initialization_with_custom_params(self): + """Test processor initialization with custom parameters""" + taskgroup_mock = MagicMock() + + processor = Processor( + taskgroup=taskgroup_mock, + id='custom-storage', + graph_host='cassandra.example.com', + graph_username='testuser', + graph_password='testpass' + ) + + assert processor.graph_host == ['cassandra.example.com'] + assert processor.username == 'testuser' + assert processor.password == 'testpass' + assert processor.table is None + + def test_processor_initialization_with_partial_auth(self): + """Test processor initialization with only username (no password)""" + taskgroup_mock = MagicMock() + + processor = Processor( + taskgroup=taskgroup_mock, + graph_username='testuser' + ) + + assert processor.username == 'testuser' + assert processor.password is None + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + async def test_table_switching_with_auth(self, mock_trustgraph): + """Test table switching logic when authentication is provided""" + taskgroup_mock = MagicMock() + mock_tg_instance = MagicMock() + mock_trustgraph.return_value = mock_tg_instance + + processor = Processor( + taskgroup=taskgroup_mock, + graph_username='testuser', + graph_password='testpass' + ) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'user1' + mock_message.metadata.collection = 'collection1' + mock_message.triples = [] + + await processor.store_triples(mock_message) + + # Verify TrustGraph was called with auth parameters + mock_trustgraph.assert_called_once_with( + hosts=['localhost'], + keyspace='user1', + table='collection1', + username='testuser', + password='testpass' + ) + assert processor.table == ('user1', 'collection1') + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + async def test_table_switching_without_auth(self, mock_trustgraph): + """Test table switching logic when no authentication is provided""" + taskgroup_mock = MagicMock() + mock_tg_instance = MagicMock() + mock_trustgraph.return_value = mock_tg_instance + + processor = Processor(taskgroup=taskgroup_mock) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'user2' + mock_message.metadata.collection = 'collection2' + mock_message.triples = [] + + await processor.store_triples(mock_message) + + # Verify TrustGraph was called without auth parameters + mock_trustgraph.assert_called_once_with( + hosts=['localhost'], + keyspace='user2', + table='collection2' + ) + assert processor.table == ('user2', 'collection2') + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + async def test_table_reuse_when_same(self, mock_trustgraph): + """Test that TrustGraph is not recreated when table hasn't changed""" + taskgroup_mock = MagicMock() + mock_tg_instance = MagicMock() + mock_trustgraph.return_value = mock_tg_instance + + processor = Processor(taskgroup=taskgroup_mock) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'user1' + mock_message.metadata.collection = 'collection1' + mock_message.triples = [] + + # First call should create TrustGraph + await processor.store_triples(mock_message) + assert mock_trustgraph.call_count == 1 + + # Second call with same table should reuse TrustGraph + await processor.store_triples(mock_message) + assert mock_trustgraph.call_count == 1 # Should not increase + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + async def test_triple_insertion(self, mock_trustgraph): + """Test that triples are properly inserted into Cassandra""" + taskgroup_mock = MagicMock() + mock_tg_instance = MagicMock() + mock_trustgraph.return_value = mock_tg_instance + + processor = Processor(taskgroup=taskgroup_mock) + + # Create mock triples + triple1 = MagicMock() + triple1.s.value = 'subject1' + triple1.p.value = 'predicate1' + triple1.o.value = 'object1' + + triple2 = MagicMock() + triple2.s.value = 'subject2' + triple2.p.value = 'predicate2' + triple2.o.value = 'object2' + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'user1' + mock_message.metadata.collection = 'collection1' + mock_message.triples = [triple1, triple2] + + await processor.store_triples(mock_message) + + # Verify both triples were inserted + assert mock_tg_instance.insert.call_count == 2 + mock_tg_instance.insert.assert_any_call('subject1', 'predicate1', 'object1') + mock_tg_instance.insert.assert_any_call('subject2', 'predicate2', 'object2') + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + async def test_triple_insertion_with_empty_list(self, mock_trustgraph): + """Test behavior when message has no triples""" + taskgroup_mock = MagicMock() + mock_tg_instance = MagicMock() + mock_trustgraph.return_value = mock_tg_instance + + processor = Processor(taskgroup=taskgroup_mock) + + # Create mock message with empty triples + mock_message = MagicMock() + mock_message.metadata.user = 'user1' + mock_message.metadata.collection = 'collection1' + mock_message.triples = [] + + await processor.store_triples(mock_message) + + # Verify no triples were inserted + mock_tg_instance.insert.assert_not_called() + + @pytest.mark.asyncio + @patch('trustgraph.storage.triples.cassandra.write.TrustGraph') + @patch('trustgraph.storage.triples.cassandra.write.time.sleep') + async def test_exception_handling_with_retry(self, mock_sleep, mock_trustgraph): + """Test exception handling during TrustGraph creation""" + taskgroup_mock = MagicMock() + mock_trustgraph.side_effect = Exception("Connection failed") + + processor = Processor(taskgroup=taskgroup_mock) + + # Create mock message + mock_message = MagicMock() + mock_message.metadata.user = 'user1' + mock_message.metadata.collection = 'collection1' + mock_message.triples = [] + + with pytest.raises(Exception, match="Connection failed"): + await processor.store_triples(mock_message) + + # Verify sleep was called before re-raising + mock_sleep.assert_called_once_with(1) \ No newline at end of file