"""Utility functions for ows-notifications endpoints.""" import json import time from ddtrace import tracer as ddtracer from oto import response from owsrequest import request from activity_detector import base_config from activity_detector.utils import errors, sentry @ddtracer.wrap(resource='task.trending_tracks_discrete', name='create_trending_track_notification') def create_trending_track_notification( date, region, dsp, pct_diff, day_streams, track_id, track_isrc, track_vendor_id, track_subaccount_id, retries=5): """Send trending track event to ows-notifications for processing. Args: date (str): YYYY-MM-DD date format when event happened region (str): Region where spike happened dsp (str): Streaming platform where spike happened pct_diff (int): Percent change from previous day of data day_streams (int): Daily streams for date spike occured track_id (int): Internal unique identifier of track track_isrc (str): Identifier of sound recording track is based on track_vendor_id (int): Internal identifier of vendor label of track track_subaccount_id (int): Internal identifier of subaccount label Returns: bool: success or failure """ if ddtracer.enabled: dd_span = ddtracer.current_span() if dd_span: dd_span.set_tag('date', date) dd_span.set_tag('region', region) dd_span.set_tag('dsp', dsp) dd_span.set_tag('pct_diff', pct_diff) dd_span.set_tag('day_streams', day_streams) dd_span.set_tag('track_id', track_isrc) dd_span.set_tag('track_isrc', track_isrc) dd_span.set_tag('track_vendor_id', track_vendor_id) dd_span.set_tag('track_subaccount_id', track_subaccount_id) body = { 'date': date, 'region': region, 'dsp': dsp, 'percent_diff': pct_diff, 'day_streams': day_streams, 'track': { 'id': track_id, 'isrc': track_isrc, 'vendor_id': track_vendor_id, 'subaccount_id': track_subaccount_id } } result = request.process( application='swf-activity-detector', environment=base_config.ENVIRONMENT, method='post', service_name='ows-notifications', path='/activity/trending_track', json=body) if result.status_code not in [200, 201, 409]: if retries > 0: if result.status_code == 429: time.sleep(120) return create_trending_track_notification( date, region, dsp, pct_diff, day_streams, track_id, track_isrc, track_vendor_id, track_subaccount_id, retries=retries - 1) else: error = {} try: error = result.json() except json.decoder.JSONDecodeError: pass error['notification'] = body sentry.send_response_to_sentry( response.create_error_response( status=result.status_code, code=errors.ERROR_CODE_OWS_NOTIFICATIONS_REQUEST, message=error), result.status_code ) return False return True @ddtracer.wrap(resource='task.playlist_placements', name='create_playlist_placement_notification') def create_playlist_placement_notification( timestamp, playlist_name, playlist_id, playlist_dsp, playlist_storeid, playlist_rank, sound_recording_isrc, tracks_list, retries=5): """Send playlist placement event to ows-notifications for processing. Args: timestamp (str): ISO datestamp when placement happened playlist_name (str): Human readable playlist title playlist_id (str): DSP specific playlist identifier playlist_dsp (str): DSP where playlist exists playlist_storeid (int): store ID for playlist DSP playlist_rank (int): Relative importance of playlist sound_recording_isrc (str): Unique identifier of sound recording tracks_list (list): track_id (int): Internal identifier of track vendor_id (int): Internal identifier of vendor label subaccount_id (int): Internal identifier of subaccount (optional) retries (int): Number of times to retry failures Returns: bool: success or failure """ if ddtracer.enabled: dd_span = ddtracer.current_span() if dd_span: dd_span.set_tag('playlist_name', playlist_name) dd_span.set_tag('playlist_id', playlist_id) dd_span.set_tag('playlist_dsp', playlist_dsp) dd_span.set_tag('playlist_storeid', playlist_storeid) dd_span.set_tag('sound_recording_isrc', sound_recording_isrc) tracks = [ { 'id': x['track_id'], 'vendor_id': x['vendor_id'], 'subaccount_id': x.get('subaccount_id') } for x in tracks_list ] body = { 'timestamp': timestamp, 'playlist': { 'id': playlist_id, 'name': playlist_name, 'dsp': playlist_dsp, 'store_id': playlist_storeid, 'rank': playlist_rank }, 'sound_recording': { 'isrc': sound_recording_isrc, 'tracks': tracks } } result = request.process( application='swf-activity-detector', environment=base_config.ENVIRONMENT, method='post', service_name='ows-notifications', path='/activity/playlist_placement', json=body) if result.status_code not in [200, 201, 409]: if retries > 0: if result.status_code == 429: time.sleep(120) return create_playlist_placement_notification( timestamp, playlist_name, playlist_id, playlist_dsp, playlist_storeid, playlist_rank, sound_recording_isrc, tracks_list, retries=retries - 1) else: error = {} try: error = result.json() except json.decoder.JSONDecodeError: pass error['notification'] = body sentry.send_response_to_sentry( response.create_error_response( status=result.status_code, code=errors.ERROR_CODE_OWS_NOTIFICATIONS_REQUEST, message=error), result.status_code ) return False return True @ddtracer.wrap(resource='task.chartmetric_spike_detector', name='create_social_spike_notification') def create_social_spike_notification( date, network, new_followers, chartmetric_id, retries=5): """Send social spike event to ows-notifications for processing. Args: date (str): YYYY-MM-DD date format when event happened network (str): Social network (youtube, instagram) new_followers (int): Increase in followers on network day-over-day chartmetric_id (int): Identifier for a participant retries (int): Number of times to retry failures Returns: bool: success or failure """ if ddtracer.enabled: dd_span = ddtracer.current_span() if dd_span: dd_span.set_tag('network', network) dd_span.set_tag('chartmetric_id', chartmetric_id) body = { 'date': date, 'network': network, 'new_followers': new_followers, 'chartmetric_id': chartmetric_id } result = request.process( application='swf-activity-detector', environment=base_config.ENVIRONMENT, method='post', service_name='ows-notifications', path='/activity/social_spike', json=body) if result.status_code not in [200, 201, 409]: if retries > 0: if result.status_code == 429: time.sleep(120) return create_social_spike_notification( date, network, new_followers, chartmetric_id, retries=retries - 1) else: error = {} try: error = result.json() except json.decoder.JSONDecodeError: pass error['notification'] = body sentry.send_response_to_sentry( response.create_error_response( status=result.status_code, code=errors.ERROR_CODE_OWS_NOTIFICATIONS_REQUEST, message=error), result.status_code ) return False return True def create_notification(notification, retries=5): """Create a notification by calling ows-notifications. Args: notification (dict): the notification data retries (int): the number of retries Returns: response.Response: 201-400-500 """ if ddtracer.enabled: resource_type = f'task.{notification["feed_name"]}' with ddtracer.trace(resource=resource_type, name='create_ows_notification') as dd_span: dd_span.set_tag('feed_name', notification['feed_name']) dd_span.set_tag('feed_id', notification['feed_id']) result = request.process( application='swf-activity-detector', environment=base_config.ENVIRONMENT, method='post', service_name='ows-notifications', path='/activity', json=notification) if result.status_code != 201 and retries > 0: time.sleep(60) return create_notification(notification, retries=retries - 1) if result.status_code != 201: try: error = result.json() error['notification'] = notification except ValueError: error = {'notification': result.raise_for_status()} sentry.send_response_to_sentry( response.create_error_response( status=result.status_code, code=errors.ERROR_CODE_OWS_NOTIFICATIONS_REQUEST, message=error), result.status_code ) return response.Response()