"""Logic for retrieving ugc video metrics.""" from typing import Any, Mapping from ddtrace import tracer from oto import response as oto_response from analytics.constants import cache, countries, store from analytics.constants.ordering import ORDER_DIRECTIONS from analytics.logic.parallel import parallel from analytics.logic.stores import add_outage_error_to_stores from analytics.queries.format import format_row from analytics.queries.ugc_video_metrics import UgcVideoMetrics, UgcVideoMetricsCount from analytics.schemas.ugc_video_metrics import UgcVideoMetricsSchema from analytics.utils import store_availability from analytics.utils.cache import cache_in_redis UGC_VIDEO_METRICS_FIELDS = [ "video_id", "asset_id", "channel_id", "track_id", "track_isrc", "views", "premium_views", "percentage", "claim_type", "claim_created_date", "claim_policy_monetize", "claim_policy_track", "claim_policy_block", "claim_status", "content_type", "video_name", "channel_name", "release_date", ] TASK_RESULT = "get_ugc_video_metrics" TASK_RESULT_COUNT = "get_ugc_video_metrics_total_results" @cache_in_redis(ttl=cache.ONE_DAY) @tracer.wrap(name="get_ugc_video_metrics") def get_ugc_video_metrics( query_params: Mapping[str, Any], permissions: Mapping[str, Any], ): """Fetch UGC video metrics. Args: query_params: Dict with distributors, store_ids, track_ids, asset_ids, video_ids, track_isrcs, countries, claim_types, claim_statuses, order_by, order_dir, limit, offset. permissions: Dict with permission_* keys. Returns: oto_response.Response with UGC video metrics payload. """ order_by = query_params["order_by"] order_dir = query_params["order_dir"] if order_by not in UGC_VIDEO_METRICS_FIELDS: raise Exception("Invalid order_by field") if order_dir.upper() not in ORDER_DIRECTIONS: raise Exception("Invalid order_dir value") store_ids = store_availability.get_query_store_ids( query_params.get("store_ids") or [], store.VIDEO_STORE ) if not store_ids: return oto_response.Response([]) claim_statuses = query_params.get("claim_statuses") or [] common_input = { **permissions, "distributors": query_params["distributors"], "country_ids": query_params.get("countries") or [], "track_ids": query_params.get("track_ids") or [], "asset_ids": query_params.get("asset_ids") or [], "video_ids": query_params.get("video_ids") or [], "track_isrcs": query_params.get("track_isrcs") or [], "claim_types": query_params.get("claim_types") or [], "transfer_product_ownership_enabled": query_params.get( "transfer_product_ownership_enabled", False ), # V_METRICS_UGC_BY_VIDEO_ASSET_*_ROLLUP is widened by the # transfer-product-ownership dbt PR; rollup_grain dedups former-owner # rows that would otherwise inflate views/premium_views post-transfer. "rollup_grain": True, } metrics_input = { **common_input, "claim_statuses": claim_statuses, "claim_statuses_blocked": "blocked" in claim_statuses, "order_by": order_by, "order_dir": order_dir, "limit": query_params["limit"], "offset": query_params["offset"], } requests = { TASK_RESULT: {"func": _fetch_metrics, "args": (metrics_input,)}, TASK_RESULT_COUNT: {"func": _fetch_count, "args": (common_input,)}, } result = parallel(requests).message if not result or not result[TASK_RESULT]: return oto_response.Response({"items": []}) return UgcVideoMetricsSchema.normalized_response( { "items": result[TASK_RESULT], "total_results": result[TASK_RESULT_COUNT], "sources": add_outage_error_to_stores( store_availability.get_video_sources() ), } ) def _fetch_metrics(query_input): rows = [format_row(row) for row in UgcVideoMetrics(query_input).execute()] for video in rows: video["claim_policy"] = _parse_claim_policy(video) return rows def _fetch_count(query_input): rows = [format_row(row) for row in UgcVideoMetricsCount(query_input).execute()] if not rows: return 0 return rows[0].get("total_results", 0) or 0 def _parse_claim(claim): if claim is None: return "" if claim.startswith("INCLUDE"): included = claim[8:].strip().split(" ") return ",".join(included) if claim.startswith("EXCLUDE"): excluded = claim[8:].strip().split(" ") result = [c for c in countries.ISO_ALPHA_2 if c not in excluded] result.sort() return ",".join(result) return "" def _parse_claim_policy(video): policy = "|".join( [ _parse_claim(video["claim_policy_monetize"]), _parse_claim(video["claim_policy_track"]), _parse_claim(video["claim_policy_block"]), ] ) del video["claim_policy_monetize"] del video["claim_policy_block"] del video["claim_policy_track"] return policy if len(policy) > 2 else None