mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-23 20:21:03 +02:00
Fix integratino tests
This commit is contained in:
parent
fe49a606cf
commit
8a7be0a3f1
2 changed files with 40 additions and 52 deletions
|
|
@ -6,23 +6,11 @@ import uuid
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
from trustgraph.base.subscriber import Subscriber
|
from trustgraph.base.subscriber import Subscriber
|
||||||
|
|
||||||
# Mock JsonSchema globally to avoid schema issues in tests
|
|
||||||
# Patch at the module level where it's imported in subscriber
|
|
||||||
@patch('trustgraph.base.subscriber.JsonSchema')
|
|
||||||
def mock_json_schema_global(mock_schema):
|
|
||||||
mock_schema.return_value = MagicMock()
|
|
||||||
return mock_schema
|
|
||||||
|
|
||||||
# Apply the global patch
|
|
||||||
_json_schema_patch = patch('trustgraph.base.subscriber.JsonSchema')
|
|
||||||
_mock_json_schema = _json_schema_patch.start()
|
|
||||||
_mock_json_schema.return_value = MagicMock()
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def mock_pulsar_client():
|
def mock_pulsar_backend():
|
||||||
"""Mock Pulsar client for testing."""
|
"""Mock Pulsar backend for testing."""
|
||||||
client = MagicMock()
|
backend = MagicMock()
|
||||||
consumer = MagicMock()
|
consumer = MagicMock()
|
||||||
consumer.receive = MagicMock()
|
consumer.receive = MagicMock()
|
||||||
consumer.acknowledge = MagicMock()
|
consumer.acknowledge = MagicMock()
|
||||||
|
|
@ -30,15 +18,15 @@ def mock_pulsar_client():
|
||||||
consumer.pause_message_listener = MagicMock()
|
consumer.pause_message_listener = MagicMock()
|
||||||
consumer.unsubscribe = MagicMock()
|
consumer.unsubscribe = MagicMock()
|
||||||
consumer.close = MagicMock()
|
consumer.close = MagicMock()
|
||||||
client.subscribe.return_value = consumer
|
backend.create_consumer.return_value = consumer
|
||||||
return client
|
return backend
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
def subscriber(mock_pulsar_client):
|
def subscriber(mock_pulsar_backend):
|
||||||
"""Create Subscriber instance for testing."""
|
"""Create Subscriber instance for testing."""
|
||||||
return Subscriber(
|
return Subscriber(
|
||||||
client=mock_pulsar_client,
|
backend=mock_pulsar_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
|
|
@ -60,14 +48,14 @@ def create_mock_message(message_id="test-id", data=None):
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_deferred_acknowledgment_success():
|
async def test_subscriber_deferred_acknowledgment_success():
|
||||||
"""Verify Subscriber only acks on successful delivery."""
|
"""Verify Subscriber only acks on successful delivery."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
schema=dict,
|
schema=dict,
|
||||||
max_size=10,
|
max_size=10,
|
||||||
|
|
@ -102,15 +90,15 @@ async def test_subscriber_deferred_acknowledgment_success():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_deferred_acknowledgment_failure():
|
async def test_subscriber_deferred_acknowledgment_failure():
|
||||||
"""Verify Subscriber negative acks on delivery failure."""
|
"""Verify Subscriber negative acks on delivery failure."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
schema=dict,
|
schema=dict,
|
||||||
max_size=1, # Very small queue
|
max_size=1, # Very small queue
|
||||||
backpressure_strategy="drop_new"
|
backpressure_strategy="drop_new"
|
||||||
|
|
@ -140,14 +128,14 @@ async def test_subscriber_deferred_acknowledgment_failure():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_backpressure_strategies():
|
async def test_subscriber_backpressure_strategies():
|
||||||
"""Test different backpressure strategies."""
|
"""Test different backpressure strategies."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
# Test drop_oldest strategy
|
# Test drop_oldest strategy
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
schema=dict,
|
schema=dict,
|
||||||
|
|
@ -187,12 +175,12 @@ async def test_subscriber_backpressure_strategies():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_graceful_shutdown():
|
async def test_subscriber_graceful_shutdown():
|
||||||
"""Test Subscriber graceful shutdown with queue draining."""
|
"""Test Subscriber graceful shutdown with queue draining."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
|
|
@ -253,14 +241,14 @@ async def test_subscriber_graceful_shutdown():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_drain_timeout():
|
async def test_subscriber_drain_timeout():
|
||||||
"""Test Subscriber respects drain timeout."""
|
"""Test Subscriber respects drain timeout."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
schema=dict,
|
schema=dict,
|
||||||
max_size=10,
|
max_size=10,
|
||||||
|
|
@ -288,12 +276,12 @@ async def test_subscriber_drain_timeout():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_pending_acks_cleanup():
|
async def test_subscriber_pending_acks_cleanup():
|
||||||
"""Test Subscriber cleans up pending acknowledgments on shutdown."""
|
"""Test Subscriber cleans up pending acknowledgments on shutdown."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
|
|
@ -342,12 +330,12 @@ async def test_subscriber_pending_acks_cleanup():
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_subscriber_multiple_subscribers():
|
async def test_subscriber_multiple_subscribers():
|
||||||
"""Test Subscriber with multiple concurrent subscribers."""
|
"""Test Subscriber with multiple concurrent subscribers."""
|
||||||
mock_client = MagicMock()
|
mock_backend = MagicMock()
|
||||||
mock_consumer = MagicMock()
|
mock_consumer = MagicMock()
|
||||||
mock_client.subscribe.return_value = mock_consumer
|
mock_backend.create_consumer.return_value = mock_consumer
|
||||||
|
|
||||||
subscriber = Subscriber(
|
subscriber = Subscriber(
|
||||||
client=mock_client,
|
backend=mock_backend,
|
||||||
topic="test-topic",
|
topic="test-topic",
|
||||||
subscription="test-subscription",
|
subscription="test-subscription",
|
||||||
consumer_name="test-consumer",
|
consumer_name="test-consumer",
|
||||||
|
|
|
||||||
|
|
@ -1,5 +1,5 @@
|
||||||
|
|
||||||
from . pubsub import PulsarClient
|
from . pubsub import PulsarClient, get_pubsub
|
||||||
from . async_processor import AsyncProcessor
|
from . async_processor import AsyncProcessor
|
||||||
from . consumer import Consumer
|
from . consumer import Consumer
|
||||||
from . producer import Producer
|
from . producer import Producer
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue