"""Sql2sf feed (covers all the syncing tables).""" from collections import defaultdict from datetime import datetime from itertools import groupby from feed_status.models.orm import sql2sf_status STATUS_COMPLETED = 'ingested' TS_FORMAT = '%Y-%m-%dT%H:%M:%S' def _parse_ts(time_str): """Iterate over acceptable formats and try to parse string. Args: time_str (str): target string to parse. Returns: datetime: result datetime object. If the format is not recognized default value is returned datetime.min. """ acceptable_formats = [TS_FORMAT, '%Y-%m-%d'] for fmt in acceptable_formats: try: return datetime.strptime(time_str, fmt) except Exception: pass return datetime.min def get_table_statuses(target_tables): """Fetch statuses for tables listed in config. Args: target_tables list(tuple): list of (schema_name, table_name). Returns: dict: tables grouped by schema_name and sorted by name. """ target_tables_set = set(target_tables) status_items = filter( lambda e: (e['schema_name'], e['table_name']) in target_tables_set, sql2sf_status.get_all_statuses() ) mappers = defaultdict( lambda: (lambda x: x), # default last_sync_timestamp=_parse_ts, synced_to_timestamp=_parse_ts, status=lambda v: v.lower() if v else 'unknown' ) # group by schema_name, sort by table_name result = [] for schema_name, tables in groupby( status_items, lambda e: e['schema_name']): # strip schema_name field and convert types tables_sanitized = [ dict( (k, mappers[k](v)) for k, v in t.items() if k != 'schema_name') for t in tables ] item = { 'schema_name': schema_name, 'tables': sorted( tables_sanitized, key=lambda e: e['table_name']) } result.append(item) return result