import base64 import logging import json import boto3 from datetime import datetime, timezone hook_events = [ "Hook-Web-Event", ] RECS_EVENTS_KINESIS_CLIENT = boto3.client( "kinesis", region_name="us-east-2", ) CLIENT_EVENT_NAME_TO_REC_EVENT_NAME = { "RemixHookSongClicked": "HooksTapRemix", "AutoRepeatPlayHook": "HooksRepeatWatch", } logger = logging.getLogger(__name__) def lambda_handler(event, context): """ Lambda handler to check if an event name is in hook_events list and log to RECS_EVENTS_KINESIS stream """ try: if "Records" not in event: return for record in event["Records"]: process_record(record) return {"statusCode": 200, "body": json.dumps("Success")} except Exception as e: print(f"Error in lambda_handler: {str(e)}") return {"statusCode": 500, "body": json.dumps(f"Error: {str(e)}")} def process_record(record): """ Process a single record and check if it should be logged to RECS_EVENTS_KINESIS Unlike the mobile hook events, web events have the event name 'Hook-Web-Event', so we can filter there. We can also filter out events if they have a hook ID. """ if "kinesis" not in record: return kinesis = record["kinesis"] if "data" not in kinesis: return payload = base64.b64decode(record["kinesis"]["data"]) payload_as_json = json.loads(payload) request_body = payload_as_json.get("request_body", "") if not request_body: return decoded_request_body = json.loads(base64.b64decode(request_body)) # Check if this is a Hook-Web-Event event_name = decoded_request_body.get("event", "") if event_name != "Hook-Web-Event": return properties = decoded_request_body.get("properties", {}) if not properties: return action_name = properties.get("actionName", "") rec_event_name = CLIENT_EVENT_NAME_TO_REC_EVENT_NAME.get(action_name, "") if not rec_event_name: return # For web events, hook ID is in the context context = properties.get("context", {}) hook_id = context.get("hookId", "") if not hook_id: return user_id = properties.get("userId", "") if not user_id: return recs_event = {} recs_event["name"] = rec_event_name recs_event["source"] = "web" recs_event["user_id"] = user_id recs_event["timestamp"] = str(datetime.now(timezone.utc).isoformat()) recs_event["properties"] = { "hook_id": hook_id, } logger.info(f"Logging {rec_event_name} event: {recs_event}") RECS_EVENTS_KINESIS_CLIENT.put_record( StreamName="rec-events-stream", Data=json.dumps(recs_event), PartitionKey=str(datetime.now(timezone.utc).isoformat()), )