import json from dataclasses import dataclass, field, asdict from datetime import datetime from typing import Dict, Any, List, Union, Callable from smelog.factory import BoundLogger from marshmallow import Schema, ValidationError from charts_messages.config import EVENT from charts_messages.constants import TRACKS_TO_PROCESS, Changelog, GLOBAL_CC_CHART_NAME, GLOBAL_CC_SETTINGS_NAME, \ UserSettings, NO_SETTINGS_HANDLER, MessageMetaType def parse_chart_country_code(delphi_chart_code: str) -> str: if delphi_chart_code == GLOBAL_CC_CHART_NAME: return GLOBAL_CC_SETTINGS_NAME return delphi_chart_code def deserialize_event( logger: BoundLogger, schema: Schema, event: Dict[str, Any] ): result = None try: result = schema.load(event) except ValidationError as err: logger.error(f"Not valid structure of incoming event: {err}") return result 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 get_changelog_key(): """Get a specific key such as additions/removals/moves to get correct event tracks""" return TRACKS_TO_PROCESS[Changelog(EVENT)] def get_tracks_to_process( event_full_data: Dict[str, Any] ) -> List[Dict[str, Union[str, Dict[str, Any]]]]: """Get additions/removals/moves tracks from event data changelog Args: event_full_data: full event data got from User Data API api/service/events/ Returns: List of tracks additions/removals/moves """ return event_full_data["data"]["changelog"].get(get_changelog_key(), []) def get_tracks_asdict( tracks: List[Dict[str, Union[str, Dict[str, Any]]]] ) -> Dict[str, Any]: """Get list of tracks isrc Args: tracks: List of additions/removals/moves tracks Returns: Dict of isrc to track data """ return {track["track"]["isrc"]: track for track in tracks} def get_event_chart_raw_info( event: Dict[str, Any] ) -> Dict[str, Any]: """Get event chart information and remove redundant field for further use Args: event: Returns: Dict[str, Any] """ chart = event["meta"] del chart["id"], chart["domain"] return chart def get_user_settings_handler( version: str ) -> Union[Callable, type(NO_SETTINGS_HANDLER)]: """Get handler for a specific user notifications settings Args: version: str - version of user settings Returns: Handler for a current version or NO_SETTINGS_HANDLER """ handler = NO_SETTINGS_HANDLER try: handler = USER_SETTING_HANDLER[UserSettings(version)] except (ValueError, KeyError): pass return handler def handle_user_notification_settings_v1( chart: Dict[str, str], user_settings_data: Dict[str, Any] ) -> bool: """Handler for User Notifications Settings v1 Args: chart: ex: { "dsp": "spotify", "country_code": "us", "type": "viral", "breakdown": "daily" } user_settings_data: account settings data ex: { "markets": ["global"], "vendors": {"apple": true, "spotify": true}, "notifications": true } } """ market, dsp = chart["country_code"], chart["dsp"] user_notification = user_settings_data["notifications"] user_vendors_activated = user_settings_data["vendors"][dsp] user_markets = user_settings_data["markets"] return all( [ user_notification, user_vendors_activated, market in user_markets ] ) def handle_user_notification_settings_v2( chart: Dict[str, str], user_settings_data: Dict[str, Any] ) -> bool: """Handler for User Notifications Settings v2 Args: chart: ex: { "dsp": "spotify", "country_code": "us", "type": "viral", "breakdown": "daily" } user_settings_data: account settings data ex: { "markets": ["global"], "categories": { "charts":{ "apple":false, "spotify":false }, "starred_tracks":{ "apple":true, "spotify":true }, "starred_playlists": { "apple":true, "spotify":true } }, "notifications":true} } """ market, dsp = chart["country_code"], chart["dsp"] event_settings_key = "starred_tracks" if EVENT in ["entry", "exit", "move"] else "charts" user_notification = user_settings_data["notifications"] user_categories_vendors_activated = user_settings_data["categories"][event_settings_key][dsp] user_markets = user_settings_data["markets"] return all( [ user_notification, user_categories_vendors_activated, market in user_markets ] ) def get_messages_status_length(statuses: Dict[str, Any]): return len(statuses.get("ok", [])), len(statuses.get("failed", [])) USER_SETTING_HANDLER = { UserSettings.V1: handle_user_notification_settings_v1, UserSettings.V2: handle_user_notification_settings_v2 } def create_post_favorites_body( isrc: List[str], include: List[str], entity_type: str = "track", device_is_active: bool = True, settings_type: str = "mobile", ): """Create a body for """ @dataclass class Entity: type: str id: list = field(default_factory=list) @dataclass class Device: is_active: bool @dataclass class Settings: type: str @dataclass class PostFavoritesBody: entity: Entity device: Device settings: Settings include: list = field(default_factory=list) return asdict( PostFavoritesBody( entity=Entity( type=entity_type, id=isrc ), device=Device( is_active=device_is_active ), settings=Settings( type=settings_type ), include=include ) ) def get_message_meta_type(): return MessageMetaType[Changelog(EVENT)] 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"]) single_event = json.loads(body.get("Message", {})) if event: result.append(single_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