"""Lambda sr-delivery-tiktok function module.""" import json import os from datetime import datetime from datetime import timezone import random import re from typing import Optional import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from src.connectors import s3_asset_delivery from . import constants from .common import logger from .common.connectors import s3_sound_recordings from .common.connectors import s3_assets from .common.connectors import s3_delivery_audit from .utils import graphql from soundrecording_utils.constants.ddex.constants import RuleService from soundrecording_utils.ddex.generate import generate_ddex from soundrecording_utils.metadata.typedload import get_loader from soundrecording_utils.metadata.types import OrchardSoundRecording import config # initialize sentry sentry_dsn_value = config.secrets_manager_client.get_cred('SENTRY_DSN') sentry_dsn: Optional[str] = os.environ.get( 'SENTRY_DSN', sentry_dsn_value if isinstance(sentry_dsn_value, str) else None) if sentry_dsn: logger.info('Initializing with sentry') sentry_sdk.init( sentry_dsn, integrations=[AwsLambdaIntegration()] ) else: logger.info('Initializing without sentry') class NoAssets(Exception): """No deliverable assets on an OrchardSoundRecording.""" pass def get_osr(sound_recording_id: str, version_id: Optional[str] = None) -> OrchardSoundRecording: """Get the metadata blob from the event message.""" # retrieve sound recording from s3 if version_id: (raw_data, _) = s3_sound_recordings.get_sound_recording_version( sound_recording_id, version_id ) else: (_, raw_data, _) = s3_sound_recordings.get_latest_sound_recording_version( sound_recording_id ) loader = get_loader() if raw_data is None: raise ValueError(f'Sound recording {sound_recording_id} not found') osr = loader.load(json.loads(raw_data), OrchardSoundRecording) return osr def _generate_batch_id() -> str: date_str = datetime.now().strftime('%Y%m%d%H%M%S') random_str = str(random.randint(0, 9999)).rjust(4, '0') return date_str + random_str def _s3_remote_foldernames(batch_id, isrc): resources_dirname = config.RESOURCES_DIRNAME.rstrip('/') return ( batch_id, batch_id + '/' + isrc, batch_id + '/' + isrc + '/' + resources_dirname, ) def _remove_always_different(ddex_str: str) -> str: pattern = r'.*?' ddex_str = re.sub(pattern, '', ddex_str, flags=re.DOTALL) pattern = r'.*?' ddex_str = re.sub(pattern, '', ddex_str, flags=re.DOTALL) return ddex_str def handler(event, context): """Lambda entry point.""" try: process_asset = config.PROCESS_ASSET perform_final_ddex_comparison = config.PERFORM_FINAL_DDEX_COMPARISON takedown: bool = event.get('delivery_type', '') == 'TAKEDOWN_DELIVERY' upload = event['upload_asset'] assert isinstance(upload, bool) sound_recording_id: str = event['sound_recording']['id'] assert isinstance(sound_recording_id, str) version_id: str = event['sound_recording'].get('version') assert isinstance(version_id, str) or version_id is None osr: OrchardSoundRecording = get_osr( sound_recording_id, version_id ) batch_id = _generate_batch_id() execution_metadata = { 'timestamp': datetime.now(timezone.utc), 'message_thread_id': 1, 'message_id': batch_id, } # select asset by metadata if upload required upload_asset = None if upload: if not osr.assets: raise NoAssets() upload_asset = osr.assets[0] # generate XML based on metadata and asset ddex: str = generate_ddex(osr, upload_asset, execution_metadata, config.APPLICATION_NAME, RuleService.TIKTOK, selected_version=config.DDEX_FILE_VERSION, takedown=takedown) # fetch previous DDEX file if final DDEX comparison enabled if perform_final_ddex_comparison: prev_ddex: Optional[str] = None sr_delivery_histories = graphql.get_delivery_histories( [sound_recording_id], constants.TIK_TOK_STORENAME ) # do duplicate check only if the previous delivery was successful if (sound_recording_id in sr_delivery_histories and # noqa:W504 sr_delivery_histories[sound_recording_id].get('status', '').lower() == 'success'): last_delivered_xml_location: Optional[str] = graphql.get_last_delivered_xml_location( # noqa:E501 sound_recording_id ) if last_delivered_xml_location: xml = s3_delivery_audit.download_sr_delivery_audit_xml( last_delivered_xml_location ) # turn XML into string if xml: prev_ddex = xml.decode('utf-8') if prev_ddex: # Remove elements that will always be different before comparison comparable_ddex = _remove_always_different(ddex) comparable_prev_ddex = _remove_always_different(prev_ddex) if comparable_ddex == comparable_prev_ddex: logger.warning('Redundant DDEX! SRID: %s, version: %s', sound_recording_id, version_id) # noqa:E501 return { 'details': { 'undelivered': True, 'reason': 'duplicate' } } # fetch actual audio asset asset_bytes = s3_assets.download_asset( upload_asset.filename, upload_asset.extension ) if upload_asset and process_asset else None # list files to transfer sr_isrc = osr.isrc (batch_dir, isrc_dir, resources_dir) = _s3_remote_foldernames(batch_id, sr_isrc) transfer_files = [ ( asset_bytes, f'{resources_dir}/{upload_asset.filename}.{upload_asset.extension}' ) ] if upload_asset and process_asset else [] xml_files = [ ( bytes(ddex.encode('utf-8')), f'{isrc_dir}/{sr_isrc}.xml' ), ( b'', f'{batch_dir}/BatchComplete_{batch_id}.xml' ) ] # upload XML file to S3 for future auditing s3_delivery_audit.write_sr_delivery_audit_xml( xml_files[0][0], xml_files[0][1], event['execution_name'] ) # Write DDEX XML to S3 s3_asset_delivery.write_file( xml_files[0][0], xml_files[0][1], event['execution_name'] ) # Write actual asset to S3 if upload_asset and process_asset: if not transfer_files or transfer_files[0][0] is None: raise NoAssets() s3_asset_delivery.write_file( transfer_files[0][0], transfer_files[0][1], event['execution_name'] ) # Write "Batch Complete" XML to S3 s3_asset_delivery.write_file( xml_files[1][0], xml_files[1][1], event['execution_name'] ) transfer_files += xml_files return { 'details': { 'batch_id': batch_id, 'filenames': [x[1] for x in transfer_files] } } # suppress entry, fail lambda function except NoAssets as e: sentry_sdk.init() raise e except Exception as e: logger.exception(str(e)) raise e