"""Lambda function module for create_tracks.""" from typing import Dict, List, Optional from common.models.state_machine.grps_ingestion_context import \ GrpsIngestionContext from common.models.state_machine.track import Track from common.schemas.state_machine_schema import StateMachineSchema from lambdacommon.graphql import graphql import config from config import graphql_gateway from src.constants.queries import GET_TRACKS_BY_UPC, SAVE_TRACKS from src.exceptions import CreateTracksException logger = config.app_logger def handler(event, context): """Create tracks handler. Retrieves existing tracks for the product UPC from Orchard. Determines which tracks need to be created and which need to be deleted. Creates new tracks, enriches context with tuids, and updates track sequence numbers. """ logger.info(f'Triggered create_tracks: {event}') sm_context = StateMachineSchema().load(event) correlation_id = sm_context.correlation_id graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, } ) try: art_relations_tracks = get_art_relations_tracks(sm_context.product.upc) product_id = sm_context.product.product_id existing_tuids = [] tuids_to_delete = [] for art_relations_track in art_relations_tracks: for context_track in sm_context.tracks: if art_relations_track['isrc'] == context_track.isrc: existing_tuids.append(int(art_relations_track['tuid'])) context_track.tuid = int(art_relations_track['tuid']) break tracks_to_create = [ t for t in sm_context.tracks if not t.tuid ] for art_relations_track in art_relations_tracks: track_tuid = int(art_relations_track['tuid']) if track_tuid not in existing_tuids: tuids_to_delete.append(track_tuid) delete_tracks(product_id, tuids_to_delete) created_tracks = create_tracks(product_id, tracks_to_create) enrich_track_tuids(sm_context, created_tracks) update_track_sequence_numbers(sm_context) sm_context.deleted_tracks = tuids_to_delete except graphql.GraphQLError as err: message = str(err) logger.error(f'GraphQL error in create_tracks: {message}') raise CreateTracksException(message) from err return StateMachineSchema().dump(sm_context) def get_art_relations_tracks(upc: str) -> List[Dict]: """Retrieve existing tracks in Orchard for a given UPC. Args: upc: Product UPC. Returns: list: List of track dicts with isrc and tuid fields. """ result = graphql_gateway.execute( GET_TRACKS_BY_UPC, {'upc': upc} ) logger.info(f'Tracks found for UPC {upc}: {result}') if result['data']['productByUpc']: return result['data']['productByUpc']['tracks'] else: return [] def delete_tracks(product_id: int, tuids: List[int]) -> None: """Delete tracks with the given tuids from a product. Args: product_id: Product ID. tuids: List of track tuids to delete. """ if not tuids: logger.info('Skipped deleting tracks') return payload = { 'delete': { 'productId': product_id, 'tracks': tuids } } graphql_gateway.execute( SAVE_TRACKS, {'data': payload} ) logger.info(f'Deleted tracks with payload {payload}') def create_tracks( product_id: int, tracks: List[Track]) -> Optional[List[Dict]]: """Create the given tracks for a product. Args: product_id: Product ID. tracks: List of Track objects to create. Returns: list: List of created track dicts with isrc and tuid, or None. """ if not tracks: logger.info('Skipped creating tracks') return None tracks_payload = [] for track in tracks: tracks_payload.append({ 'isrc': track.isrc, 'trackName': track.track_name, 'volumeNumber': track.volume, 'explicit': track.explicit or 'N' }) payload = { 'create': { 'productId': product_id, 'tracks': tracks_payload } } logger.info(f'Creating tracks with payload {payload}') result = graphql_gateway.execute( SAVE_TRACKS, {'data': payload} ) return result['data']['saveTracks'] def enrich_track_tuids( context: GrpsIngestionContext, new_tracks: Optional[List[Dict]]) -> None: """Populate context tracks with tuids returned from track creation. Args: context: State machine context containing tracks. new_tracks: List of track dicts returned after creation. """ if not new_tracks: logger.info('Skipped populating track tuids') return for new_track in new_tracks: context_track = next( ( ct for ct in context.tracks if ct.isrc == new_track['isrc'] ), None ) if context_track: context_track.tuid = int(new_track['tuid']) def update_track_sequence_numbers(context: GrpsIngestionContext) -> None: """Update sequence numbers of tracks in the product. Args: context: State machine context containing tracks and product. """ tracks_payload = [] for context_track in context.tracks: tracks_payload.append({ 'trackNumber': context_track.sequence_number, 'volumeNumber': context_track.volume, 'tuid': context_track.tuid }) payload = { 'updatePosition': { 'productId': context.product.product_id, 'tracks': tracks_payload } } logger.info( f'Updating track sequence numbers with payload {payload}' ) graphql_gateway.execute( SAVE_TRACKS, {'data': payload} )