import inspect from collections import defaultdict from typing import Callable, Dict, List, Optional, Tuple from server import config from server.constants.delphi.cleaning import DataCleaningMode from server.constants.delphi.skipping import FIELDS_SKIPPING_MAP, ITEM_SKIPPING_MAP, SkippingRule from server.constants.delphi.streams.include import STREAMS_EXTRA_FIELDS_FILTERS, STREAMS_EXTRA_FIELDS_KEYS from server.utils.args.getter_setter import get_arg from server.utils.skippers import skip_extra_data, skip_item, skip_nullable_data def merge_cleaning_rules( rules: List[Dict[str, Dict[str, Optional[Tuple[str]]]]] ) -> Dict[str, Dict[str, Optional[Tuple[str]]]]: """Combine cleaning rules by dsp and keys. Returns combined result.""" merged = defaultdict(dict) for rule in rules: for dsp, dsp_rule in rule.items(): for field, keys in dsp_rule.items(): actual = merged[dsp].get(field, set()) if actual is None: continue if keys is None: merged[dsp][field] = None else: merged[dsp][field] = set(keys) | actual return merged def clean( mode: DataCleaningMode = DataCleaningMode(config.DELPHI_STREAMS_CLEANING_MODE), response_items_node: str or None = "items", include_arg_name: str = "include", dsp_key: str = "dsp", item_skipping_rules: Tuple[SkippingRule] or None = ITEM_SKIPPING_MAP, fields_skipping_rules: Tuple[SkippingRule] or None = FIELDS_SKIPPING_MAP, extra_keys: Dict[str, Tuple[str]] = STREAMS_EXTRA_FIELDS_KEYS, extra_fields_filters: Dict[str, Dict[str, Dict[str, Optional[Tuple[str]]]]] = STREAMS_EXTRA_FIELDS_FILTERS, ): """Clean Delphi streams data, remove redundant data. Args: mode: Mode of cleaning, response_items_node: key for getting the result from response dictionary. include_arg_name: name of wrapped function parameter containing 'include' request options. dsp_key: key for getting dsp value from the result item. item_skipping_rules: rules to remove the whole item from the response. fields_skipping_rules: rules to remove particular fields from the response item. extra_keys: map of extra (controlled by 'include') data keys in the result item by dsp. extra_fields_filters: map of filters by include types and dsp for extra keys in the result item. Returns: Cleaned result. """ def inner(f: Callable): if not mode: return f args_spec = inspect.getfullargspec(f) async def wrapped(*args, **kwargs) -> List: result = await f(*args, **kwargs) items = result.get(response_items_node, []) if response_items_node else result filter_rules = {} if mode == DataCleaningMode.FULL: includes = get_arg(args, kwargs, args_spec, include_arg_name, from_default=True) or [] filter_rules = merge_cleaning_rules([extra_fields_filters[include] for include in includes]) for i in range(len(items) - 1, -1, -1): if skip_item(items, i, item_skipping_rules): continue item = items[i] skip_nullable_data(item, fields_skipping_rules) if mode == DataCleaningMode.FULL: skip_extra_data(item, filter_rules, extra_keys, dsp_key) return result wrapped.__signature__ = inspect.signature(f) return wrapped return inner def clean_video_analytics( mode: DataCleaningMode = DataCleaningMode(config.DELPHI_VIDEO_ANALYTICS_CLEANING_MODE), response_items_node: str or None = "items", inner_node: str or None = "metrics", include_arg_name: str = "only", ): """Clean Delphi video analytics data. Differs from the /streams 'clean' decorator: 1. Clean all None/null fields without mapping. 2. Drop data fields that can be null depending on the 'only' argument, not ID fields as for /streams. 3. 'Only' is used to understand if we need to clean, if it is not set (full data, the largest size) than exit. 4. Need to work with inner objects, different data structure. So, as Delphi creates each new endpoint in a completely different way, but for now we have only two streams/views endpoints, it is not clear how we can merge these two decorators. TODO: Let's return to that question when we would have at least 3. Args: mode: Cleaning mode. response_items_node: Key to get result list from response dictionary. inner_node: Inner data node, path to clean. include_arg_name: Result fields filter argument name. Returns: Cleaned result. """ def inner(f: Callable): if not mode: return f args_spec = inspect.getfullargspec(f) async def wrapped(*args, **kwargs) -> List: result = await f(*args, **kwargs) includes = get_arg(args, kwargs, args_spec, include_arg_name, from_default=True) if not includes: return result items = result.get(response_items_node, []) if response_items_node else result for item in items: if inner_node: item = item.get(inner_node, {}) for key in list(item.keys()): if item[key] is None: del item[key] return result wrapped.__signature__ = inspect.signature(f) return wrapped return inner