from dataclasses import dataclass from datetime import timedelta from typing import Dict, Iterable, List, Tuple from google.cloud.bigtable.column_family import MaxVersionsGCRule from google.cloud.bigtable.row import DirectRow from google.cloud.bigtable.row_data import PartialRowData, PartialRowsData from google.cloud.bigtable.row_filters import RowFilterUnion from google.cloud.bigtable.row_set import RowSet from google.protobuf.pyext._message import Message from delphi_api.bigtable.attributes import BigTableAttribute from delphi_api.bigtable.service import BigTableService from delphi_api.const import ENVIRONMENT from delphi_api.errors import BigTableError, Codes from delphi_api.utils import DateUtils from delphi_api.v2.constants import COLUMN_FAMILY_ID, SEP, UTF8 from delphi_api.v2.enums import ColumnName @dataclass class CountriesStruct: countries: Dict[str, Message] class BigTableModel: """Abstract base class for BigTable models. Children should set ``Meta.table_name``""" class Meta: table_name = None country_codes = [] # DSPs can have "sub-dsps" like amazon (prime, unlimited, ads). set in implementing models len_sub_dsps = 1 def __init__(self): self.attribute_values = {} @property def service(self) -> BigTableService: return BigTableService(self.Meta.table_name) def _check_table_exists(self): """To be called before attempting queries. BigTable requires admin client for ``.exists()`` Raises: :class:`BigTableError`: Exception if table does not exist """ if not self.service.admin_table.exists(): raise BigTableError({ 'code': Codes.database_table_missing.value, 'description': 'BigTable database table does not exist.' }) @staticmethod def get_row_key(keys: Iterable[str]) -> bytes: """Concatenate ``keys`` with the key separator :const:`delphi_api.v2.constants.SEP` """ return bytes(f'{SEP}'.join(keys), UTF8) @classmethod def from_raw_data(cls, data: PartialRowData): """Deserializes Bigtable row result data to model Args: data (:class:`PartialRowData`): Bigtable row from query results Returns: :class:`BigTableModel`: an instance of this (or child) class from a single row of data """ if data is None: # pragma: no cover raise ValueError('No data received to construct object') klass = cls() for cls_prop, attribute in cls.__dict__.items(): if not isinstance(attribute, BigTableAttribute): continue try: value = data.cell_value(attribute.column_family, bytes(attribute.attr_name, UTF8)) except KeyError: # No data for this column, skip it. continue else: setattr(klass, cls_prop, attribute.deserialize(value)) return klass def batch_get(self, item_keys: List[Tuple[str, str]], filter_: RowFilterUnion = None, limit=0): """Get a list of rows from BigTable via an iterable of tuples of ``item_keys`` Raises: :class:`BigTableError`: Exception for any query-related errors Returns: List[:class:`BigTableModel`]: List of loaded BigTableModel types """ row_set = RowSet() for keys in item_keys: key = BigTableModel.get_row_key(keys) row_set.add_row_key(key) return self._read_rows_from_set(limit=limit, filter_=filter_, row_set=row_set) def batch_get_range(self, item_keys: List[Tuple[str, str]], start_date: str, end_date: str, filter_: RowFilterUnion = None, end_key_prefix: bool = True, limit=0): """Get a range of rows from BigTable via list of prefix ``item_keys`` Raises: :class:`BigTableError`: Exception for any query-related errors Args: item_keys: This is a list of tuples of row key prefixes start_date: Date to be suffixed to the row key for the ``start_key`` end_date: Date to be suffixed to the row key for the ``end_key`` filter_: Any additional BigTable RowFilter end_key_prefix: If our end key is a prefix, add one day to end date in query range limit: Defaults to unlimited Returns: List[:class:`BigTableModel`]: List of loaded BigTableModel types """ if end_key_prefix: # If our end key is a prefix, add one day to end date (avoids determining end key) end_date = DateUtils.as_str(DateUtils.from_str(end_date) + timedelta(days=1)) end_inclusive = not end_key_prefix row_set = RowSet() prefix_keys: Tuple[str, str] for prefix_keys in item_keys: start_key_params = prefix_keys + (start_date,) end_key_params = prefix_keys + (end_date,) start_key = BigTableModel.get_row_key(start_key_params) end_key = BigTableModel.get_row_key(end_key_params) row_set.add_row_range_from_keys(start_key=start_key, end_key=end_key, start_inclusive=True, end_inclusive=end_inclusive) return self._read_rows_from_set(limit=limit, filter_=filter_, row_set=row_set, end_inclusive=end_inclusive) def batch_get_range_multiset(self, item_keys: Iterable[Tuple[bytes, bytes]], filter_: RowFilterUnion = None, limit=0, end_inclusive=False): """Get a range of rows from BigTable via list of pre-built ``item_keys``. Note these are complete start and end keys in bytes. Returns: List[:class:`BigTableModel`]: List of loaded BigTableModel types """ row_set = RowSet() for start_key, end_key in item_keys: row_set.add_row_range_from_keys(start_key=start_key, end_key=end_key, start_inclusive=True, end_inclusive=end_inclusive) return self._read_rows_from_set(limit=limit, filter_=filter_, row_set=row_set) def _read_rows_from_set(self, limit: int, row_set: RowSet, filter_: RowFilterUnion = None, end_inclusive=False): """Wrapper for BigTable's ``read_rows()`` specific to using a :class:`RowSet` Raises: :class:`BigTableError`: Exception for any query-related errors Returns: List[:class:`BigTableModel`]: List of loaded BigTableModel types """ try: read_rows: PartialRowsData = self.service.table.read_rows( limit=limit, filter_=filter_, row_set=row_set, end_inclusive=end_inclusive) return [self.from_raw_data(row) for row in read_rows] except Exception as e: # pragma: no cover raise BigTableError({ 'code': Codes.database_query_error.value, 'description': str(e), }) @property def country_stats(self) -> CountriesStruct: """Property for legacy models where countries are individual Bigtable columns""" countries = {} for country in self.Meta.country_codes: stats = getattr(self, country, None) if stats is None: continue countries[country] = stats return CountriesStruct(countries=countries) # ----------------------------------- # # --- Test related helper methods --- # def create_table(self, column_families: dict = None): # pragma: no cover """Create the table associated with the table name of this model's instance. Notes: This convenient helper method is currently only used for unit test setup """ if ENVIRONMENT != 'test': # pragma: no cover raise EnvironmentError('create_table() was called from non-test environment. Aborting.') if not self.service.admin_table.exists(): default_column_families = {COLUMN_FAMILY_ID: MaxVersionsGCRule(1)} column_families = column_families if column_families else default_column_families return self.service.admin_table.create(column_families=column_families) def delete_table(self): # pragma: no cover """Deletes the table associated with this model's instance. Notes: This convenient helper method is currently only used for unit test setup. Raises: :class:`EnvironmentError`: Additional fail-safe if environment is not test """ if ENVIRONMENT != 'test': # pragma: no cover raise EnvironmentError('delete_table() was called from non-test environment. Aborting.') if self.service.admin_table.exists(): return self.service.table.delete() def save(self, row_key: bytes = None): # pragma: no cover """Save the model instance's row data. Notes: This convenient helper method is currently only used for unit test setup. Raises: :class:`BigTableError`: Exception if table does not exist """ if ENVIRONMENT != 'test': # pragma: no cover raise EnvironmentError('save() was called from non-test environment. Aborting.') self._check_table_exists() isrc = self.attribute_values.get(ColumnName.ISRC.value) date = self.attribute_values.get(ColumnName.DATE.value) playlist_id = self.attribute_values.get(ColumnName.PLAYLIST_ID.value) if not date: # pragma: no cover raise AttributeError('Missing attributes to build key') if row_key is None: # very basic (v2) construction of row key attrs = [] if playlist_id: attrs.append(playlist_id) if isrc: attrs.append(isrc) attrs.append(date) keys = [k.decode(UTF8) for k in attrs] row_key = BigTableModel.get_row_key(keys) row: DirectRow = self.service.admin_table.direct_row(row_key) for col, value in self.attribute_values.items(): row.set_cell(COLUMN_FAMILY_ID, bytes(col, UTF8), value) return row.commit() # --- end: Test related helper methods --- # # ---------------------------------------- #