import json from datetime import datetime from typing import Any, Dict, List, Optional, Union from marshmallow import Schema, ValidationError from smelog.factory import BoundLogger from apollo_playlists_messages.config import EVENT, EVENT_HOLDER, FULL_APP_NAME from apollo_playlists_messages.schemas.events_raw import InputRawEvent def is_supported(event: dict) -> bool: return event.get("code", "").lower() == EVENT.value and event.get("app", "").lower() == EVENT_HOLDER def is_actual(event: dict, current_dt: datetime = None) -> bool: now = current_dt or datetime.utcnow() return (now - event["created_at"]).total_seconds() <= event["ttl"] def parse_raw_event( logger: BoundLogger, raw_event: dict, schema: Schema, current_dt: datetime = None ) -> Optional[dict]: if not is_supported(raw_event): logger.warning(f"{FULL_APP_NAME}.parse_raw_event got unsupported event, skipped\n{raw_event}.") return try: event = schema.load(raw_event) except ValidationError as ex: logger.warning( f"{FULL_APP_NAME}.parse_raw_event got invalid structure event, skipped\n" f"Errors: {ex.messages}\n" f"Event: {raw_event}." ) return if not is_actual(event, current_dt=current_dt): logger.warning(f"{FULL_APP_NAME}.parse_raw_event got outdated {event['id']} event, skipped\n{event}.") return return event def parse_raw_input(logger: BoundLogger, input: Union[List[Dict[str, Any]], Dict[str, Any]]) -> List[dict]: current_dt = datetime.utcnow() events = [] records = input.get("Records", []) schema = InputRawEvent() for record in records: body = json.loads(record.get("body", '""')) or {} raw_event = json.loads(body.get("Message", '""')) or {} event = parse_raw_event(logger, raw_event, schema, current_dt=current_dt) if event: events.append(event) if not events: logger.error(f"{FULL_APP_NAME} got 0/{len(records)} raw events parsed. Nothing to process.") else: logger.info(f"{FULL_APP_NAME} got {len(events)}/{len(records)} raw events successfully parsed.") return events def parse_full_events(logger: BoundLogger, events: List[dict], schema: Schema) -> List[dict]: parsed_events = [] for event in events: try: parsed_event = schema.load(event) parsed_events.append(parsed_event) except ValidationError as ex: logger.error(f"{FULL_APP_NAME}.parse_full_events got some broken full events:\n" f"{ex.messages}") if not parsed_events: logger.error(f"{FULL_APP_NAME} got 0/{len(events)} full events parsed. Nothing to process.") else: logger.info(f"{FULL_APP_NAME} got {len(parsed_events)}/{len(events)} full events successfully parsed.") return parsed_events