"""release_approval tasks.""" import asyncio import datetime from ddtrace import tracer as ddtracer from garcon import task from activity_detector.connectors.mysql import ar_db_session from activity_detector.flows.release_approval import config from activity_detector.utils import ows_notifications GET_INGESTED_RELEASES = """ SELECT r.release_id, r.upc, r.display_upc, r.release_name, ai.name AS artist_name, v.vendor_id, r.subaccount_id FROM releases r JOIN artist_info ai ON ai.artist_id = r.artist_id JOIN vendor v ON v.vendor_id = ai.vendor_id WHERE r.release_status = 'in_content' AND DATE_FORMAT(r.ingestion_completed, '%Y-%m-%d') = '{day}' AND v.owner not in ( 'awal', 'AWAL Core Accounts', 'AWAL-Core', 'AWAL-US', 'AWAL-UK', 'AWALrev-NY', 'AWALrev-UK', 'AWALrev-DE', 'HIFI-AWAL', 'AWL-ORC-US', 'AWL-ORC-UK' ) AND v.company_brand_id != 5 AND r.not_for_distribution = 'N' """ @task.decorate(timeout=3600) @ddtracer.wrap(resource='task.release_approval', name='detect_ingested_releases') def detect_ingested_releases(activity): """Detect ingested releases. Args: activity (ActivityWorker): The activity worker. Returns: None """ loop = asyncio.new_event_loop() loop.run_until_complete( create_notifications( group_releases( get_ingested_releases()))) loop.close() @ddtracer.wrap(resource='task.release_approval', name='get_ingested_releases') def get_ingested_releases(): """Get ingested releases from art_relations. Args: None Returns: [dict]: The ingested releases. """ today = datetime.datetime.utcnow().date() yesterday = today - datetime.timedelta(days=1) query = GET_INGESTED_RELEASES.format(day=yesterday.strftime('%Y-%m-%d')) with ar_db_session() as session: rows = session.execute(query).fetchall() if ddtracer.enabled: with ddtracer.current_span() as dd_span: dd_span.set_tag('releases_count', len(rows)) result = [ { 'release_id': row['release_id'], 'upc': row['upc'], 'display_upc': row['display_upc'], 'release_name': row['release_name'], 'artist_name': row['artist_name'], 'vendor_id': row['vendor_id'], 'subaccount_id': row['subaccount_id'] } for row in rows ] return result @ddtracer.wrap(resource='task.release_approval', name='group_releases') def group_releases(releases): """Group releases by label. Args: releases ([dict]): The releases to group. Returns: dict: The releases grouped by label. """ releases_by_label = {} for release in releases: vendor_id = release['vendor_id'] subaccount_id = release['subaccount_id'] if subaccount_id: key = 'subaccount_{}'.format(subaccount_id) else: key = 'vendor_{}'.format(vendor_id) if key in releases_by_label: releases_by_label[key].append(release) else: releases_by_label[key] = [release] return releases_by_label def create_batched_notifications(notifications): """ Send a batch of notifications asynchronously. Args: notifications (list): List of notification objects to be sent. Returns: None """ for notification in notifications: ows_notifications.create_notification(notification) @ddtracer.wrap(resource='task.label_release_approval', name='create_notifications') async def create_notifications(releases_by_label): """Create a notification for each label by calling ows-notifications. Args: releases_by_label (dict): The releases grouped by label. Returns: None """ notifications = [] for key, releases in releases_by_label.items(): notification = { 'feed_name': config.FEED_NAME, 'feed_id': key, 'payload': { 'actor': 'The Orchard Activity Detector', 'verb': 'Detected', 'object': 'Approved Releases', 'approved_releases': releases } } notifications.append(notification) if ddtracer.enabled: with ddtracer.current_span() as dd_span: dd_span.set_tag('notifications_count', len(notifications)) while notifications: notifications_batch = notifications[:15] # Take the first 15 notifications as a batch notifications = notifications[15:] # and remove the sent notifications from the list create_batched_notifications(notifications_batch) await asyncio.sleep(1) # Sleep for 1 second between each batch