"""Handlers for SalesForce update messages.""" from typing import Dict from kafka_utils.consumer.source.mapping import EventSourceMessage from kafka_utils.consumer.deserializer.simple_json import JSONDeserializer from src.models.fan import ProfileUpdatedV1 from src.models.fan import ProfileDeletedV1 from src.models.fan import SForceFanUpdate OP_DELETE = 'delete' OP_UPDATE = 'update' def get_operation_type(msg: Dict) -> str: """Get operation type from message key.""" if msg.get('deletedAt'): return OP_DELETE elif msg.get('updatedAt'): return OP_UPDATE else: raise Exception('Unable to get operation type from message.') def prepare_message(msg: Dict, operation: str) -> SForceFanUpdate: """Prepare a message depending on the operation type.""" if operation == OP_UPDATE: fan_update = ProfileUpdatedV1(**msg) return SForceFanUpdate.from_profile_update(fan_update) elif operation == OP_DELETE: fan_delete = ProfileDeletedV1(**msg) return SForceFanUpdate.from_profile_delete(fan_delete) def read_source_event(event: Dict) -> Dict: """Get message from EventSource.""" json_deserialzier = JSONDeserializer() _, msk_message = next(iter(EventSourceMessage(event))) return json_deserialzier.deserialize(msk_message.value)