"""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