"""Unit tests for repository connector.""" from unittest.mock import MagicMock import pytest from src.connectors.repository import Repository from src.enums import AbacusOutboxStatus from src.schemas import AbacusOutbox @pytest.fixture def mock_connection(): """Mock DB connection.""" return MagicMock() @pytest.fixture def repository(mock_connection): """Repository instance.""" return Repository(mock_connection) class TestRepository: """Test Repository.""" def test_get_pending_events(self, repository, mock_connection): """Test fetching and parsing pending events.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor # Mock DB row return mock_cursor.fetchall.return_value = [ { 'abacus_outbox_id': 1, 'target_type': 'test', 'target_id': 100, 'event_type': 'test.event', 'correlation_id': 'abc', 'details': '{}', 'status': 'pending', 'retry_count': 0, 'max_retries': 3, 'next_retry_at': None, 'error_message': None, 'created_at': '2023-01-01 00:00:00', 'last_modified': '2023-01-01 00:00:00', 'processed_at': None, } ] events = repository.get_pending_events(limit=5) assert len(events) == 1 assert isinstance(events[0], AbacusOutbox) assert events[0].abacus_outbox_id == 1 assert events[0].status == AbacusOutboxStatus.PENDING # Verify query arguments mock_cursor.execute.assert_called_once() args = mock_cursor.execute.call_args[0] assert args[1] == (AbacusOutboxStatus.PENDING, AbacusOutboxStatus.FAILED, 5) def test_update_event_status(self, repository, mock_connection): """Test updating event status.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.execute.return_value = 1 rows = repository.update_event_status( event_id=1, status=AbacusOutboxStatus.COMPLETED, processed_at='2023-01-01 10:00:00', ) assert rows == 1 mock_cursor.execute.assert_called_once() query, params = mock_cursor.execute.call_args[0] assert 'UPDATE abacus_outbox SET status = %s' in query assert 'processed_at = %s' in query assert 'WHERE abacus_outbox_id = %s' in query assert 'AND status IN (%s, %s)' in query # Optimistic locking assert params == [ AbacusOutboxStatus.COMPLETED, '2023-01-01 10:00:00', 1, AbacusOutboxStatus.PENDING, AbacusOutboxStatus.FAILED, ] def test_update_event_status_with_error(self, repository, mock_connection): """Test updating event status with error message.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.execute.return_value = 1 rows = repository.update_event_status( event_id=1, status=AbacusOutboxStatus.FAILED, error_message='EventBridge throttled', next_retry_at='2023-01-01 10:01:00', ) assert rows == 1 query, params = mock_cursor.execute.call_args[0] assert 'UPDATE abacus_outbox SET status = %s' in query assert 'error_message = %s' in query assert 'retry_count = retry_count + 1' in query assert 'next_retry_at = %s' in query assert 'WHERE abacus_outbox_id = %s' in query assert 'AND status IN (%s, %s)' in query # Optimistic locking assert params == [ AbacusOutboxStatus.FAILED, 'EventBridge throttled', '2023-01-01 10:01:00', 1, AbacusOutboxStatus.PENDING, AbacusOutboxStatus.FAILED, ] def test_update_event_status_all_fields(self, repository, mock_connection): """Test updating event status with all optional fields.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.execute.return_value = 1 rows = repository.update_event_status( event_id=123, status=AbacusOutboxStatus.FAILED, error_message='Timeout', next_retry_at='2023-01-01 10:05:00', processed_at='2023-01-01 10:00:00', ) assert rows == 1 query, params = mock_cursor.execute.call_args[0] assert 'UPDATE abacus_outbox SET status = %s' in query assert 'error_message = %s' in query assert 'retry_count = retry_count + 1' in query assert 'next_retry_at = %s' in query assert 'processed_at = %s' in query assert 'WHERE abacus_outbox_id = %s' in query def test_commit(self, repository, mock_connection): """Test commit method.""" repository.commit() mock_connection.commit.assert_called_once() def test_rollback(self, repository, mock_connection): """Test rollback method.""" repository.rollback() mock_connection.rollback.assert_called_once() def test_get_pending_events_empty(self, repository, mock_connection): """Test fetching pending events when none exist.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.fetchall.return_value = [] events = repository.get_pending_events(limit=10) assert events == [] mock_cursor.execute.assert_called_once() def test_get_pending_events_multiple(self, repository, mock_connection): """Test fetching multiple pending events.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.fetchall.return_value = [ { 'abacus_outbox_id': 1, 'target_type': 'test1', 'target_id': 100, 'event_type': 'test.event1', 'correlation_id': 'abc', 'details': '{}', 'status': 'pending', 'retry_count': 0, 'max_retries': 3, 'created_at': '2023-01-01 00:00:00', 'processed_at': None, }, { 'abacus_outbox_id': 2, 'target_type': 'test2', 'target_id': 200, 'event_type': 'test.event2', 'correlation_id': 'def', 'details': None, 'status': 'failed', 'retry_count': 1, 'max_retries': 5, 'created_at': '2023-01-01 01:00:00', 'processed_at': None, }, ] events = repository.get_pending_events(limit=10) assert len(events) == 2 assert events[0].abacus_outbox_id == 1 assert events[0].status == AbacusOutboxStatus.PENDING assert events[1].abacus_outbox_id == 2 assert events[1].status == AbacusOutboxStatus.FAILED assert events[1].retry_count == 1 def test_update_event_status_optimistic_locking(self, repository, mock_connection): """Test that update includes optimistic locking WHERE clause.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor mock_cursor.execute.return_value = 1 repository.update_event_status( event_id=1, status=AbacusOutboxStatus.COMPLETED, ) query, params = mock_cursor.execute.call_args[0] # Verify optimistic locking: only update PENDING or FAILED events assert 'WHERE abacus_outbox_id = %s' in query assert 'AND status IN (%s, %s)' in query assert params == [ AbacusOutboxStatus.COMPLETED, 1, AbacusOutboxStatus.PENDING, AbacusOutboxStatus.FAILED, ] def test_update_event_status_already_completed(self, repository, mock_connection): """Test that update returns 0 rows when event already completed.""" mock_cursor = MagicMock() mock_connection.cursor.return_value.__enter__.return_value = mock_cursor # Simulate no rows affected (event already completed by another consumer) mock_cursor.execute.return_value = 0 rows = repository.update_event_status( event_id=1, status=AbacusOutboxStatus.COMPLETED, ) assert rows == 0 mock_cursor.execute.assert_called_once()