"""Model for GetStream interaction.""" from getstream_connector import get_client from stream import exceptions as stream_exceptions from config import secrets_manager_client class GetStream: """Wrapper class for GetStream integrations.""" def __init__(self): """Initialize connection.""" self.client = self._init_client() @property def saved_headers(self): """Get saved headers.""" return self.client.saved_headers def _init_client(self): """Return an initialized Stream client. Returns: stream.Client: the stream client """ stream_key = self._load_secret('STREAM_API_KEY') stream_secret = self._load_secret('STREAM_API_SECRET') stream_region = self._load_secret('STREAM_API_REGION') return get_client(stream_key, stream_secret, location=stream_region) def add_activity(self, activity): """Add activity to a feed. Args: activity (dict): details of events and destination feed Returns: None """ feed = self.client.feed( activity['feed_group'], activity['feed_id'] ) feed.add_activity(activity['event']) def filter_activities(self, activities): """Query GetStream in batch for activities already added. Args: activities (list): dict with following keys feed_group (str): GetStream feed group to send activity to feed_id (str): GetStream feed id to send activity to event (dict): event metadata, requires foreign_id and time Returns: list: activites that are NOT in GetStream already, same format in """ # see PLATFORM-2391 # max batch is 10 in lambda connector anyway, but be careful if len(activities) > 10: raise Exception('unexpectedly large batch of activities') # format foreign ids and times to bulk query GetStream id_tuples = [ (x['event']['foreign_id'], x['event']['time']) for x in activities ] response = self.client.get_activities(foreign_id_times=id_tuples) def _primary_activity_key(activity_event): """Identity activity dict in a globally unique way.""" return f"{activity_event['foreign_id']}:{activity_event['time']}" # get primary keys of events already in GetStream's system past_keys = [ _primary_activity_key(x) for x in response['results'] ] # filter out events already existing in GetStream return [ x for x in activities if _primary_activity_key(x['event']) not in past_keys ] @staticmethod def _load_secret(name): """Read secret config value.""" result = secrets_manager_client.get_cred(name) if not result: raise stream_exceptions.ApiKeyException(f'unable to load {name}') return result