""" Prepare Asset Search Logic ========================= This logic takes the label audit object from the label queue, gets the releases for that label and adds them to the asset search queue """ from labelaudit.tasks import asset_search_tasks from labelaudit.logic import audit_release from labelaudit.logic import audit_persistence def handle_releases(audit_id): """Iterates through the paginated releases and calls queue handler methods Args: audit_id (int): ID of the audit record for which to set releases """ audit_releases = audit_release.set_audit_releases(audit_id) if audit_releases.errors: raise Exception(audit_releases.errors) num_qualified_releases = audit_releases.message # get the first page of releases to get pagination info releases = audit_release.fetch_audit_releases(audit_id, 0) if releases.errors: raise Exception(releases.errors) # Update the audit report status to in_progress audit_persistence.update_report_status(audit_id, 'in_progress') releases_first_page = releases.message # get data for pagination limit = releases_first_page['pagination']['limit'] total_records = releases_first_page['pagination']['totalRecords'] add_releases_to_queue(releases_first_page['youtubeAuditReleases']) if num_qualified_releases != total_records: raise Exception('Number of releases set does not equal number fetched') offset = limit while offset < total_records: # fetch collection of qualified youtube audit releases for label releases_by_page_response = audit_release.fetch_audit_releases( audit_id, offset) if releases_by_page_response.errors: raise Exception(releases_by_page_response.errors) # add each page of youtube audit releases to the asset search queue releases_by_page = releases_by_page_response.message add_releases_to_queue( releases_by_page['youtubeAuditReleases']) offset += limit def add_releases_to_queue(releases): """Sends releases to asset search queue in batches by page Args: releases (list): A list of qualified releases to add to the queue """ # add to the asset search queue via celery instead of using a producer asset_search_tasks.process_releases(releases)