"""Methods to operate with the State.""" from itertools import groupby from operator import attrgetter from typing import Sequence, cast from lambdacommon.common_config import logger import config from ..connectors.dynamodb import get_client, get_table from ..errors.store_state import StoreStateLoadingError, StoreStateProcessingError, StoreStateSavingError from ..models.store import Store from ..models.store_state import StoreState from ..schemas import StoreStateDBSchema from ..utils import has_duplicates from .feed_status import get_feed_statuses __all__ = ['get_latest_store_states', 'get_updated_states', 'load_store_states', 'save_store_states'] EXCLUDED_FEEDS = (16,) def save_store_states(states: Sequence[StoreState]) -> None: """Save a list of state objects to DynamoDB in a single batch operation. Args: states (Sequence[StoreState]): A list of State objects to save. """ if not states: raise StoreStateSavingError('No store states to save.') if has_duplicates([state.store for state in states]): raise StoreStateSavingError('Store states are not unique.') schema = StoreStateDBSchema() table = get_table(config.STORE_STATE_TABLE_NAME) try: with table.batch_writer() as batch: for state in states: batch.put_item(Item=schema.dump(state)) logger.debug(f'Saved {len(states)} states.') except Exception as e: raise StoreStateSavingError('Could not save states in batch.') from e def load_store_states(stores: Sequence[Store]) -> list[StoreState]: """Load multiple states from DynamoDB in a single batch operation. Args: stores (Sequence[Store]): A list of stores to load the state for. Returns: A list of the found StoreState objects by Store. """ if not stores: raise StoreStateLoadingError('No stores to load.') if has_duplicates(stores): raise StoreStateLoadingError('Stores are not unique.') schema = StoreStateDBSchema() primary_key_name = str(schema.fields['store'].data_key) keys_to_get = [{primary_key_name: store.value} for store in stores] client = get_client() try: response = client.batch_get_item(RequestItems={config.STORE_STATE_TABLE_NAME: {'Keys': keys_to_get}}) items = response.get('Responses', {}).get(config.STORE_STATE_TABLE_NAME, []) result = schema.load(items, many=True) logger.debug(f'Loaded {len(result)} states.') return result except Exception as e: raise StoreStateLoadingError('Could not load states in batch.') from e def get_latest_store_states(stores: Sequence[Store]) -> list[StoreState]: """Get a sequence of StoreState objects from feed states. This logic adapted from https://github.com/theorchard/frontend-insights/blob/c1d98ab36e0b68ff4d4ad5ffbd70267008b8dd40/src/apollo/queries/analyticsFeedsStatus.ts Args: stores (Sequence): A list of stores to load the state for. Returns: A list of the found StoreState objects by Store. """ if not stores: raise StoreStateLoadingError('No stores to load.') if has_duplicates(stores): raise StoreStateLoadingError('Stores are not unique.') feed_statuses = get_feed_statuses() intermediate_sources = [] for status in feed_statuses: if status.feed.store.store_id not in stores: continue if status.feed.id in EXCLUDED_FEEDS: continue intermediate_sources.append( StoreState(store=Store(status.feed.store.store_id), last_available_date=status.store_high_water_mark) ) store_states: list[StoreState] = [] sorted_sources = sorted(intermediate_sources, key=attrgetter('store')) for _store, group in groupby(sorted_sources, key=attrgetter('store')): earliest_state = cast(StoreState, min(group, key=attrgetter('last_available_date'))) store_states.append(earliest_state) return store_states def get_updated_states(new: Sequence[StoreState], existed: Sequence[StoreState]) -> list[StoreState]: """Compare two state sequences by last available date. Args: new (Sequence[StoreState]): A sequence of new states to compare. existed (Sequence[StoreState]): A sequence of existing states to compare. Returns: A sequence of updated states. """ if not new: return [] if has_duplicates(new, key=attrgetter('store')): raise StoreStateProcessingError('New store states are not unique.') if has_duplicates(existed, key=attrgetter('store')): raise StoreStateProcessingError('Existed store states are not unique.') existed_state_lookup: dict[Store, StoreState] = {state.store: state for state in existed} updated_states: list[StoreState] = [] for new_state in new: existing_state = existed_state_lookup.get(new_state.store) if existing_state is None or new_state.last_available_date > existing_state.last_available_date: updated_states.append(new_state) return updated_states