"""Test event producer.""" from unittest.mock import MagicMock from unittest.mock import patch from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer from kafka_utils.testing.unit.fixture_plugin import MockProducer @patch('kafka_utils.producer.event.Producer', new=MockProducer) def test_event_producer(): """Test event producer.""" event_producer = EventProducer( 'foo', StringSerializer(), SimpleJSONSerializer(), 'SSL', {'foo': 'bar'}) assert event_producer mock_callback = MagicMock() event_producer.produce('foo', 'bar', {'foo': 'bar'}, mock_callback) assert mock_callback.call_count == 1 with event_producer as ep: ep.produce('foo', 'bar', {'foo': 'bar'}, mock_callback) assert mock_callback.call_count == 2 @patch('kafka_utils.producer.event.Producer', new=MockProducer) def test_event_producer_no_auto_flush(): """Test event producer without flushing messages.""" event_producer = EventProducer('foo') assert event_producer mock_callback = MagicMock() event_producer.produce('foo', 'bar', '{"foo": "bar"}', mock_callback, auto_flush=False) assert mock_callback.call_count == 0 event_producer.producer.assert_messages([ { 'topic': 'foo', 'value': '{"foo": "bar"}', 'key': 'bar' } ]) event_producer.producer.assert_call_count('flush', 0) event_producer.producer.flush() assert mock_callback.call_count == 1 event_producer.producer.assert_call_count('flush', 1) @patch('kafka_utils.producer.event.Producer', new=MockProducer) def test_event_producer_no_serializers(): """Test event producer without passing in serializers.""" event_producer = EventProducer('foo') assert event_producer mock_callback = MagicMock() event_producer.produce('foo', 'bar', '{"foo": "bar"}', mock_callback) event_producer.producer.assert_messages([ { 'topic': 'foo', 'value': '{"foo": "bar"}', 'key': 'bar', } ]) @patch('kafka_utils.producer.event.Producer', new=MockProducer) def test_event_producer_produce_serializers(): """Test event producer passing serializers in the product method.""" event_producer = EventProducer('foo') assert event_producer mock_callback = MagicMock() event_producer.produce( 'foo', {'baz': 'qux'}, {'foo': 'bar'}, mock_callback, key_serializer=SimpleJSONSerializer(), value_serializer=SimpleJSONSerializer() ) event_producer.producer.assert_messages([{ 'topic': 'foo', 'key': b'{"baz": "qux"}', 'value': b'{"foo": "bar"}', }] ) @patch('kafka_utils.producer.event.Producer', new=MockProducer) def test_event_producer_without_callback(mocker): """Test event producer without callback.""" event_producer = EventProducer( 'foo', StringSerializer(), SimpleJSONSerializer(), 'SSL') produce_spy = mocker.spy(event_producer.producer, 'produce') event_producer.produce('foo', 'bar', {'foo': 'bar'}) produce_spy.assert_called_once_with( topic='foo', key=b'bar', value=b'{"foo": "bar"}', on_delivery=None)