"""Test event handling.""" import base64 from contextlib import nullcontext as does_not_raise import json from content_utils.exceptions import IneligibleEventError from kafka_utils.consumer.message import debezium from kafka_utils.consumer.source.mapping import EventSourceMessage import pytest from src.logic.event_handling import ReviewQueueItemExplicitSkip @pytest.mark.parametrize(( 'test_description', 'before_status', 'after_status', 'before_locked_by_user_id', 'after_locked_by_user_id', 'before_locked_until_datetime', 'after_locked_until_datetime', 'operation', 'expected_raise', 'expected_raise_message', 'expected_result' ), [ ( 'locked_by_user_id change', 'new', 'new', None, 'test-locker', None, None, 'u', does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': {'product_id': 4180642, 'locked_by_user_id': 'test-locker'}, 'operation_type': 'update_index' } ), ( 'locked_until_datetime change None to value', 'new', 'new', None, None, None, 1232123212321212, 'u', does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': { 'product_id': 4180642, 'locked_until_datetime': '2009-01-16T16:26:52.321212Z' }, 'operation_type': 'update_index' } ), ( 'locked_until_datetime change value to None', 'new', 'new', None, None, 1232123212321212, None, 'u', does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': {'product_id': 4180642, 'locked_until_datetime': None}, 'operation_type': 'update_index' } ), ( 'locked_by_user_id change and locked_until_datetime change', 'new', 'new', None, 'test-locker', None, 1232123212321212, 'u', does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': { 'product_id': 4180642, 'locked_by_user_id': 'test-locker', 'locked_until_datetime': '2009-01-16T16:26:52.321212Z' }, 'operation_type': 'update_index' } ), ( 'status before not new', 'asdf', 'new', None, None, None, None, 'u', pytest.raises(IneligibleEventError), 'skip operation: u, before: {"id": 5340, "product_id": 4180642, "status": "asdf", ' '"created_datetime": 1232123212321232, "locked_by_user_id": null, ' '"locked_until_datetime": null}, after: {"id": 5340, "product_id": 4180642, ' '"status": "new", "created_datetime": 1232123212321232, "locked_by_user_id": ' 'null, "locked_until_datetime": null}', None ), ( 'status after not complete', 'new', 'asdf', None, None, None, None, 'u', pytest.raises(IneligibleEventError), 'skip operation: u, before: {"id": 5340, "product_id": 4180642, "status": "new", ' '"created_datetime": 1232123212321232, "locked_by_user_id": null, ' '"locked_until_datetime": null}, after: {"id": 5340, "product_id": 4180642, ' '"status": "asdf", "created_datetime": 1232123212321232, "locked_by_user_id": ' 'null, "locked_until_datetime": null}', None ), ( 'status new -> complete', 'new', 'complete', None, None, None, None, 'u', does_not_raise(), 'None', { 'operation_type': 'deindex', 'payload': None, 'product_id': 4180642, 'review_queue_id': 5340 } ), ( 'Add back to index when applying -> new', 'applying', 'new', None, None, None, None, 'u', does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': { 'id': 5340, 'product_id': 4180642, 'status': 'new', 'created_datetime': 1232123212321232, 'locked_by_user_id': None, 'locked_until_datetime': None }, 'operation_type': 'add_to_index' } ), ]) def test_parse_message_payload( test_description, before_status, after_status, before_locked_by_user_id, after_locked_by_user_id, before_locked_until_datetime, after_locked_until_datetime, operation, expected_raise, expected_raise_message, expected_result ): """Test parse_message_payload.""" from src.logic.event_handling import parse_message_payload review_queue_id = 5340 product_id = 4180642 created_datetime = 1232123212321232 kafka_event = { 'records': { 'cdc.contentReview.reviewQueue-0': [{ 'topic': 'cdc.contentReview.reviewQueue', 'value': base64.b64encode(json.dumps({ 'payload': { 'before': { 'id': review_queue_id, 'product_id': product_id, 'status': before_status, 'created_datetime': created_datetime, 'locked_by_user_id': before_locked_by_user_id, 'locked_until_datetime': before_locked_until_datetime }, 'after': { 'id': review_queue_id, 'product_id': product_id, 'status': after_status, 'created_datetime': created_datetime, 'locked_by_user_id': after_locked_by_user_id, 'locked_until_datetime': after_locked_until_datetime }, 'op': operation, } }).encode('ascii')).decode('ascii') }] } } _, msk_message = next(iter(EventSourceMessage(kafka_event))) db_msg = debezium.DebeziumMessage( msk_message.value, msk_message.topic, allowed_event_ops=['c', 'u', 'd'], ) response = None with expected_raise as er: response = parse_message_payload(db_msg) assert response == expected_result assert str(getattr(er, 'value', None)) == expected_raise_message def test_parse_message_payload_explicit_skip(monkeypatch): """Test parse_message_payload explicit skip with isolated mock.""" from src import config from src.logic.event_handling import parse_message_payload review_queue_id = 5340 product_id = 4180642 created_datetime = 1232123212321232 monkeypatch.setattr(config, 'SKIP_REVIEW_QUEUE_ITEMS', [review_queue_id]) kafka_event = { 'records': { 'cdc.contentReview.reviewQueue-0': [{ 'topic': 'cdc.contentReview.reviewQueue', 'value': base64.b64encode(json.dumps({ 'payload': { 'before': { 'id': review_queue_id, 'product_id': product_id, 'status': 'new', 'created_datetime': created_datetime, 'locked_by_user_id': None, 'locked_until_datetime': None }, 'after': { 'id': review_queue_id, 'product_id': product_id, 'status': 'complete', 'created_datetime': created_datetime, 'locked_by_user_id': None, 'locked_until_datetime': None }, 'op': 'u', } }).encode('ascii')).decode('ascii') }] } } _, msk_message = next(iter(EventSourceMessage(kafka_event))) db_msg = debezium.DebeziumMessage( msk_message.value, msk_message.topic, allowed_event_ops=['c', 'u', 'd'], ) with pytest.raises(ReviewQueueItemExplicitSkip) as er: parse_message_payload(db_msg) assert str(getattr(er, 'value', None)) == 'Explicit Skip' @pytest.mark.parametrize(( 'test_description', 'after_status', 'expected_raise', 'expected_raise_message', 'expected_result' ), [ ( 'Add to index', None, does_not_raise(), 'None', { 'product_id': 4180642, 'review_queue_id': 5340, 'payload': { 'id': 5340, 'product_id': 4180642, 'status': None, 'created_datetime': 1232123212321232, 'locked_by_user_id': None, 'locked_until_datetime': None }, 'operation_type': 'add_to_index' } ), ( 'No-op when create already complete', 'complete', pytest.raises(IneligibleEventError), 'skip operation: c, before: null, after: {"id": 5340, "product_id": 4180642, ' '"status": "complete", "created_datetime": 1232123212321232, "locked_by_user_id": ' 'null, "locked_until_datetime": null}', None ), ]) def test_parse_message_payload_create_operations( test_description, after_status, expected_raise, expected_raise_message, expected_result ): """Test parse_message_payload for create operations with empty before payload.""" from src.logic.event_handling import parse_message_payload review_queue_id = 5340 product_id = 4180642 created_datetime = 1232123212321232 kafka_event = { 'records': { 'cdc.contentReview.reviewQueue-0': [{ 'topic': 'cdc.contentReview.reviewQueue', 'value': base64.b64encode(json.dumps({ 'payload': { 'before': None, 'after': { 'id': review_queue_id, 'product_id': product_id, 'status': after_status, 'created_datetime': created_datetime, 'locked_by_user_id': None, 'locked_until_datetime': None }, 'op': 'c', } }).encode('ascii')).decode('ascii') }] } } _, msk_message = next(iter(EventSourceMessage(kafka_event))) db_msg = debezium.DebeziumMessage( msk_message.value, msk_message.topic, allowed_event_ops=['c', 'u', 'd'], ) response = None with expected_raise as er: response = parse_message_payload(db_msg) assert response == expected_result assert str(getattr(er, 'value', None)) == expected_raise_message