"""Lambda fanout-event function module.""" import datetime import json from lambdacommon.common_config import logger import pytz import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from config import secrets_manager_client from src.models.getstream import GetStream from src.models import graphql_router from src.models import ows_users sentry_on = False if secrets_manager_client: sentry_dsn = secrets_manager_client.get_cred('SENTRY_DSN') if sentry_dsn: sentry_sdk.init( sentry_dsn, integrations=[AwsLambdaIntegration()] ) sentry_on = True logger.info(f'Initializing with Sentry {sentry_on}') # initialize getstream client getstream = GetStream() def handler(event, context): """Lambda entry point.""" try: logger.info(event) # parse message body messages = format_messages(event) submitted = len(messages) # filter out events already in GetStream messages = getstream.filter_activities(messages) processed = len(messages) # check if profile has access, then add to profile feed added = 0 failed = 0 for message in messages: if access_granted(message): try: getstream.add_activity(message) logger.info( f'Successfully added event: {json.dumps(message["event"], default=str)} ' f'to GetStream' ) except Exception: failed += 1 logger.exception('An error occurred while adding activity') continue added += 1 # output results results = { 'submitted': submitted, 'processed': processed, 'added': added, 'failed': failed } logger.info(results) return results except Exception as e: logger.exception(str(e)) raise def format_messages(event): """Format SQS event into messages. Args: event (dict): SQS events with 1 to 10 messages Returns list: messages in event, loaded from JSON """ date_format = '%Y-%m-%dT%H:%M:%S' formatted = [] for message in event.get('Records', []): body = json.loads(message['body']) time = datetime.datetime.strptime(body['event']['time'], date_format) if time.utcoffset() is None: time = pytz.utc.localize(time) body['event']['time'] = time formatted.append(body) return formatted def access_granted(message): """Perform ownership checks on message. Args: message (dict): profile scoped activity event Returns: bool: if messages should be forwarded or not """ # extract standard event information verb = message['event']['verb'] profile_type, profile_id = message['feed_id'].split('_') identity_response = ows_users.get_identity(profile_id, profile_type) if identity_response.status_code == 404: logger.info(f'404 identity response for {profile_type} {profile_id}. Skipping Fanout...') return False identity_response = identity_response.json() identity_id = identity_response['id'] # perform ownership filtering checks has_access = True if verb == 'playlist_placements': has_access = graphql_router.profile_can_access_sound_recording( profile_id, profile_type, identity_id, message['event']['object']['sound_recording']['isrc'] ) elif verb == 'trending_tracks': has_access = graphql_router.profile_can_access_track( profile_id, profile_type, identity_id, message['event']['object']['track']['isrc'], message['event']['object']['track']['id'] ) if not has_access: logger.info( f'identity: {identity_id} with brand: {identity_response.get("default_brand")} ' f'and profile: {profile_type} {profile_id} does not have access ' f'to this {verb}. Skipping Fanout...' ) return has_access