"""trending_tracks_discrete workflow tasks.""" from datetime import datetime from datetime import timedelta from ddtrace import tracer as ddtracer from garcon import task from snowflake_connector.snowflake_conn import get_session from snowflake_connector.snowflake_conn import SQLLoader from activity_detector import base_config from activity_detector.flows.trending_tracks_discrete import config from activity_detector.utils import dynamodb from activity_detector.utils import ows_notifications sql_loader = SQLLoader(__file__) @task.decorate(timeout=1000) def bootstrap(activity, date, run_id): """Startup flow with parameters from exec command.""" if not date: date = (datetime.today() - timedelta( days=config.DEFAULT_SHIFT_DAYS)).strftime('%Y-%m-%d') return dict( target_date=date ) @task.decorate(timeout=300) def check_dynamo_status(activity, date): """Check the activity status in DynamoDB. Args: activity (ActivityWorker): The activity worker date (str): YYYY-MM-DD target date """ res = dynamodb.get_status(config.ACTIVITY_TYPE_NAME, date) status = res and res['status'] if status and status == config.STATUS_PROCESSED_NOTIF_SENT: return {'should_run': False} return {'should_run': True} def _load_query(filename): return sql_loader.load_query(filename).format(env=base_config.ENVIRONMENT) @task.decorate(timeout=7200) def detect_spikes(activity, date): """Find trending track spikes and send data to ows-notifications. Args: activity (ActivityWorker): The activity worker date (str): YYYY-MM-DD target date Returns: dict: resulting number of failures when sending to API """ return _detect_spikes(activity, date) @ddtracer.wrap(resource='task.trending_tracks_discrete', name='detect_trending_track_spikes') def _detect_spikes(activity, date): """Find trending track spikes to send them to ows-notifications. @ddtracer.wrap decorator is incompatible with garcon task runner and causes failures in the trending_tracks_discrete workflow execution. To enable full tracing without impacting the workflow, all related logic was moved here instead. """ check_activity_sql = _load_query('check_activity') insert_activity_sql = _load_query('insert_activity') get_activity_sql = _load_query('get_activity') store_ids = config.STORE_IDS rows = [] with get_session() as session: check_rows = session.execute( check_activity_sql, { 'date': date, 'store_ids': store_ids } ).fetchall() if len(check_rows) != len(store_ids): found_store_ids = [x['store_id'] for x in check_rows] missing_store_ids = set(store_ids) - set(found_store_ids) raise Exception(f'Missing data for store(s) {missing_store_ids} for {date}!') # noqa:E501 session.execute( insert_activity_sql, { 'date': date, 'store_ids': store_ids, 'min_spike_score': config.MIN_SPIKE_SCORE, 'min_days_of_data': config.MIN_DAYS_OF_DATA, 'min_streams': config.MIN_STREAMS, 'max_percent_diff': config.MAX_PERCENT_DIFF, 'top_markets_number': config.TOP_MARKETS_NUMBER } ) session.execute('commit') rows = session.execute( get_activity_sql, { 'date': date } ).fetchall() activity.logger.info(f'Found {len(rows)} trending tracks discrete for {date}') failures = 0 for row in rows: success = ows_notifications.create_trending_track_notification( date, row['countryname'], row['storename'].lower(), row['pct_diff'], row['total_streams'], row['track_unique_id'], row['isrc'], row['labelid'], row['subaccountid'] ) if not success: failures += 1 return dict( failures=failures ) @task.decorate(timeout=300) def set_dynamo_status(activity, date, failures): """Set status of workflow run according to params. Args: activity (ActivityWorker): The activity worker date (str): YYYY-MM-DD target date failures (int): number of API request failures """ status = config.STATUS_PROCESSED_NOTIF_SENT if not failures \ else config.STATUS_PROCESSED dynamodb.set_status( config.ACTIVITY_TYPE_NAME, date, status)