import datetime import json import os import sys import time import boto3 from neo4j import GraphDatabase LOG_FILE = os.path.join(os.path.dirname(__file__), 'fingerprint-logs.txt') LOG_INFO = 'INFO' LOG_ERROR = 'ERROR' ENVIRONMENT = os.environ.get('ENVIRONMENT', 'qa') NEO4J_URL = os.environ.get('NEO4J_URL', '***') NEO4J_USERNAME = os.environ.get('NEO4J_USERNAME', '***') NEO4J_PASSWORD = os.environ.get('NEO4J_PASSWORD', '***') NEO4J_MAX_RETRY_TIME = 30 ORCHARD_ASSET_LIMIT = int(os.environ.get('ORCHARD_ASSET_LIMIT', 100)) ORCHARD_ASSET_OFFSET = str(os.environ.get('ORCHARD_ASSET_OFFSET', '0')) ORCHARD_ASSET_EXT = str(os.environ.get('ORCHARD_ASSET_EXT', 'flac')) GET_ORCHARD_ASSET_WITHOUT_ACR = ''' MATCH (oa:OrchardAsset) WHERE NOT (oa)-[:FINGERPRINTED_AS]->() AND oa.id > $provided_offset AND oa.extension = $provided_ext RETURN oa ORDER BY oa.id LIMIT $provided_limit ''' CHECK_ORCHARD_ASSET_WITH_ACR = ''' MATCH (oa:OrchardAsset)-[:FINGERPRINTED_AS]->(:ACRID) WHERE oa.id = $provided_id RETURN oa ''' def log(log_type, log_msg, final_log=False): with open(LOG_FILE, 'a') as log_file: ts = datetime.datetime.now().isoformat() msg = '[{}] [{}] {}\n'.format(ts, log_type, log_msg) log_file.write(msg) if final_log: log_file.write('----------\n') def _get_driver(): """Connect to Neo4J and return Neo4J driver.""" try: driver = GraphDatabase.driver( NEO4J_URL, auth=(NEO4J_USERNAME, NEO4J_PASSWORD), encrypted=True, max_transaction_retry_time=NEO4J_MAX_RETRY_TIME) return driver except Exception as e: raise e def _check_orchard_asset_processed(neo4j_driver, asset_id): with neo4j_driver.session() as session: nodes = session.run( CHECK_ORCHARD_ASSET_WITH_ACR, provided_id=asset_id ) result = bool(nodes.data()) session.close() return result def _get_orchard_asset_to_process(neo4j_driver, asset_id): """Retrieve OrchardAsset to process.""" with neo4j_driver.session() as session: assets = [] nodes = session.run( GET_ORCHARD_ASSET_WITHOUT_ACR, provided_limit=ORCHARD_ASSET_LIMIT, provided_offset=asset_id, provided_ext=ORCHARD_ASSET_EXT ) results = [record for record in nodes.data()] for result in results: asset = { 'id': result['oa'].get('id'), 'filename': result['oa'].get('filename'), 'extension': result['oa'].get('extension') } assets.append(asset) session.close() return assets try: log(LOG_INFO, 'Starting fingerprint script for {}'.format(ENVIRONMENT)) neo4j_driver = _get_driver() gt_asset_id = ORCHARD_ASSET_OFFSET lambda_client = boto3.client('lambda', region_name='us-east-1') orchard_assets = [] asset_id_to_check = None while True: # get last asset from previous batch if orchard_assets: asset_id_to_check = orchard_assets[-1].get('id') gt_asset_id = asset_id_to_check # retrieve new batch of assets log(LOG_INFO, 'Fetching new batch of assets > {}'.format(gt_asset_id)) orchard_assets = _get_orchard_asset_to_process(neo4j_driver, gt_asset_id) # noqa:E501 if not orchard_assets: log(LOG_INFO, 'No assets to fingerprint') break # verify previous batch is done if asset_id_to_check: while True: if _check_orchard_asset_processed(neo4j_driver, asset_id_to_check): # noqa:E501 break log(LOG_INFO, 'OrchardAsset with id: {} not fingerprinted yet, waiting....'.format(asset_id_to_check)) # noqa:E501 time.sleep(15) # process batch of assets for asset in orchard_assets: payload = json.dumps({'asset': asset}).encode('utf-8') lambda_client.invoke( FunctionName='{}-lambda-acr-asset-fingerprinter'.format(ENVIRONMENT), # noqa:E501 InvocationType='Event', Payload=payload ) log(LOG_INFO, 'OrchardAsset with id: {} sent for fingerprint'.format(asset.get('id'))) # noqa:E501 neo4j_driver.close() log(LOG_INFO, 'Fingerprint script done', True) sys.exit() except Exception as e: log(LOG_ERROR, 'Unexpected error: {}'.format(e), True) sys.exit(e)