"""Configuration and marshaling of DelphiFeed statuses.""" from collections import defaultdict import datetime import time from flask_restful import fields from flask_restful import marshal_with from flask_restful import Resource from feed_status import config from feed_status.models import delphi_feed from feed_status.models import feed_monitor from feed_status.util import generate_default_date_dict_range delphi_feed_status_history_fields = { 'data': fields.Nested({ 'feed_id': fields.String(attribute='feed_id'), 'feed_name': fields.String(attribute='feed_name'), 'metrics': fields.Raw, 'delivery_date': fields.String(attribute='delivery_date'), 'status': fields.Raw, 'status_details': fields.Raw, 'last_modified_utc': fields.DateTime( attribute='last_modified_utc', dt_format='iso8601')}), 'dates': fields.List(fields.String) } class DelphiFeedStatus(Resource): """Provides response for /delphi_feed//status.""" @marshal_with(delphi_feed_status_history_fields) def get(self, feed_name, date=None): """Marshal specific delphiFeed status. Args: feed_name (str): specific feed of interest. Returns: dict: json payload for feed status. """ days_count = config.DASHBOARD_DISPLAYED_DAYS if date: to_date = datetime.date.fromisoformat(date) else: to_date = datetime.date.today() delta = datetime.timedelta(days=days_count) from_date = to_date - delta active_feeds = delphi_feed.get_active_feeds(feed_name) # make request to DB start = time.time() db_feeds = delphi_feed.get_all_delphi_feed_history(from_date, to_date) end = time.time() print(f'Feed Execution time: {end - start}') # {feed_id: {date: [feeds]}} feeds_map = defaultdict(dict) for db_feed in db_feeds: feed_id = ('delphi-' f'{db_feed["data_source_name"]}-' f'{db_feed["licensor_name"]}-' f'{db_feed["report_name"]}') report_date = db_feed['date'] feeds_map[feed_id].setdefault(report_date, []).append(db_feed) for active_feed in active_feeds: self._process_activity_feed( active_feed, feeds_map, from_date, to_date) dates = generate_default_date_dict_range(from_date, to_date) dates = list(dates.keys()) return { 'data': active_feeds, 'dates': dates, 'metrics': [] } def _process_activity_feed( self, active_feed, feeds_map, from_date, to_date): db_feeds = feeds_map.get(active_feed.feed_id) status_map = generate_default_date_dict_range(from_date, to_date) status_details = {d: dict() for d in status_map.keys()} if db_feeds: # iterate on days for report_date, db_feed_list in db_feeds.items(): report_date_str = report_date.strftime('%Y-%m-%d') # ignore cancelled feeds = [f for f in db_feed_list if f['status'] != 'CANCELLED'] if not feeds: status_map[report_date_str] = 'cancelled' continue # if there are errors in any context, # then mark uow as failed status = feeds[0]['status'] context_statuses = feeds[0]['context_statuses'] if status == 'ACTIVE' and context_statuses: has_errors = ( any(map(lambda s: s['status'] in ('ON_HOLD', 'FAILED'), context_statuses))) if has_errors: status = 'FAILED' # Handle multiple uows per day, e.g. for AE. # If all units are active, icon = red X, # if at least one unit is completed, icon is = orange V, # if all units are completed, icon = green V if len(feeds) > 1: all_completed = all(map( lambda f: f['status'] == 'COMPLETE', feeds)) all_active = all(map( lambda f: f['status'] == 'ACTIVE', feeds)) if all_completed: status = 'COMPLETE' elif all_active: status = 'ACTIVE' else: status = 'ingested_warn' complete_status = status if status in ('COMPLETE', 'MIN_COMPLETE'): status = delphi_feed.INGESTION_STATUS_INGESTED status_map[report_date_str] = status.lower() if len(feeds) == 1 and feeds[0]['is_force_complete']: complete_status = 'FORCE_COMPLETE' # complete, min_complete, force_complete last_updated = (feeds[0]['last_updated_at'] .strftime('%Y-%m-%d %H:%M:%S')) status_details[report_date_str] = { 'complete_status': complete_status.lower(), 'updated': last_updated } else: status_details = None active_feed.status = status_map active_feed.status_details = status_details active_feed.metrics = feed_monitor.get_metric_status( active_feed, active_feed.status)