import copy import decimal import os import sys import time import uuid import boto3 ISRC = 'isrc' TIMESTAMP = 'timestamp' USER = 'user' BACKFILL_USER = '179' OPCODE = 'opcode' CORRELATION_ID = 'correlation_id' PROD_ENVIRONMENT = 'prod' TERRITORIES = 'territories' TAKEDOWN_OPCODE = 'TAKEDOWN' REMOVE_OPCODE = 'REMOVE' env = os.environ.get('Environment', PROD_ENVIRONMENT) dynamodb_resource = boto3.resource( 'dynamodb', 'us-east-1') active_table = dynamodb_resource.Table( '{}-{}'.format(env, 'masters_active')) audit_table = dynamodb_resource.Table( '{}-{}'.format(env, 'masters_audit_new')) correlation_id = uuid.uuid1() isrc_processed = [] isrc_not_processed = [] def takedown_isrc(item, tuid): """Remove a tuid from DynamoDB registry item. If there is only one tuid for all territories - remove the item from registry. Args: item (dict): Masters Registry record from DynamoDB tuid (int): Track unique ID """ isrc = item['isrc'] unique_tuids = set() deleted_territories = {} updated_territories = {} for terr, tuid_obj in item.get('territories').items(): updated_tuids = [] if type(tuid_obj) is dict: tuid_obj = [tuid_obj] for mr_tuid_obj in tuid_obj: mr_tuid = mr_tuid_obj.get('tuid') unique_tuids.add(mr_tuid) if mr_tuid != tuid: updated_tuids.append(mr_tuid_obj) else: deleted_territories[terr] = tuid if updated_tuids: updated_territories[terr] = updated_tuids # If we have only one unique tuid - remove the whole record completely if len(unique_tuids) == 1 and tuid in unique_tuids: active_table.delete_item(Key={ISRC: isrc}) opcode = TAKEDOWN_OPCODE # Otherwise update the existing record (if we removed any territories) elif len(deleted_territories) > 0: updated_item = copy.deepcopy(item) updated_item['territories'] = updated_territories updated_item['timestamp'] = decimal.Decimal( str(time.time() * 1000)) active_table.put_item(Item=updated_item) opcode = REMOVE_OPCODE # Save audit record if we made any changes if len(deleted_territories) > 0: audit_table.put_item( Item={ ISRC: '{isrc}'.format(isrc=isrc), TIMESTAMP: decimal.Decimal(str(time.time() * 1000)), USER: BACKFILL_USER, OPCODE: opcode, CORRELATION_ID: str(correlation_id), TERRITORIES: deleted_territories}) print('ISRC: {},{} processed.'.format(isrc, tuid)) if __name__ == '__main__': filename = sys.argv[1] with open(filename, 'r') as f: lines = f.readlines() for line in lines: isrc, tuid = line.split(',') isrc = isrc.strip() tuid = int(tuid.strip()) response = active_table.get_item(Key={ISRC: isrc}) item = response.get('Item') if item and item.get('territories'): takedown_isrc(item, tuid) isrc_processed.append((isrc, tuid)) else: isrc_not_processed.append((isrc, tuid)) print('processed: {}'.format(len(isrc_processed))) print(isrc_processed) print('not processed: {}'.format(len(isrc_not_processed))) print(isrc_not_processed)