import datetime import logging import os import boto3 from boto3.dynamodb.types import TypeSerializer from tadas.platform import config STATUS_INGESTED = 'INGESTED' STATUS_NOT_INGESTED = 'NOT_INGESTED' STATUS_NOT_AVAILABLE = 'NOT_AVAILABLE' STATUS_DOWNLOADED = 'DOWNLOADED' ALL_STATUSES = { STATUS_INGESTED, STATUS_NOT_INGESTED, STATUS_NOT_AVAILABLE, STATUS_DOWNLOADED, } logger = logging.getLogger(__name__) serializer = TypeSerializer() def set_overall_status(feed_name, date, overall_status, details=None): """Set overall ingestion status for a (feed_name, date) row in DynamoDB.""" if overall_status not in ALL_STATUSES: raise ValueError( f'Unsupported overall status: {overall_status}. ' f'Available statuses: {ALL_STATUSES}' ) TADAS_UPDATE_DYNAMODB_STATUS = config.get('TADAS_UPDATE_DYNAMODB_STATUS') if not TADAS_UPDATE_DYNAMODB_STATUS: logger.info(f'Skipping DynamoDB status update {TADAS_UPDATE_DYNAMODB_STATUS=}') return logger.info(f'Setting overall status for {feed_name=} {date=} {overall_status=} {details=}') attributes = {} now = datetime.datetime.now(datetime.UTC).replace(microsecond=0) attributes['updated_at'] = now.isoformat() if details: attributes['details'] = details update_dynamodb_status( feed_name=feed_name, date=date, status=overall_status, attributes=attributes, dynamodb_table=os.environ['FEED_INGESTION_TABLE'], aws_region=os.environ.get('AWS_REGION', 'us-east-1'), ) def update_dynamodb_status( feed_name, date, status, dynamodb_table, aws_region, attributes=None): if not date: raise ValueError('no date') # validating date format. Will raise ValueError if invalid. datetime.datetime.strptime(date, '%Y-%m-%d') dynamodb = boto3.client('dynamodb', region_name=aws_region) key = { 'feed_name': {'S': feed_name}, 'date': {'S': date}, } attributes = attributes or {} attributes['status'] = status expression_attribute_names = {f'#{k}': k for k in attributes.keys()} expression_attribute_values = { f':{k}': serializer.serialize(v) for k, v in attributes.items() } update_expression = ', '.join(f'#{k} = :{k}' for k in attributes.keys()) update_expression = f'SET {update_expression}' dynamodb.update_item( TableName=dynamodb_table, Key=key, UpdateExpression=update_expression, ExpressionAttributeNames=expression_attribute_names, ExpressionAttributeValues=expression_attribute_values, )