"""Catalog ingestion helper.""" import csv from datetime import datetime import io from typing import List from typing import NamedTuple from typing import Union from uuid import uuid4 import boto3 """ Note: If defaults other than `None` are used for fields in NamedTuple classes (such as a datetime for timestamp) then that default value may get cached between lambda invocations resulting in undefined behaviour. For information around how lambda runtimes work: https://docs.aws.amazon.com/lambda/latest/dg/runtimes-context.html """ class CatalogIngestion(NamedTuple): """Represents a catalog ingestion.""" table_name = 'catalog_ingestion' column_order = [ 'state_machine_name', 'state_machine_execution_name', 'catalog_ingestion_source_id', 'timestamp', 's3_bucket_name', 's3_key_name', 'ingest_format' ] state_machine_name: str state_machine_execution_name: str catalog_ingestion_source_id: int s3_bucket_name: str s3_key_name: str ingest_format: str timestamp: str class CatalogIngestionValidationResult(NamedTuple): """Represents a catalog ingestion validation result.""" table_name = 'catalog_ingestion_validation_result' column_order = [ 'state_machine_name', 'state_machine_execution_name', 'validation_rule_id', 'response', 'message' ] state_machine_name: str state_machine_execution_name: str validation_rule_id: int = None response: str = None message: str = None class CatalogIngestionAction(NamedTuple): """Represents a catalog ingestion action.""" table_name = 'catalog_ingestion_action' column_order = [ 'state_machine_name', 'state_machine_execution_name', 'timestamp', 'action', 'entity_type', 'project_code', 'project_id', 'project_name', 'upc', 'release_id', 'release_name', 'vendor_catalog_number', 'isrc', 'tuid', 'track_sequence_number', 'track_volume_number', 'track_name', 'result', 'message' ] state_machine_name: str state_machine_execution_name: str action: str entity_type: str result: str timestamp: str project_code: str = None project_id: int = None project_name: str = None upc: str = None release_id: int = None release_name: str = None vendor_catalog_number: str = None isrc: str = None tuid: int = None track_sequence_number: int = None track_volume_number: int = None track_name: str = None message: str = None class CatalogIngestionSession: """Class representing a catalog ingestion session. Example: Adding a single model to this session:: session = CatalogIngestionSession( 'bucket', 'foo/bar/snowflake/') validation_result = CatalogIngestionValidationResult( state_machine_name='foo', state_machine_execution_name='bar') session.add(validation_result) Adding multiple models to this session:: session = CatalogIngestionSession( 'bucket', 'foo/bar/snowflake/') results = [ CatalogIngestionValidationResult( state_machine_name='foo', state_machine_execution_name='bar' ), CatalogIngestionValidationResult( state_machine_name='foo', state_machine_execution_name='bar' ) ] session.add(results) """ def __init__(self, snowflake_s3_bucket: str, snowflake_s3_path: str): """Construct a catalog ingestion session.""" self.snowflake_s3_bucket = snowflake_s3_bucket self.snowflake_s3_path = snowflake_s3_path self.data = { 'catalog_ingestion': [], 'catalog_ingestion_action': [], 'catalog_ingestion_validation_result': [] } def add(self, models: Union[List[NamedTuple], NamedTuple]): """Add a constructed catalog ingestion model to the session.""" if not isinstance(models, list): models = [models] for model in models: self.data.get(model.table_name).append(model) def save(self): """Save the data in this session as TSVs to S3.""" for models in self.data.values(): if models: stream = io.StringIO() writer = csv.DictWriter( stream, models[0].column_order, delimiter='\t' ) for model in models: writer.writerow(model._asdict()) s3_client = boto3.client('s3') path = ( f'{self.snowflake_s3_path}/{models[0].table_name}' f'/{uuid4()}' ) s3_client.put_object( Bucket=self.snowflake_s3_bucket, Key=path, Body=stream.getvalue() ) def log_catalog_action( bucket, location, asset, action, asset_type, result, msg): """Log an attempt to process artwork.""" # Catalog ingestion db session = CatalogIngestionSession(bucket, location) if asset.asset_type == 'audio': isrc = asset.product.track.isrc tuid = asset.product.track.tuid track_sequence_number = asset.product.track.track_sequence_number track_volume_number = asset.product.track.track_volume_number track_name = asset.product.track.track_name else: isrc = None tuid = None track_sequence_number = None track_volume_number = None track_name = None action_model = CatalogIngestionAction( state_machine_execution_name=asset.execution_name, state_machine_name=asset.state_machine_name, action=action, entity_type=asset_type, result=result, upc=asset.product.upc, project_id=asset.product.project_id, project_code=asset.product.project_code, project_name=asset.product.project_name, release_id=asset.product.product_id, release_name=asset.product.release_name, isrc=isrc, tuid=tuid, track_sequence_number=track_sequence_number, track_volume_number=track_volume_number, track_name=track_name, message=msg, timestamp=datetime.utcnow().isoformat() + 'Z' ) session.add(action_model) session.save() def log_catalog_session(bucket, location, file_key, ingest_format, execution_name, state_machine_name, source_id): """Log a catalog ingestion session.""" # Catalog ingestion db session = CatalogIngestionSession(bucket, location) ingestion_model = CatalogIngestion( state_machine_execution_name=execution_name, state_machine_name=state_machine_name, catalog_ingestion_source_id=source_id, s3_bucket_name=bucket, s3_key_name=file_key, ingest_format=ingest_format, timestamp=datetime.utcnow().isoformat() + 'Z' ) session.add(ingestion_model) session.save()