"""set-main-release-date.""" from datetime import date, datetime from typing import Dict, List, NamedTuple import uuid import config from constants import queries from constants.set_main_release_date import ( ALTAFONTE_CUTOVER_DATE, ALTAFONTE_SOUNDCLOUD_DATE, ALTAFONTE_SPOTIFY_DATE, RC_JSON_NAME) from ddex_ingester_common.constants.ddex_providers import ( ALTAFONTE, SME, SOM_LIVRE_VENDOR_ID) from ddex_ingester_common.constants.release_type import VIDEO_RELEASE_TYPES from ddex_ingester_common.constants.status import ( IN_CONTENT, TRANSFER_TO_CONTENT) from ddex_ingester_common.helpers.s3_ddex import load_ddex_json from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.models.s3.deal import Deal as S3Deal from ddex_ingester_common.models.state_machine.body import ( Body as StateMachineContext) from ddex_ingester_common.release_correction.release_correction_diffs import ( ReleaseCorrectionDiffDetail) from ddex_ingester_common.release_correction.release_corrections_s3 import ( create_blank_rc_json, write_rc_json) from ddex_ingester_common.schemas.s3_schema import DealSchema from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.schemas.state_machine_schema import ( StateMachineSchema ) from lambdacommon.graphql import graphql logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Lambda entry point.""" parsed_ddex = S3Schema().load(load_ddex_json(event)) context = StateMachineSchema().load(event) correlation_id = context.correlation_id or str(uuid.uuid4()) context.correlation_id = correlation_id logging_utils.update_logger_correlation_id(logger, correlation_id) logging_utils.update_logger_with_message_ids( logger, parsed_ddex.message_id, parsed_ddex.message_thread_id, parsed_ddex.execution_name ) graphql_conn = graphql.GraphQLConnector( config.GRAPHQL_GATEWAY_URL, config.APPLICATION_NAME) graphql_conn.set_headers(config.GRAPHQL_HEADERS) # Skip the lambda to avoid removing the timed release if context.product.vendor_id == SOM_LIVRE_VENDOR_ID: if timed_release_exists( context, graphql_conn, check_release_schedule=True ): upc = parsed_ddex.product.upc logger.info(f'Timed release found for product {upc}. Exiting') return StateMachineSchema().dump(context) log_deals = [DealSchema().dump(deal) for deal in parsed_ddex.deals or []] logger.info(f'Received deals: {log_deals}') deals = parsed_ddex.deals release_dates = retrieve_release_dates(deals, context.ddex_provider) sales_start_date = release_dates.get('sales_start_date') territory_dates = release_dates.get('territory_dates') release_date = parsed_ddex.product.original_release_date \ if parsed_ddex.product.original_release_date else sales_start_date if not release_date: raise ValueError('Release date not found.') if context.product.release_type in VIDEO_RELEASE_TYPES: update_video_release_date( graphql_conn, context, release_date, sales_start_date) else: update_audio_release_date( graphql_conn, context, release_date, sales_start_date) context.product.sale_start_date = str_to_time(sales_start_date) # CCM-3561 We no longer set timed release from DDEX deliveries # call update timed release date # update_product_timed_release( # context, # graphql_conn, # context.product.product_id, # release_dates.get('timed_release_date'), # sales_start_date) # call replace territory dates replace_product_territory_dates( graphql_conn, context.product.product_id, territory_dates) return StateMachineSchema().dump(context) def filter_release_deals(deals: List[S3Deal]) -> S3Deal: """Filter deals based on release references and deal terms.""" filtered_deals = [] for reference in ['R0', 'R1']: for deal in deals: if reference in deal.release_references and deal.deal_terms: filtered_deals.append(deal) if filtered_deals: return filtered_deals[0] # Fallback to deals with only deal terms filtered_deals = [ deal for deal in deals if deal.deal_terms] if filtered_deals: return filtered_deals[0] raise ValueError('No valid release deals found.') def retrieve_deal_dates(deal: S3Deal, ddex_provider: str) -> (str, str): """Retrieve earliest sales_start_date and set timed_release_date.""" sales_start_date = None timed_release_date = None for deal_term in deal.deal_terms: start_date = deal_term.start_date start_date_time = deal_term.start_date_time if not start_date and not start_date_time: continue if not start_date: # Strip time from datetime start_date = start_date_time[0:10] # A date time with format '00:00:00.000Z' would indicate a timed # release and '00:00:00.000' would indicate a non-timed release. # Timed release date is the latest deal term start date set_timed_release_date = ( start_date_time.endswith('Z') and ( not timed_release_date or timed_release_date < start_date_time) # noqa ) if set_timed_release_date: timed_release_date = start_date_time # For non SME we want diff timed release and sales start date # So we never set sale start date with the date of a timed release if ddex_provider != SME and start_date_time.endswith('Z'): continue if not sales_start_date or sales_start_date > start_date: sales_start_date = start_date if timed_release_date and not sales_start_date: sales_start_date = timed_release_date[0:10] return sales_start_date, timed_release_date def retrieve_territory_with_differing_dates( filtered_deal: S3Deal, sales_start_date: str) -> Dict: """Retrieve territories when start date != sales start date.""" territory_dates = {} for deal_term in filtered_deal.deal_terms: start_date = deal_term.start_date start_date_time = deal_term.start_date_time if not start_date and not start_date_time: continue if not start_date: # Strip time from datetime start_date = start_date_time[0:10] if start_date == sales_start_date: continue if deal_term.territories: for territory in deal_term.territories: territory_date = territory_dates.get(territory) if not territory_date or territory_date < start_date: territory_dates[territory] = start_date return territory_dates def filter_preorder_deals(deal: S3Deal) -> S3Deal: """Filter out preorder deal terms from release deal.""" # We have a filtered deal, weed out pre_orders deal.deal_terms = [deal_term for deal_term in deal.deal_terms if not deal_term.pre_order] if not deal.deal_terms: raise ValueError('Only pre-order deals exist in R0 or R1.') return deal def retrieve_release_dates(deals: List[S3Deal], ddex_provider: str) -> Dict: """Retrieve the latest start date from deal terms.""" # Search for release deals R0 and R1 first, then fallback to any deal filtered_deal = filter_release_deals(deals) # We have a filtered deal, weed out pre_orders filtered_deal = filter_preorder_deals(filtered_deal) # Find earliest sales start date and set timed release date sales_start_date, timed_release_date = retrieve_deal_dates( filtered_deal, ddex_provider) # Gather territory dates that are not the same as sales start date territory_dates = retrieve_territory_with_differing_dates( filtered_deal, sales_start_date) return { 'sales_start_date': sales_start_date, 'timed_release_date': timed_release_date, 'territory_dates': territory_dates, } def update_audio_release_date( graphql_conn: graphql.GraphQLConnector, context: StateMachineContext, release_date: str, sales_start_date: str): """Update audio product release date and sale start date.""" if context.ddex_provider == ALTAFONTE: CUTOVER_DATE_OBJ = str_to_time(ALTAFONTE_CUTOVER_DATE) if str_to_time(sales_start_date) < CUTOVER_DATE_OBJ: # noqa sales_start_date = ALTAFONTE_CUTOVER_DATE payload = { 'data': { 'productId': context.product.product_id, 'releaseDate': release_date, 'saleStartDate': sales_start_date, } } logger.info(f'Updating product release date with payload: {payload}') graphql_conn.execute( queries.UPDATE_PRODUCT_RELEASE_DATE, payload ) def update_video_release_date( graphql_conn: graphql.GraphQLConnector, context: StateMachineContext, release_date: str, sales_start_date: str): """Update video product release date and sale start date.""" new_release = True original_release_date = None status = context.product.status provider = context.ddex_provider if provider == SME and status not in [IN_CONTENT, TRANSFER_TO_CONTENT]: if sales_start_date is None: sales_start_date = release_date # Videos cannot have release_date in the past, use today if needed if str_to_time(sales_start_date).date() < date.today(): sales_start_date = str(date.today()) # Check if we should consider this a re-release if release_date is not None and release_date != sales_start_date: original_release_date = release_date new_release = False payload = { 'data': { 'update': { 'productId': str(context.product.product_id), 'originalReleaseDate': original_release_date, 'releaseDate': sales_start_date, 'newRelease': new_release } } } logger.info(f'Updating video release date with payload: {payload}') graphql_conn.execute( queries.UPDATE_VIDEO_PRODUCT_RELEASE_DATE, payload ) return if provider == ALTAFONTE: # Non-SME logic decided by SWITCH-2793 scenarios CUTOVER_DATE_OBJ = str_to_time(ALTAFONTE_CUTOVER_DATE) if release_date != sales_start_date: original_release_date = release_date new_release = False if str_to_time(sales_start_date) < CUTOVER_DATE_OBJ: # Scenarios 1-2 original_release_date = release_date release_date = ALTAFONTE_CUTOVER_DATE # Because we're updating the release date here, we also # need to be check if it's a new release date or not. new_release = ( str_to_time(release_date).date() < date.today() ) else: if release_date != sales_start_date: # noqa # Scenario 4 original_release_date = release_date release_date = sales_start_date else: # Scenario 3 original_release_date = None new_release = True is_past_date = ( str_to_time(release_date).date() < date.today() ) if is_past_date and status not in [IN_CONTENT, TRANSFER_TO_CONTENT]: # noqa # Release date in the past, set the release date to today # Make sure we don't change the release date of existing products new_release = False original_release_date = release_date release_date = str(date.today()) logger.info( f'Release date in the past, setting the release date to today.' f'Previous release date: {original_release_date}. ' f'New release date: {release_date}.') if provider == SME and status in [IN_CONTENT, TRANSFER_TO_CONTENT]: update_diff_details = [] graphql_response = check_for_product(context, graphql_conn) current_release_date = graphql_response.get('releaseDate') if release_date and current_release_date != release_date: logger.info(f'Found difference for Video field Release Date' f'Current field values: {current_release_date}' f'Update attempt field values: {release_date}') update_diff_details.append(ReleaseCorrectionDiffDetail( field_name='releaseDate', isrc=None, old=current_release_date, new=release_date)) if update_diff_details: event_details =\ {'bucket': context.bucket, 'key': context.key} s3_release_corrections = {'changes': update_diff_details} create_blank_rc_json(event_details, RC_JSON_NAME) write_rc_json(event_details, s3_release_corrections, RC_JSON_NAME) # noqa return # Skip the update since it will fail with the error: # {'status': 400, 'statusText': 'BAD REQUEST', 'body': {'code': 'invalid_status', 'message': 'Invalid release status'}} # noqa: E501 elif status in [IN_CONTENT, TRANSFER_TO_CONTENT]: logger.info(f'Skip release date update for {status} video') return payload = { 'data': { 'update': { 'productId': str(context.product.product_id), 'originalReleaseDate': original_release_date, 'releaseDate': release_date, 'newRelease': new_release } } } logger.info(f'Updating video release date with payload: {payload}') graphql_conn.execute( queries.UPDATE_VIDEO_PRODUCT_RELEASE_DATE, payload ) def check_for_product( context: StateMachineContext, graphql_conn: graphql.GraphQLConnector) -> Dict: """Check for existence of a product.""" upc = context.product.upc logger.info(f'Running get product with upc: {upc}') payload = { 'upc': upc } result = graphql_conn.execute( queries.GET_PRODUCT_BY_UPC, payload )['data']['productByUpc'] logger.info(f'Get product result: {result}') return result def update_product_timed_release( context, graphql_conn, product_id, timed_release_date, sales_start_date): """Update product timed release date. Args: graphql_conn (graphql.GraphQLConnector): Graphql connector product_id (int): Product ID timed_release_date (datetime.datetime): Timed release date - format YYYY-MM-DDTHH:MM:SSZ. The presence of the last letter Z indicates whether or not it's GMT Timed release date always ends in 'Z', so it's always in UTC time """ # We don't support staggered releases on the DDEX # We only have staggered releases hardcoded for some stores for Altafonte staggered_payload = get_altafonte_staggered_release( context, sales_start_date ) if not timed_release_date: if staggered_payload: payload = { 'data': { 'productId': product_id, 'releaseSchedule': { 'staggered': staggered_payload, 'timed': [], }, } } logger.info(f'Set staggered release with payload: {payload}') graphql_conn.execute(queries.SET_PRODUCT_RELEASE_SCHEDULE, payload) return if context.ddex_provider == ALTAFONTE: CUTOVER_DATE_OBJ = str_to_time(ALTAFONTE_CUTOVER_DATE) TIMED_DATE_TRIM = ( str_to_time(timed_release_date[:-1], frmt='%Y-%m-%dT%H:%M:%S')) if TIMED_DATE_TRIM < CUTOVER_DATE_OBJ: # If we have a Timed Release before Cutover date, disregard it. if staggered_payload: pld = { 'data': { 'productId': product_id, 'releaseSchedule': { 'staggered': staggered_payload, 'timed': [], }, } } logger.info(f'Set staggered release with payload: {pld}') graphql_conn.execute(queries.SET_PRODUCT_RELEASE_SCHEDULE, pld) return # This config value was used before the query, use as reference if needed timed_release_store_ids = config.TIMED_RELEASE_STORE_IDS timed_release_store_ids = [ item['id'] for item in graphql_conn.execute( queries.GET_TIMED_RELEASE_STORES, {} )['data']['deliveryStoresV2']['items'] ] payload_store_ids = [ { 'deliveryStore': {'id': store_id}, 'saleDateTime': timed_release_date } for store_id in timed_release_store_ids ] payload = { 'data': { 'productId': product_id, 'releaseSchedule': { 'staggered': staggered_payload, 'timed': payload_store_ids, }, } } logger.info(f'Set product timed release with payload: {payload}') graphql_conn.execute(queries.SET_PRODUCT_RELEASE_SCHEDULE, payload) def timed_release_exists(context, graphql_conn, check_release_schedule=False): """Check if a timed release exists for this product.""" product = check_for_product(context, graphql_conn) timed_release_exists = ( product.get('timedRelease') and product.get('timedRelease', {}).get('timeOfDayProduct') ) # Check the new timed release implementation release_schedule_exists = ( product.get('releaseSchedule') and product.get('releaseSchedule', {}).get('timed') ) if check_release_schedule: return release_schedule_exists or timed_release_exists return timed_release_exists def replace_product_territory_dates( graphql_conn, product_id, territory_dates): """Update product territory date. Args: graphql_conn (graphql.GraphQLConnector): Graphql connector product_id (int): Product ID territory_dates (dict): Territory dates """ payload = { 'productId': product_id, 'territoryDates': [] } data = {} for country_code, start_date in territory_dates.items(): if not data.get(start_date): data[start_date] = [] data[start_date].append(country_code) for territory_date, country_codes in data.items(): payload['territoryDates'].append( { 'countryCodes': sorted(country_codes), 'saleStartDate': territory_date, } ) logger.info(f'Replacing product territory dates with payload: {payload}') graphql_conn.execute( queries.REPLACE_PRODUCT_TERRITORY_DATES, {'data': payload} ) def get_altafonte_staggered_release(context, sales_start_date): """Get Spotify and Soundcloud staggered release dates for Altafonte.""" if context.ddex_provider != ALTAFONTE: return [] staggered_payload = [] if str_to_time(sales_start_date) < str_to_time(ALTAFONTE_SPOTIFY_DATE): staggered_payload.append({ 'deliveryStore': {'id': '286'}, 'saleDate': ALTAFONTE_SPOTIFY_DATE }) if str_to_time(sales_start_date) < str_to_time(ALTAFONTE_SOUNDCLOUD_DATE) \ and context.product.release_type not in VIDEO_RELEASE_TYPES: staggered_payload.append({ 'deliveryStore': {'id': '723'}, 'saleDate': ALTAFONTE_SOUNDCLOUD_DATE }) return staggered_payload def str_to_time(date, frmt='%Y-%m-%d'): """Convert a string to a datetime object.""" return datetime.strptime(date, frmt) class ReleaseCorrectionValuePair(NamedTuple): """Contains the current and incoming values for a DDEX field.""" new: any old: any