"""Helpers.""" import json def load_file_json(filename): """Load json from file.""" with open(f'tests/unit/mock_event_data/{filename}') as event_file: read_data = event_file.read() return json.loads(read_data) def mock_product_review_event_json(): """Mock data.""" return { 'operation': { 'type': 'create', 'context': 'new', 'timestamp': 1655390557 }, 'payload': { 'product_id': 1, 'review_queue_id': 1, 'submission_type': 'New', 'review_note': '', 'user_id': 'linus' } } def delivery_callback(err, msg): """Handle produce result.""" if err is not None: raise Exception(err) print('SUCCESS!') print('KEY: ', msg.key()) print('TOPIC: ', msg.topic()) print('PARTITION:', msg.partition()) print('OFFSET: ', msg.offset()) def debug_msk_message( batch_key, msk_message, key_deserializer=None, value_deserializer=None ): """Print kafka message properties.""" print('SUCCESS!') print('BATCH_KEY: ', batch_key) print('TOPIC: ', msk_message.topic) print('OFFSET: ', msk_message.offset) if key_deserializer is not None: key = key_deserializer.deserialize(msk_message.key, msk_message.topic, 'key') print('KEY: ', key) else: print('KEY: ', msk_message.key) if value_deserializer is not None: value = value_deserializer.deserialize(msk_message.value, msk_message.topic, 'value') print('VALUE: ', value) else: print('VALUE: ', msk_message.value)