"""Logic code.""" import json import os import subprocess from tempfile import TemporaryDirectory from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from lambdacommon.common_config import logger from src import config from src.utils import download_from_s3_url from src.utils import FinalAsset from src.utils import KafkaResultEvent BINARY_PATH = os.path.join( os.path.dirname(__file__), config.EXTRACTOR_MUSIC_BINARY ) string_serializer = StringSerializer() producer = EventProducer( bootstrap_servers=config.KAFKA_BROKERS, key_serializer=string_serializer, value_serializer=string_serializer, security_protocol='SSL') def download_and_process(bucket, key): """Download file from S3 and process it through the library.""" with TemporaryDirectory() as dir: source_file_path = download_from_s3_url( bucket=bucket, key=key, dir=dir, ) size = os.path.getsize(source_file_path) logger.info(f'Downloaded file {source_file_path}, size {size} bytes') output_file_path = os.path.join(dir, 'result.json') command = [BINARY_PATH, source_file_path, output_file_path] logger.info(f'Going to run command: {command}') completed_process = subprocess.run( command, capture_output=True, timeout=config.SUBPROCESS_TIMEOUT) logger.info(f'Result: {completed_process}') with open(output_file_path) as file: return json.load(file) def produce_kafka_event(asset_final: FinalAsset, lowlevel_data: dict): """Produce kafka result event.""" with producer: event_value = KafkaResultEvent( final_asset_id=asset_final.asset_final_id, asset_upload_id=asset_final.asset_upload_id, bucket=asset_final.bucket, filename=asset_final.filename, lowlevel_data=lowlevel_data, ) producer.produce( topic=config.KAFKA_RESULT_TOPIC, event_key=f'asset-final-{asset_final.asset_final_id}', event_value=event_value.json(), auto_flush=False)