import json from datetime import datetime from typing import Any, Dict, List, Optional, Union from marshmallow import Schema, ValidationError from smelog.factory import BindableLogger from apollo_messages_views.config import EVENT_TYPES, FULL_APP_NAME from apollo_messages_views.constants.base import EventView from apollo_messages_views.schemas.base import BaseRawMessage from apollo_messages_views.schemas.factories import get_schema def parse_raw_event( logger: BindableLogger, raw_event: dict, schema: Schema, ) -> Optional[dict]: event_type = raw_event.get("meta", {}).get("type", "").lower() if event_type not in EVENT_TYPES: logger.warning( f"{FULL_APP_NAME}.parse_raw_event got unsupported event type, skipped\n" f"Errors: got event_type={event_type}, but allowed values are: {EVENT_TYPES}\n" f"Event: {raw_event}." ) return try: event = schema.load(raw_event) except ValidationError as ex: logger.warning( f"{FULL_APP_NAME}.parse_raw_event got invalid event, skipped\n" f"Errors: {ex.messages}\n" f"Event: {raw_event}." ) return return event def parse_raw_input( logger: BindableLogger, input: Union[List[Dict[str, Any]], Dict[str, Any]], current_dt: datetime = None ) -> List[dict]: current_dt = current_dt or datetime.utcnow() events = [] records = input.get("Records", []) schema = BaseRawMessage() schema.context["current_dt"] = current_dt for record in records: body = json.loads(record.get("body", {})) raw_event = json.loads(body.get("Message", {})) event = parse_raw_event(logger, raw_event, schema) 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: BindableLogger, events: List[dict], event_view: EventView = None) -> List[dict]: parsed_events = [] for event in events: try: parsed_event = get_schema(event["meta"]["type"], event_view=event_view and event_view.value).load(event) parsed_events.append(parsed_event) except ValidationError as ex: logger.error( f"{FULL_APP_NAME}.parse_full_events got broken full event:\nErrors: {ex.messages}\n" f"Event: {event}" ) 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