"""Logic for retrieving video metrics.""" import json from typing import Any, Mapping from ddtrace import tracer from analytics.constants import cache from analytics.handler_utils import user_has_full_access from analytics.logic import data_availability from analytics.logic.permissions import has_full_access_on_permissions from analytics.logic.stores import add_outage_error_to_stores from analytics.queries.format import format_row from analytics.queries.videos import AllTimeVideoMetrics, VideoMetrics from analytics.schemas.video_metrics import AllTimeVideoMetricsSchema from analytics.utils import store_availability from analytics.utils.cache import cache_in_redis def lowercase_keys(d): if isinstance(d, dict): return {k.lower(): lowercase_keys(v) for k, v in d.items()} elif isinstance(d, list): return [lowercase_keys(i) for i in d] else: return d @tracer.wrap(name="get_video_metrics_bulk") def get_video_metrics_bulk(query_params, permissions): """Return metrics for video_ids. Args: query_params (dict): Query parameters permissions (dict): Permissions filter """ # remove everything except of label_ids, since only vendors has access to videos if not has_full_access_on_permissions(permissions) and permissions.get("label_ids"): permissions["artist_ids"] = [] permissions["subaccount_ids"] = [] permissions["label_participant_ids"] = [] elif not has_full_access_on_permissions(permissions) and not permissions.get( "permission_label_ids" ): permissions["permission_artist_ids"] = [] permissions["permission_subaccount_ids"] = [] permissions["permission_label_participant_ids"] = [] permissions["permission_label_ids"] = [ 999999999999999999999999999 ] # dummy label_id value # we need to add a dummy value # because there is a _has_full_access check downstream sources = add_outage_error_to_stores(store_availability.get_video_sources()) empty_item = { "aggregate_video_metrics": { "video_id": None, "ad_fill_rate": 0, "download_activity_date": "1900-01-01", "average_view_duration_percentage": 0, "average_view_duration_seconds": 0, "comments": 0, "cpm": 0, "estimated_monetized_playbacks": 0, "estimated_partner_ad_revenue": 0, "estimated_partner_premium_revenue": 0, "impressions_ctr": 0, "premium_views": 0, "premium_watch_time_minutes": 0, "rpm": 0, "views": 0, "watch_time_minutes": 0, "content_type": "", }, "video_metrics": [], "sources": sources, } if query_params["start_date"] and query_params["days"]: start_date, end_date = data_availability.get_date_range( query_params["start_date"], int(query_params["days"]) ) else: start_date, end_date = data_availability.get_date_range( data_availability.HIGHWATERMARK_DATE, days=28, downloads=False, videos=True ) query_params["start_date"] = start_date.strftime("%Y-%m-%d") query_params["end_date"] = end_date.strftime("%Y-%m-%d") if query_params.get("country_ids"): query_params["table_name"] = "V_VIEWS_BY_VIDEO_COUNTRY_FEED_DISTRIBUTOR_DAILY" else: query_params["table_name"] = "VIEWS_BY_VIDEO_FEED_DISTRIBUTOR_DAILY" query = VideoMetrics({**query_params, **permissions}) result = [format_row(data_point) for data_point in query.execute()] result = [ { json.loads(i["result"])["AGGREGATE"]["VIDEO_ID"]: { k: v for k, v in json.loads(i["result"]).items() } } for i in result ] response = {} for video_id in query_params["video_ids"]: empty_item["aggregate_video_metrics"]["video_id"] = video_id response[video_id] = lowercase_keys( { "aggregate_video_metrics": empty_item["aggregate_video_metrics"], "video_metrics": empty_item["video_metrics"], "sources": empty_item["sources"], } ) for item in result: for video_id in query_params["video_ids"]: if video_id in item: timeseries = [ i for i in item[video_id]["TIMESERIES"] if i["VIDEO_ID"] == video_id ] response[video_id] = lowercase_keys( { "aggregate_video_metrics": item[video_id]["AGGREGATE"], "video_metrics": sorted( timeseries, key=lambda x: x["DOWNLOAD_ACTIVITY_DATE"], ), "sources": sources, } ) return response @cache_in_redis(ttl=cache.ONE_DAY) def get_all_time_video_metrics( query_params: Mapping[str, Any], permissions: Mapping[str, Any], ): """Return all time metrics for video_id.""" response_body = { "video_id": query_params["video_id"], "all_time_video_metrics": {}, "sources": add_outage_error_to_stores(store_availability.get_video_sources()), } # videos are visible only to employees and vendors if ( not user_has_full_access(permissions) and not permissions["permission_label_ids"] ): return AllTimeVideoMetricsSchema.normalized_response(response_body) if not query_params["store_ids"]: query_params["store_ids"] = store_availability.get_video_store_ids() else: query_params["store_ids"] = list( set(query_params["store_ids"]).intersection( store_availability.get_video_store_ids() ) ) if not query_params["store_ids"]: return AllTimeVideoMetricsSchema.normalized_response(response_body) if query_params.get("country_ids"): query_params["table_name"] = "VIEWS_BY_VIDEO_COUNTRY_FEED_DISTRIBUTOR_ROLLUP" else: query_params["table_name"] = "VIEWS_BY_VIDEO_FEED_DISTRIBUTOR_ROLLUP" result = AllTimeVideoMetrics({**query_params, **permissions}).execute().fetchone() if result: response_body["all_time_video_metrics"] = format_row(result) return AllTimeVideoMetricsSchema.normalized_response(response_body)