"""Integration tests for lambda function.""" import json import os from datetime import timezone from unittest.mock import MagicMock, patch import pymysql import pytest from src import app from src.enums import AbacusOutboxStatus DB_HOST = os.environ.get('MYSQL_DB_HOST', 'mysql') DB_USER = os.environ.get('MYSQL_USER', 'royalties') DB_PASS = os.environ.get('MYSQL_PASSWORD', '1234') DB_NAME = os.environ.get('MYSQL_DATABASE', 'royalty_accounting') DB_PORT = int(os.environ.get('MYSQL_DB_PORT', 3306)) @pytest.fixture(scope='module') def db_conn(): """Create a database connection.""" conn = pymysql.connect( host=DB_HOST, user=DB_USER, password=DB_PASS, database=DB_NAME, port=DB_PORT, cursorclass=pymysql.cursors.DictCursor, autocommit=True, ) yield conn conn.close() @pytest.fixture(scope='function') def seed_data(db_conn): """Seed data for the test.""" with db_conn.cursor() as cursor: cursor.execute('DELETE FROM abacus_outbox') # Insert pending events query = """ INSERT INTO abacus_outbox (target_type, target_id, event_type, correlation_id, details, status) VALUES (%s, %s, %s, %s, %s, %s) """ data = [ ( 'test_entity', 101, 'test.created', 'corr-1', json.dumps({'id': 101}), AbacusOutboxStatus.PENDING, ), ( 'test_entity', 102, 'test.updated', 'corr-2', json.dumps({'id': 102}), AbacusOutboxStatus.PENDING, ), ] cursor.executemany(query, data) yield # Cleanup with db_conn.cursor() as cursor: cursor.execute('DELETE FROM abacus_outbox') @patch('src.app.get_events_client') def test_lambda_handler_success(mock_get_client, seed_data, db_conn): """Test successful lambda execution with DB integration.""" # Mock EventBridge client mock_eb = MagicMock() mock_eb.put_events.return_value = {'FailedEntryCount': 0} mock_get_client.return_value = mock_eb # Patch config to use the test DB credentials with patch.dict( 'config.MYSQL_CONFIG', { 'host': DB_HOST, 'user': DB_USER, 'password': DB_PASS, 'database': DB_NAME, 'port': DB_PORT, }, ): context = MagicMock() # Execute Handler response = app.handler({}, context) # Verify Response assert response['failed'] == 0 assert response['processed'] == 2 assert response['skipped'] == 0 assert response['total'] == 2 # Verify DB updates with db_conn.cursor() as cursor: cursor.execute('SELECT * FROM abacus_outbox ORDER BY target_id ASC') events = cursor.fetchall() assert len(events) == 2 assert events[0]['status'] == AbacusOutboxStatus.COMPLETED assert events[0]['processed_at'] is not None assert events[1]['status'] == AbacusOutboxStatus.COMPLETED assert events[1]['processed_at'] is not None # Verify EventBridge calls assert mock_eb.put_events.call_count == 2 # Check first call details call_args = mock_eb.put_events.call_args_list[0] entry = call_args.kwargs['Entries'][0] detail = json.loads(entry['Detail']) assert entry['Source'] == 'abacus.outbox' assert entry['DetailType'] == 'test.created' assert detail['metadata']['target_id'] == 101 assert detail['metadata']['correlation_id'] == 'corr-1' assert detail['data'] == {'id': 101} # Verify timestamps are included in metadata assert 'created_at' in detail['metadata'] assert 'processed_at' in detail['metadata'] assert detail['metadata']['created_at'] is not None assert detail['metadata']['processed_at'] is not None # Verify ISO 8601 format (should contain 'T' and timezone) assert 'T' in detail['metadata']['created_at'] assert 'T' in detail['metadata']['processed_at'] @patch('src.app.get_events_client') def test_lambda_timestamps_end_to_end(mock_get_client, seed_data, db_conn): """Test that timestamps flow correctly from DB through to EventBridge.""" # Mock EventBridge client mock_eb = MagicMock() mock_eb.put_events.return_value = {'FailedEntryCount': 0} mock_get_client.return_value = mock_eb # Get the created_at timestamp from DB before processing with db_conn.cursor() as cursor: cursor.execute('SELECT created_at FROM abacus_outbox WHERE target_id = 101') # Patch config to use the test DB credentials with patch.dict( 'config.MYSQL_CONFIG', { 'host': DB_HOST, 'user': DB_USER, 'password': DB_PASS, 'database': DB_NAME, 'port': DB_PORT, }, ): context = MagicMock() # Execute Handler app.handler({}, context) # Extract EventBridge metadata call_args = mock_eb.put_events.call_args_list[0] entry = call_args.kwargs['Entries'][0] detail = json.loads(entry['Detail']) metadata = detail['metadata'] # Verify created_at matches DB timestamp assert metadata['created_at'] is not None # Verify processed_at is after created_at (measuring latency) created_at_str = metadata['created_at'] processed_at_str = metadata['processed_at'] # Both should be valid ISO 8601 timestamps assert 'T' in created_at_str assert 'T' in processed_at_str # In most cases, processed_at should be >= created_at # (unless the event was created just now and processed immediately) from datetime import datetime created_at = datetime.fromisoformat(created_at_str) processed_at = datetime.fromisoformat(processed_at_str) # Ensure both are timezone-aware for comparison if created_at.tzinfo is None: created_at = created_at.replace(tzinfo=timezone.utc) if processed_at.tzinfo is None: processed_at = processed_at.replace(tzinfo=timezone.utc) assert processed_at >= created_at