""" Task status management functions via DynamoDB. Similar to garcon_feed_status.py but used to check individual task statuses to see if they can be safely skipped. This uses the same DynamoDB items created by the feed_status functions, but uses a different attribute. There are no pre-set constants of task statuses, and these functions are not considered a replacement for feed_status functions. """ import datetime import logging import boto3 from garcon_contrib.dynamo_feed_status import config from garcon_contrib.dynamo_feed_status import garcon_feed_status logger = logging.getLogger(__name__) FIELD_NEW_CONTEXTS = 'new_contexts' FIELD_COMPLETED = 'completed' FIELD_CONTEXT_STATUS = 'context_status' FIELD_CONTEXT_REQUIRED = 'required' CONTEXT_STATUS_NOT_AVAILABLE = 'NOT_AVAILABLE' CONTEXT_STATUS_IN_PROGRESS = 'IN_PROGRESS' CONTEXT_STATUS_PROCESSED = 'PROCESSED' aws_session = boto3.Session() def get_completed_tasks(feed_name, datestamp): """Return a list of completed tasks for execution lookup params. This does not have different behavior if the task does not exist in DynamoDB. Consider a non existing task execution as simply not having any completed tasks instead of an exception. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. Returns: list: List of task names, empty list if no tasks are completed. """ status_item = garcon_feed_status._get_item(feed_name, datestamp) if not status_item: return [] # this is called with _status at the end because the # garcon_feed_status.set_status function appends it automatically. tasks = status_item.get('completed_tasks_status', '').split(',') return tasks if any(tasks) else [] def is_completed_task(feed_name, datestamp, task_name): """Check DynamoDB for task completion status. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. task_name (str): Name of task. Returns: bool: True if task_name is found in marker attribute of DynamoDB item. """ return task_name in get_completed_tasks(feed_name, datestamp) def mark_completed_task(feed_name, datestamp, task_name): """Update DynamoDB item with task completion status note. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. task_name (str): Name of task. """ completed_tasks = get_completed_tasks(feed_name, datestamp) if task_name not in completed_tasks: completed_tasks.append(task_name) garcon_feed_status.set_status( feed_name, datestamp, 'completed_tasks', ','.join(completed_tasks)) def get_values(feed_name, datestamp, field): """Return a list of values by field. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. field (str): Filed name in DynamoDB. Returns: list: List of task names, empty list if no tasks are completed. """ status_item = garcon_feed_status._get_item(feed_name, datestamp) if not status_item: return [] result = status_item.get(field, '').split(',') return result if any(result) else [] def set_values(feed_name, datestamp, field, values): """Update DynamoDB item filed with values. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. field (str): Filed name in DynamoDB. values (list): List of values which need to save. """ garcon_feed_status.set_status( feed_name, datestamp, field, ','.join(values)) def get_item(feed_name, datestamp): """Return a DynamoDB Item. Args: feed_name (str): Feed name of workflow execution for status updates. datestamp (str): Workflow generated date. Returns: Dict: An item if found, empty dict if no item is found. """ status_item = garcon_feed_status._get_item(feed_name, datestamp) if not status_item: return {} return status_item def _get_ingestion_table(): dynamodb = aws_session.resource('dynamodb', region_name='us-east-1') return dynamodb.Table(config.feed_ingestion_table) def is_completed_overall_job(feed_name, date): """Check if overall job is completed. Args: feed_name (str): Feed name of workflow execution for status updates. date (str): Workflow generated date. Returns: Bool: completed status. Default is True. """ key = { 'feed_name': feed_name, 'date': date } item_dict = _get_ingestion_table().get_item(Key=key) item = item_dict.get('Item', {}) if not item: return True else: status = item.get(FIELD_COMPLETED, True) return status def mark_completed_overall_job(feed_name, date, value): """Mark overall job as completed. Args: feed_name (str): Feed name of workflow execution for status updates. date (str): Workflow generated date. value (bool): flag if job is completed or not. """ attr_name = FIELD_COMPLETED key = { 'feed_name': feed_name, 'date': date } update_expression = 'SET #{} = :{}'.format(attr_name, attr_name) expression_attribute_names = { '#{}'.format(attr_name): '{}'.format(attr_name) } expression_attribute_values = { ':{}'.format(attr_name): value } _get_ingestion_table().update_item( Key=key, UpdateExpression=update_expression, ExpressionAttributeNames=expression_attribute_names, ExpressionAttributeValues=expression_attribute_values ) def is_completed_report(feed_name, date): """Check if report is completed. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. Returns: bool: True if all contexts for report are in completed state. """ contexts = get_report_contexts(feed_name, date) if contexts: return all(val[FIELD_CONTEXT_STATUS] == CONTEXT_STATUS_PROCESSED for val in contexts.values()) else: return True def create_report_contexts(feed_name, date, contexts, contexts_config, optional_config): """Create contexts for report. For each report set initial state, required status and other attributes. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. contexts (list): List of contexts to create. contexts_config (list): Configuration for required/optional contexts. optional_config (bool): configuration for required/optional contexts. """ key = { 'feed_name': feed_name, 'date': date } contexts_value = {} for context in contexts: if optional_config: is_required = not (context in contexts_config) else: is_required = context in contexts_config contexts_value[context] = { FIELD_CONTEXT_STATUS: CONTEXT_STATUS_NOT_AVAILABLE, FIELD_CONTEXT_REQUIRED: is_required } update_expression = 'SET #contexts = :contexts' expression_attribute_names = { '#contexts': 'contexts' } expression_attribute_values = { ':contexts': contexts_value } _get_ingestion_table().update_item( Key=key, UpdateExpression=update_expression, ExpressionAttributeNames=expression_attribute_names, ExpressionAttributeValues=expression_attribute_values ) def get_report_contexts(feed_name, date): """Get contexts for report. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. Returns: dict of contexts. Example: { 'US': { 'context_status': 'NOT_AVAILABLE', 'required': true }, 'PL': { ... } } """ key = { 'feed_name': feed_name, 'date': date } item_dict = _get_ingestion_table().get_item(Key=key) if item_dict: item = item_dict.get('Item', {}) contexts = item.get('contexts', {}) return contexts return {} def update_report_context_status(feed_name, date, context, status): """Update status for context. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. context (str): context name. status (str): context status, for instance, PROCESSED. """ key = { 'feed_name': feed_name, 'date': date } update_expression = 'SET contexts.#sub.context_status = :value' expression_attribute_names = { '#sub': context } expression_attribute_values = { ':value': status } _get_ingestion_table().update_item( Key=key, UpdateExpression=update_expression, ExpressionAttributeNames=expression_attribute_names, ExpressionAttributeValues=expression_attribute_values ) def mark_report_processed_contexts(feed_name, date): """Change contexts from in_progress to processed. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. Returns: list of processed contexts. """ res = [] # Try to find better way for bulk update. for context in get_report_in_progress_contexts(feed_name, date): update_report_context_status( feed_name, date, context, CONTEXT_STATUS_PROCESSED) res.append(context) return res def get_report_in_progress_contexts(feed_name, date): """Get contexts in 'in progress' state. Args: feed_name (str): Feed name of workflow execution. date (str): Workflow generated date. Returns: list of contexts. """ contexts = get_report_contexts(feed_name, date) return [context for context, val in contexts.items() if val[FIELD_CONTEXT_STATUS] == CONTEXT_STATUS_IN_PROGRESS] def add_newcontext(feed_name, date, value): """Add context to set of new_contexts. Args: feed_name (str): Feed name of workflow execution for status updates. date (str): Workflow generated date. value (str): context, e.g. vendor id for AppleMusic. """ key = { 'feed_name': feed_name, 'date': date } update_expression = 'ADD {} :elements'.format(FIELD_NEW_CONTEXTS) expression_attribute_values = {':elements': {value}} _get_ingestion_table().update_item( Key=key, UpdateExpression=update_expression, ExpressionAttributeValues=expression_attribute_values ) def delete_newcontexts(feed_name, date): """Remove new_contexts attribute for item. Args: feed_name (str): Feed name of workflow execution for status updates. date (str): Workflow generated date. """ key = { 'feed_name': feed_name, 'date': date } update_expression = 'REMOVE {}'.format(FIELD_NEW_CONTEXTS) _get_ingestion_table().update_item( Key=key, UpdateExpression=update_expression ) def get_newcontexts(feed_name, date): """Get new_contexts attribute for item. Args: feed_name (str): Feed name of workflow execution for status updates. date (str): Workflow generated date. Returns: set: Set of new contexts. """ key = { 'feed_name': feed_name, 'date': date } item_dict = _get_ingestion_table().get_item(Key=key) item = item_dict.get('Item', {}) if not item: return set() else: contexts = item.get(FIELD_NEW_CONTEXTS, set()) return contexts def soft_reload_update_status_clearing_tasks( context_date: str, feed_name: str, overall_status=garcon_feed_status.STATUS_POPULATED_RAW_TABLE, clear_completed_tasks=( 'create_staging_fact', 'load_staging_fact', 'load_fact_data', 'load_aggregated_skips_and_saves', 'set_overall_status_INGESTED', ), ): """ Prepare feed_status for soft reload. Soft reload is a feature to rerun all transformations for the runs which have already been processed. This method sets overall_status and clears tasks which need to be restarted. """ assert feed_name assert datetime.datetime.strptime(context_date, '%Y-%m-%d') assert overall_status item = garcon_feed_status._get_item(feed_name=feed_name, date=context_date) if not item: logging.info(f'No feed status for {feed_name=} {context_date=}') return completed_tasks = item.get('completed_tasks_status', '').split(',') tasks_to_clear_set = set(clear_completed_tasks) filtered_completed_tasks = [ task for task in completed_tasks if task not in tasks_to_clear_set ] attributes = { 'completed_tasks_status': ','.join(filtered_completed_tasks), } logger.info( f'Soft reload: updating status for {feed_name=} {context_date=} ' f'to {overall_status=} and {attributes}') garcon_feed_status.set_overall_status( feed_name=feed_name, date=context_date, overall_status=overall_status, attributes=attributes ) def soft_reload_update_status_clearing_after( context_date: str, feed_name: str, overall_status=garcon_feed_status.STATUS_POPULATED_RAW_TABLE, clear_after_task='set_overall_status_POPULATED_RAW_TABLE', ): """ Prepare feed_status for soft reload. Soft reload is a feature to rerun all transformations for the runs which have already been processed. This method sets overall_status to POPULATED_RAW_TABLE and clears completion status for tasks after specified task. """ assert feed_name assert datetime.datetime.strptime(context_date, '%Y-%m-%d') assert overall_status item = garcon_feed_status._get_item(feed_name=feed_name, date=context_date) if not item: logger.info(f'No feed status for {feed_name=} {context_date=}') return completed_tasks = item.get('completed_tasks_status', '').split(',') if clear_after_task not in completed_tasks: logger.warning( f'Soft reload: did not {clear_after_task=} in {completed_tasks=}') filtered_completed_tasks = completed_tasks else: last_entry_index = completed_tasks.index(clear_after_task) + 1 filtered_completed_tasks = completed_tasks[:last_entry_index] attributes = { 'completed_tasks_status': ','.join(filtered_completed_tasks), } logger.info( f'Soft reload: updating status for {feed_name=} {context_date=} ' f'to {overall_status=} and {attributes}') garcon_feed_status.set_overall_status( feed_name=feed_name, date=context_date, overall_status=overall_status, attributes=attributes )