import json from datetime import datetime from typing import Dict, Any, List, Union, Tuple from smelog.factory import BoundLogger from marshmallow import Schema, ValidationError from charts_feed_messages.utils.logger_messages import get_event_schema_logger_message_builder def deserialize_event( logger: BoundLogger, schema: Schema, event: Dict[str, Any] ) -> Union[Dict[str, Any], None]: """Deserialize single incoming Message Created event using passed Marshmallow schema and log if event has incorrect structure Args: logger: Logger schema: Marshmallow schema event: Incoming Message Created event Returns: deserialized Message Created Event or None if """ result = None try: result = schema.load(event) except ValidationError as err: logger.error( get_event_schema_logger_message_builder(schema)(event, err) ) return result def handle_events_deserialization( logger: BoundLogger, schema: Schema, events: List[Dict[str, Any]] ) -> List[Dict[str, Any]]: """Handle deserialization of incoming event(s) array one by one We have to loop through events to easily log events with incorrect structure and continue without fail Args: logger: Logger schema: Marshmallow schema events: List of Incoming Message Created event(s) Returns: List of deserialized events """ actual = [] for event in events: deserialized_event = deserialize_event(logger, schema, event) if deserialized_event: actual.append(deserialized_event) return actual def check_received_event( received_event: Dict[str, Any] ) -> bool: """Check if event is actual and needs to be processed Args: received_event: event that triggered handler Returns: bool """ current_dt = datetime.utcnow() event_dt = received_event["created_at"] event_ttl = received_event["ttl"] return (current_dt - event_dt).total_seconds() <= event_ttl def handle_check_received_events( logger: BoundLogger, events: List[Dict[str, Any]] ) -> List[int]: """Handle actuality of incoming deserialized events array and log events that are not actual anymore Args: logger: Logger events: List of deserialized events Returns: List of deserialized event id(s) """ result = [] for event in events: if not check_received_event(event): logger.warning(f"{event['id']} - Event is no more actual") continue result.append(event["id"]) return result def get_messages_status_length(statuses: Dict[str, Any]) -> Tuple[int, int]: """Get counts of success and failed pushed messages Args: statuses: Dict of message statuses ex: "ok": [ { "index": int, # index in request 'data' list "id": int } ], "failed": [ { "index": int, # index in request 'data' list "error": { "type": str, "data": Optional[dict] } } ] Returns: count of successfully and failures of sent messages """ return len(statuses.get("ok", [])), len(statuses.get("failed", [])) def handle_diff_source_events(logger: BoundLogger, event: Any) -> List[Dict[str, Any]]: try: result = [] records = event.get("Records", []) for record in records: body = json.loads(record["body"]) event = json.loads(body.get("Message", {})) if event: result.append(event) except AttributeError: logger.info(f"Got an event from EVENT BRIDGE: \n" f"{event} \n" f"With {len(event)} messages to process") return event logger.info(f"Got an SQS event: \n" f"EVENT:\n" f"{event} \n" f"With {len(result)} messages to process \n" f"Messages:\n" f"{result}") return result