"""Tasks for handling overall feed status in DynamoDB.""" from garcon import task from garcon.contrib.dynamo_feed_status import \ feed_status_ingestion as feed_status from analytics_aggregation.util import common as common_utils @task.decorate(timeout=1000) def set_overall_status(activity, feed_name, date_range, status): """Explicitly set the overall status of a feed. Given a date_range, set status for each date in range. Args: activity (ActivityWorker): The activity worker. date_range (str): Job date range in YYYY-MM-DD_YYYY-MM-DD format. feed_name (str): Name of the feed. status (str): Status constant in util.feed_status. """ activity.logger.info( 'Setting status for feed: {feed_name} ' 'with date_range: {date_range} ' 'to: {status}, via a task'.format( feed_name=feed_name, date_range=date_range, status=status)) start_date, end_date = date_range.split('_') for day in common_utils.list_of_dates(start_date, end_date): feed_status.set_overall_status(feed_name, day, status) activity.logger.info( 'Set status for: {}'.format(day)) @task.decorate(timeout=100) def check_workflow_status(activity, feed_name, date_range_as_str, reload): """Check if workflow already successfully executed for a given date range. date_range_as_str is in the applicable case (which is normal non-reingesting run for one day) will be constisting from the one day, e.g. 2017-01-01_2017-01-01 or a date range when reprocessing. Args: activity (ActivityWorker): The activity worker. feed_name (str): The workflow name. date_range_as_str (str): String in YYYY-MM-DD_YYYY-MM-DD format. reload (bool): Flag to force reload and ignore all the statuses. """ if date_range_as_str and not reload: activity.logger.info( 'Checking if workflow was already successfully executed for this ' 'date range: {}'.format(date_range_as_str)) start_date, end_date = date_range_as_str.split('_') range_ingested = True for day in common_utils.list_of_dates(start_date, end_date): range_ingested = feed_status.get_overall_status( feed_name, day) == feed_status.STATUS_INGESTED if not range_ingested: activity.logger.info( 'No successful previous run found for {day}, proceeding ' 'to run range: {date_range}'.format( day=day, date_range=date_range_as_str)) break if range_ingested: return common_utils.exit_message( 'Data was already aggregated for {} range'.format( date_range_as_str))