"""Logical operations for orchard sound recordings.""" import re from typing import Optional from connector_neo4j import Neo4jSession from sound_recordings import config from sound_recordings.models import orchard_sound_recordings from sound_recordings.models import ows_track from sound_recordings.utils import convert from sound_recordings.utils import s3 as s3_util class MultipleAcrOsrRelationships(Exception): """OrchardAsset found with > 1 ACR or OSR relationship.""" pass class AssetNotFound(Exception): """Unable to find OrchardAsset with Track id.""" pass class TrackISRCNotSet(Exception): """Track does not have ISRC property.""" pass class OrchardSoundRecordingNotFound(Exception): """OSR not found.""" pass class OrchardSoundRecordingWithoutAssets(Exception): """OSR without assets.""" pass class PrimaryTrackNotRelated(Exception): """Primary Track not related to the OSR.""" pass @Neo4jSession(transaction=True, use_v2=True, database=config.NEO4J_DATABASE_NAME) def upsert(track_id): """Upsert orchard sound recording based on track id. Args: track_id (int): track id to match Returns: str: uuid of orchard sound recording bool: any updates made """ track_info = orchard_sound_recordings._get_track_info(track_id) if not track_info: raise AssetNotFound() _validate_track_data(track_info) osr_ids = track_info.get('osr_ids') if len(osr_ids) != 0: return (osr_ids[0], False) # claim new ISRC is Track.isrc exists as OrchardSoundRecording.isrc osr_isrc = track_info['isrc'] if not orchard_sound_recordings._isrc_available(osr_isrc): osr_isrc = ows_track.claim_isrc() ( sr_id, nodes_created, rels_created ) = orchard_sound_recordings._create_osr(osr_isrc, track_id, track_info) return ( sr_id, (nodes_created or rels_created) ) @Neo4jSession(transaction=True, use_v2=True, database=config.NEO4J_DATABASE_NAME) def update(osr_id, primary_track_id): """Update orchard sound recording based on sr id. Args: osr_id (UUID): orchard sound recording id primary_track_id (int): primary track id to update Returns: dict: orchard sound recording """ fetched_osr = orchard_sound_recordings.fetch(orchard_sound_recording_ids=[osr_id]) if not fetched_osr: raise OrchardSoundRecordingNotFound() osr = fetched_osr[str(osr_id)] assets = osr.get('assets') if not assets: raise OrchardSoundRecordingWithoutAssets() existing_track_ids = [] for asset in assets: existing_track_ids += assets[asset]['tuids'] if primary_track_id not in existing_track_ids: raise PrimaryTrackNotRelated() updated_osr = orchard_sound_recordings.update(osr_id, primary_track_id) if not updated_osr: return updated_osr return convert.to_snake_case_dict(updated_osr.data().get('soundRecording')) def _validate_track_data(track_data): """Validate Track / OrchardSoundRecording data. Args: track_data (dict): track data """ if not track_data.get('isrc'): raise TrackISRCNotSet() if len(track_data.get('assets')) == 0 or len(track_data.get('acr_ids')) == 0: raise AssetNotFound() if len(track_data.get('acr_ids')) > 1 or len(track_data.get('osr_ids')) > 1: raise MultipleAcrOsrRelationships() def fetch(orchard_sound_recording_ids=[], track_ids=[], upcs=[], product_ids=[], project_ids=[], asset_ids=[], track_isrcs=[], include_deleted=False, include_transfer_to_content=False, term='', include_inactive=False): # noqa:E501 """Fetch orchard sound recording based on ids, set default tuid if needed. If term is provided, it will override orchard_sound_recording_ids param Args: orchard_sound_recording_ids (list): UUIDs to filter for orchard sound recordings track_ids (list): ints to filter by track ids upcs (list): ints to filter by product upcs product_ids (list): ints to filter by product ids project_ids (list): ints to filter by project ids assets_ids (list): UUIDs to filter by orchard assets track_isrcs (list): strings to filter by track isrcs include_deleted (Bool): flag to include or not soft deleted asset include_transfer_to_content (Bool): flag to include products in this status term (str): term to search, cannot be used with orchard_sound_recording_ids include_inactive (Bool): flag to include inactive tracks Returns: list: orchard sound recordings w/ asset data """ osr_ids = [] if term: cleaned_term = convert.escape_term(term) cleaned_term = cleaned_term.strip() found_osr_ids = orchard_sound_recordings.search_by_isrc(cleaned_term) if include_inactive: found_osr_ids += orchard_sound_recordings.search_by_isrc(cleaned_term, inactive=True) # noqa:E501 if not found_osr_ids: return found_osr_ids osr_ids = found_osr_ids else: osr_ids = orchard_sound_recording_ids sound_recording_map = orchard_sound_recordings.fetch( orchard_sound_recording_ids=osr_ids, track_ids=track_ids, upcs=upcs, product_ids=product_ids, project_ids=project_ids, asset_ids=asset_ids, track_isrcs=track_isrcs, include_deleted=include_deleted, include_transfer_to_content=include_transfer_to_content, include_inactive=include_inactive ) # process results, business logic and flatten from dict to list results = [] for sr_id, sr in sound_recording_map.items(): # merge all tuids into single set tuids = set([ tuid for _, asset in sr['assets'].items() for tuid in asset['tuids'] ]) # choose default tuid if not set explictly or tuid not part of connected Track if tuids: if not sr['primary_track_id'] or sr['primary_track_id'] not in tuids: sr['primary_track_id'] = min(tuids) else: sr['primary_track_id'] = None # flatten assets from dict to list sr['assets'] = [ asset for _, asset in sr['assets'].items() ] # flatten tracks from dict to list for asset in sr['assets']: asset['tracks'] = [ track for _, track in asset['tracks'].items() ] results.append(sr) return results def fetch_product(product_id, track_isrcs, include_transfer_to_content=False): """ Fetch orchard sound recording matches based on product ID or track ISRCs. Args: product_id (int): Product ID track_isrcs (list): strings to filter by track isrcs include_transfer_to_content (bool): flag to include products in this status Returns: list: asset data """ sound_recording_map = orchard_sound_recordings.fetch_product( product_id=product_id, track_isrcs=track_isrcs, include_transfer_to_content=include_transfer_to_content ) # process results, business logic and flatten from dict to list results = [] for sr_id, sr in sound_recording_map.items(): # merge all tuids into single set tuids = set([ tuid for _, asset in sr['assets'].items() for tuid in asset['tuids'] ]) # choose default tuid if not set explictly or tuid not part of connected Track if tuids: if not sr['primary_track_id'] or sr['primary_track_id'] not in tuids: sr['primary_track_id'] = min(tuids) else: sr['primary_track_id'] = None # flatten assets from dict to list sr['assets'] = [ asset for _, asset in sr['assets'].items() ] # flatten tracks from dict to list for asset in sr['assets']: asset['tracks'] = [ track for _, track in asset['tracks'].items() ] results.append(sr) return results def fetch_version(osr_id, version_id): """Fetch orchard sound recording based on id and version. Args: osr_id (str): The Orchard Sound Recording ID version_id (str): The version ID for the Orchard Sound Recording Returns: dict: Orchard sound recording data for a specific version """ return s3_util.get_sound_recording_version(osr_id, version_id) def fetch_full_delivery_history( osr_ids, execution_type=[], service=[], event_type=[], ack_status=[], limit=None, offset=None, start_date=None ): """Fetch orchard sound recording full delivery history based on id. Args: osr_ids (list): The Orchard Sound Recording IDs execution_type (list): List of execution types to filter by service (list): List of UGC services to filter by event_type (list): List of event types to filter by ack_status (list): List of ack statuses to filter by (e.g. ['success', 'awaited', 'error', 'missing']) Returns: [dict]: Orchard sound recording full delivery history data """ delivery_history = orchard_sound_recordings.fetch_full_delivery_history( osr_ids, execution_type, service, event_type, ack_status, limit, offset, start_date) if not delivery_history: return delivery_history for delivery in delivery_history: delivery_message = delivery['message'] delivery_sfn_id = delivery['sfn_execution_id'] if delivery['event_type'] != 'not_eligible': delivery['delivery_xml_signed_url'] = _get_delivery_file_signed_url( delivery_sfn_id, delivery_message, delivery['service']) else: delivery['delivery_xml_signed_url'] = None if delivery['service'] != 'TikTok (Audio Fingerprinting)': delivery['ack'] = 'missing' delivery['ack_message'] = None return delivery_history def _get_delivery_file_signed_url(delivery_sfn_id, message, service: Optional[str] = None): # noqa:E501 """Get delivery history event xml file link. Args: delivery_sfn_id (str): sfn execution id message (dict): event message service (Optional[str]): UGC service name Returns: str: signed url of the delivery file """ if message['details'] is None: return None if message['status'] == 'error': return None delivery_files = message['details']['filenames'] r = re.compile('^((?!BatchComplete).)*\\.xml$') xml_path = None for path in delivery_files: if r.match(path): xml_path = path break if not xml_path: return xml_path youtube_path_prefix = '' if service and service.strip().lower().startswith('youtube'): youtube_path_prefix = 'youtube/' normalized_xml_path = xml_path.lstrip('/') path = f'{youtube_path_prefix}{delivery_sfn_id}/{normalized_xml_path}' return s3_util.get_signed_url(path) @Neo4jSession(transaction=True, use_v2=True, database=config.NEO4J_DATABASE_NAME) def touch(resource_id, resource_type, modified_before, limit): """Update orchard sound recording based on connected resource and date. Args: resource_id (int): Integer unique ID of resource resource_type (str): Object type of resource modified_before (Datetime): Datetime to filter limit (int): maximum number of records to touch Returns: dict: number of updated osr """ touched_osr = orchard_sound_recordings.touch_osr( resource_id, resource_type, modified_before, limit) return { 'nodes_updated': touched_osr }