"""Unit tests for processor module.""" from datetime import datetime, timedelta, timezone from unittest.mock import MagicMock import pytest from freezegun import freeze_time from src.enums import AbacusOutboxStatus from src.processor import OutboxProcessor from src.schemas import AbacusOutbox from src.utils import calc_retry_delta, try_json @pytest.fixture def mock_repository(): """Mock repository.""" return MagicMock() @pytest.fixture def mock_eb_connector(): """Mock EventBridge connector.""" return MagicMock() @pytest.fixture def processor(mock_repository, mock_eb_connector): """Outbox processor instance.""" return OutboxProcessor(mock_repository, mock_eb_connector) @pytest.fixture def sample_event(): """Sample AbacusOutbox.""" return AbacusOutbox( abacus_outbox_id=123, target_type='file_upload', target_id=456, event_type='file_upload.completed', correlation_id='uuid-123', details={'some': 'data'}, status=AbacusOutboxStatus.PENDING, retry_count=0, max_retries=3, created_at=datetime(2023, 1, 1, 10, 0, 0, tzinfo=timezone.utc), processed_at=None, ) class TestOutboxProcessor: """Test OutboxProcessor.""" def test_process_completed_event_skipped( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that already completed events are skipped (idempotency).""" # Set event status to COMPLETED sample_event.status = AbacusOutboxStatus.COMPLETED # Process the event result = processor.process([sample_event]) # Verify counts assert result.total == 1 assert result.skipped == 1 assert result.processed == 0 assert result.failed == 0 # Verify EventBridge was NOT called mock_eb_connector.put_event.assert_not_called() # Verify DB update was NOT called mock_repository.update_event_status.assert_not_called() mock_repository.commit.assert_not_called() def test_process_success( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test successful event processing.""" # Mock successful DB update mock_repository.update_event_status.return_value = 1 # Process the event result = processor.process([sample_event]) # Verify counts assert result.total == 1 assert result.processed == 1 assert result.failed == 0 assert result.skipped == 0 # Verify EventBridge call mock_eb_connector.put_event.assert_called_once() call_args = mock_eb_connector.put_event.call_args assert call_args.kwargs['detail_type'] == 'file_upload.completed' assert call_args.kwargs['metadata'].target_id == 456 assert call_args.kwargs['data'] == {'some': 'data'} # Verify DB update mock_repository.update_event_status.assert_called_once() assert ( mock_repository.update_event_status.call_args.kwargs['status'] == AbacusOutboxStatus.COMPLETED ) mock_repository.commit.assert_called_once() def test_process_missing_correlation_id( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test generating correlation ID if missing.""" sample_event.correlation_id = '' # Empty string to simulate missing mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) call_args = mock_eb_connector.put_event.call_args correlation_id = call_args.kwargs['metadata'].correlation_id assert correlation_id assert correlation_id != '' assert len(correlation_id) == 36 # UUID length def test_process_details_string_json( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test processing event with details as JSON string.""" sample_event.details = '{"foo": "bar"}' mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) # Verify details was parsed call_args = mock_eb_connector.put_event.call_args assert call_args.kwargs['data'] == {'foo': 'bar'} def test_process_details_invalid_json( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test processing event with invalid JSON details string.""" sample_event.details = '{invalid_json}' mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) # Verify details remains as string (unwrapped) call_args = mock_eb_connector.put_event.call_args assert call_args.kwargs['data'] == '{invalid_json}' def test_process_failure_and_retry( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test processing failure handles retry logic.""" mock_eb_connector.put_event.side_effect = Exception('EventBridge error') mock_repository.update_event_status.return_value = 1 with freeze_time('2023-01-01 12:00:00'): response = processor.process([sample_event]) assert response.total == 1 assert response.processed == 0 assert response.failed == 1 # Verify DB update for failure mock_repository.update_event_status.assert_called_once() kwargs = mock_repository.update_event_status.call_args.kwargs assert kwargs['status'] == AbacusOutboxStatus.FAILED assert kwargs['error_message'] == 'EventBridge error' # Verify next_retry_at calculation # Base: 30 * (1.618^0) = 30s. Jitter +/- 5% (1.5s). Range: 28.5s - 31.5s next_retry_at = kwargs['next_retry_at'] diff = ( next_retry_at - datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc) ).total_seconds() assert 28.5 <= diff <= 31.5 mock_repository.commit.assert_called_once() def test_process_failure_db_update_fails( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test failure when updating event status also fails.""" mock_eb_connector.put_event.side_effect = Exception('EventBridge error') mock_repository.update_event_status.side_effect = Exception('DB Error') # Should raise the primary exception (EventBridge error) with DB error as context result = processor.process([sample_event]) assert result.total == 1 assert result.failed == 1 assert result.processed == 0 # Verify EventBridge was called (before error update failure) mock_eb_connector.put_event.assert_called_once() # Verify error update was attempted mock_repository.update_event_status.assert_called_once() def test_batch_limit_respect(self, processor, mock_repository, sample_event): """Test that processor processes multiple events in iterator.""" # Create 5 events events = [sample_event for _ in range(5)] mock_repository.update_event_status.return_value = 1 response = processor.process(events) assert response.total == 5 assert response.processed == 5 assert response.failed == 0 def test_process_mixed_success_and_failure( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test batch with both successful and failed events.""" success_event = sample_event fail_event = AbacusOutbox( abacus_outbox_id=124, target_type='file_upload', target_id=457, event_type='file_upload.failed', correlation_id='uuid-124', details={'error': 'data'}, status=AbacusOutboxStatus.PENDING, retry_count=0, max_retries=3, created_at=datetime(2023, 1, 1, 11, 0, 0, tzinfo=timezone.utc), processed_at=None, ) # Make second put_event fail mock_eb_connector.put_event.side_effect = [None, Exception('Network error')] mock_repository.update_event_status.return_value = 1 response = processor.process([success_event, fail_event]) assert response.total == 2 assert response.processed == 1 assert response.failed == 1 assert mock_eb_connector.put_event.call_count == 2 assert mock_repository.commit.call_count == 2 def test_process_none_details( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test processing event with None details.""" sample_event.details = None mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) call_args = mock_eb_connector.put_event.call_args assert call_args.kwargs['data'] == {} def test_process_empty_string_details( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test processing event with empty string details.""" sample_event.details = '' mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) # Empty string is falsy, so code treats it as None and returns {} call_args = mock_eb_connector.put_event.call_args assert call_args.kwargs['data'] == {} def test_process_none_correlation_id( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test generating correlation ID when None.""" sample_event.correlation_id = None mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) call_args = mock_eb_connector.put_event.call_args correlation_id = call_args.kwargs['metadata'].correlation_id assert correlation_id is not None assert len(correlation_id) == 36 # UUID length def test_process_retry_count_increments( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that retry count increments on failure.""" sample_event.retry_count = 2 mock_eb_connector.put_event.side_effect = Exception('Throttled') mock_repository.update_event_status.return_value = 1 with freeze_time('2023-01-01 12:00:00'): processor.process([sample_event]) # Verify retry calculation uses retry_count=2 # Base: 30 * (1.618^2) = 78.53s. Jitter +/- 5% (3.93s). Range: ~74.6 - 82.5s kwargs = mock_repository.update_event_status.call_args.kwargs next_retry_at = kwargs['next_retry_at'] diff = ( next_retry_at - datetime(2023, 1, 1, 12, 0, 0, tzinfo=timezone.utc) ).total_seconds() assert 74.0 <= diff <= 83.0 def test_process_completed_sets_processed_at( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that processed_at timestamp is set on success.""" mock_repository.update_event_status.return_value = 1 with freeze_time('2023-01-01 15:30:45'): processor.process([sample_event]) kwargs = mock_repository.update_event_status.call_args.kwargs assert 'processed_at' in kwargs assert kwargs['processed_at'] == datetime( 2023, 1, 1, 15, 30, 45, tzinfo=timezone.utc ) def test_process_error_sets_processed_at( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that processed_at timestamp is set on error.""" mock_eb_connector.put_event.side_effect = Exception('Error') mock_repository.update_event_status.return_value = 1 with freeze_time('2023-01-01 16:45:30'): processor.process([sample_event]) # Check both processed_at and next_retry_at are set kwargs = mock_repository.update_event_status.call_args.kwargs # Note: processed_at is used internally but next_retry_at is passed to update assert 'next_retry_at' in kwargs next_retry_at = kwargs['next_retry_at'] # Should be roughly 30 seconds after freeze time expected = datetime(2023, 1, 1, 16, 45, 30, tzinfo=timezone.utc) diff = (next_retry_at - expected).total_seconds() assert 28.5 <= diff <= 31.5 def test_eventbridge_metadata_includes_created_at( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that created_at from outbox is passed to EventBridge metadata.""" sample_event.created_at = datetime(2023, 5, 10, 14, 30, 0, tzinfo=timezone.utc) mock_repository.update_event_status.return_value = 1 processor.process([sample_event]) # Verify EventBridge metadata includes created_at in ISO 8601 format call_args = mock_eb_connector.put_event.call_args metadata = call_args.kwargs['metadata'] assert hasattr(metadata, 'created_at') assert metadata.created_at == '2023-05-10T14:30:00+00:00' def test_eventbridge_metadata_includes_processed_at( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that processed_at is generated and passed to EventBridge metadata.""" mock_repository.update_event_status.return_value = 1 with freeze_time('2023-06-15 10:45:30'): processor.process([sample_event]) # Verify EventBridge metadata includes processed_at in ISO 8601 format call_args = mock_eb_connector.put_event.call_args metadata = call_args.kwargs['metadata'] assert hasattr(metadata, 'processed_at') assert metadata.processed_at == '2023-06-15T10:45:30+00:00' def test_eventbridge_timestamps_are_iso8601_strings( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that timestamps in EventBridge metadata are ISO 8601 formatted strings.""" sample_event.created_at = datetime(2023, 12, 25, 9, 15, 45, tzinfo=timezone.utc) mock_repository.update_event_status.return_value = 1 with freeze_time('2023-12-25 10:00:00'): processor.process([sample_event]) call_args = mock_eb_connector.put_event.call_args metadata = call_args.kwargs['metadata'] # Verify both timestamps are strings in ISO 8601 format assert isinstance(metadata.created_at, str) assert isinstance(metadata.processed_at, str) assert metadata.created_at == '2023-12-25T09:15:45+00:00' assert metadata.processed_at == '2023-12-25T10:00:00+00:00' def test_timestamps_different_for_created_vs_processed( self, processor, mock_repository, mock_eb_connector, sample_event ): """Test that created_at and processed_at can have different values.""" # Event created in the past sample_event.created_at = datetime(2023, 1, 1, 8, 0, 0, tzinfo=timezone.utc) mock_repository.update_event_status.return_value = 1 # Process at a later time with freeze_time('2023-1-1 12:00:00'): processor.process([sample_event]) call_args = mock_eb_connector.put_event.call_args metadata = call_args.kwargs['metadata'] # Verify timestamps are different assert metadata.created_at == '2023-01-01T08:00:00+00:00' assert metadata.processed_at == '2023-01-01T12:00:00+00:00' assert metadata.created_at != metadata.processed_at class TestCalcRetryDelta: """Test calc_retry_delta helper function.""" def test_retry_count_zero(self): """Test retry delta for first retry.""" delta = calc_retry_delta(0) # Base: 30 * (1.618^0) = 30s. Jitter +/- 5% (1.5s) assert 28.5 <= delta.total_seconds() <= 31.5 def test_retry_count_one(self): """Test retry delta for second retry.""" delta = calc_retry_delta(1) # Base: 30 * (1.618^1) = 48.54s. Jitter +/- 5% (2.43s) assert 46.0 <= delta.total_seconds() <= 51.0 def test_retry_count_three(self): """Test retry delta for fourth retry.""" delta = calc_retry_delta(3) # Base: 30 * (1.618^3) = 127.13s. Jitter +/- 5% (6.36s) assert 120.0 <= delta.total_seconds() <= 134.0 def test_retry_delta_returns_timedelta(self): """Test that function returns timedelta object.""" delta = calc_retry_delta(0) assert isinstance(delta, timedelta) def test_retry_delta_always_positive(self): """Test that retry delta is always positive even with negative jitter.""" # Run multiple times to test jitter randomness for _ in range(10): delta = calc_retry_delta(0) assert delta.total_seconds() > 0 class TestTryJson: """Test try_json helper function.""" def test_valid_json_string(self): """Test parsing valid JSON string.""" result = try_json('{"key": "value"}') assert result == {'key': 'value'} def test_invalid_json_string(self): """Test parsing invalid JSON string returns original.""" result = try_json('{invalid}') assert result == '{invalid}' def test_empty_string(self): """Test parsing empty string.""" result = try_json('') assert result == '' def test_non_string_dict(self): """Test non-string dict passes through.""" input_dict = {'already': 'parsed'} result = try_json(input_dict) assert result == input_dict assert result is input_dict # Same object def test_non_string_list(self): """Test non-string list passes through.""" input_list = [1, 2, 3] result = try_json(input_list) assert result == input_list assert result is input_list def test_non_string_none(self): """Test None passes through.""" result = try_json(None) assert result is None def test_non_string_int(self): """Test integer passes through.""" result = try_json(42) assert result == 42 def test_json_array_string(self): """Test parsing JSON array string.""" result = try_json('[1, 2, 3]') assert result == [1, 2, 3] def test_json_nested_structure(self): """Test parsing nested JSON structure.""" json_str = '{"outer": {"inner": "value"}, "list": [1, 2]}' result = try_json(json_str) assert result == {'outer': {'inner': 'value'}, 'list': [1, 2]}