"""Handlers for SalesForce Subscription 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.subscription import ProfileSubscriptionsUpdatedV1 from src.models.subscription import SForceSubscriptionUpdate OP_UPDATE = 'update' def get_operation_type(msg: Dict) -> str: """Get operation type from message key.""" if msg.get('updatedAt'): return OP_UPDATE return None def prepare_message(msg: Dict, crm_id: str) -> SForceSubscriptionUpdate: """Prepare a message depending on the operation type.""" subscription_update = ProfileSubscriptionsUpdatedV1(**msg) return SForceSubscriptionUpdate.from_profile_update(subscription_update, crm_id) 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)