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_delphi_events.config import EVENT_TYPES from apollo_delphi_events.constants import APP_NAME, DOMAIN, EventType from apollo_delphi_events.schemas.input import EVENT_TYPE_TO_INPUT_SCHEMA def get_input_schema(event: dict) -> Schema: return EVENT_TYPE_TO_INPUT_SCHEMA[EventType(event["meta"]["event_type"])] def is_supported(event: dict) -> bool: domain = event.get("meta", {}).get("domain") event_type = event.get("meta", {}).get("event_type") return domain == DOMAIN and event_type in EVENT_TYPES 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_input_event( logger: BoundLogger, raw_event: dict, parsed_ids: set, current_dt: datetime = None ) -> Optional[dict]: if not is_supported(raw_event): logger.info(f"{APP_NAME}.parse_input_event got unsupported event, skipped\n{raw_event}.") return schema = get_input_schema(raw_event) try: event = schema.load(raw_event) except ValidationError as ex: logger.warning( f" {APP_NAME}.parse_input_event got invalid structure event, skipped\n" f" Errors: {ex.messages}\n" f" Event: {raw_event}." ) return event_id = event["id"] if event_id in parsed_ids: logger.warning(f" {APP_NAME}.parse_input_event got chunk-duplicated {event_id} event, skipped\n{event}.") return if not is_actual(event, current_dt=current_dt): logger.warning(f" {APP_NAME}.parse_input_event got outdated {event_id} event, skipped\n{event}.") return parsed_ids.add(event_id) return event def parse_input( logger: BoundLogger, input: Union[List[Dict[str, Any]], Dict[str, Any]], current_dt: datetime = None ) -> List[dict]: current_dt = current_dt or datetime.utcnow() events = [] parsed_events_ids = set() # to provide events uniqueness within the received chunk: Delphi can send duplicates records = input.get("Records", []) for record in records: body = json.loads(record.get("body", "{}")) raw_event = json.loads(body.get("Message", "{}")) event = parse_input_event(logger, raw_event, parsed_events_ids, current_dt=current_dt) if event: events.append(event) if not events: logger.error(f" {APP_NAME} got 0/{len(records)} parsed. Nothing to process.") else: logger.info(f" {APP_NAME} got {len(events)}/{len(records)} successfully parsed.") return events