"""Main lambda module.""" import base64 import gzip import json from datetime import datetime from http import HTTPStatus from typing import Any, Dict, List, Set from urllib.parse import unquote from redis import StrictRedis from smelog.factory import BoundLogger from sqlalchemy.dialects.mysql import insert from apollo_main_db.apollo.models import UserMarket from apollo_main_db.push_notifications.models import UserDeviceToken from handle_user_changes import admin_api, apollo_api, config, intercom_api from handle_user_changes.constants import DSP_APP_NAME from handle_user_changes.main_db import session_scope def decode_log_events(encoded_data: str) -> dict: """Decode log events from Lambda arguments. Args: encoded_data (str): Base64 encoded and gzipped log events. Returns: dict: Event logs data. """ decoded_data = base64.b64decode(encoded_data) str_data = gzip.decompress(decoded_data) data = json.loads(str_data) return data def transform_to_intercom_user( user: Dict[str, Any], home_market_id: int or None, job_category_id: int or None, markets: Dict[int, str], job_categories: Dict[int, str], ) -> Dict[str, Any]: """Generate intercom user profile data based on Auth0 user log info. Args: user (Dict[str, Any]): Auth0 user log information. home_market_id (int or None): User's home market ID. job_category_id (int or None): User's job category ID. markets (dict): Id, Name mapper. job_categories (dict): Id, Name mapper. Returns: Dict[str, Any]: Intercom profile information. """ intercom_user = { 'user_id': user['user_id'].replace('auth0|', ''), 'email': user['email'], 'name': user['name'], 'created_at': user['created_at'], 'custom_attributes': { 'Suspended': user.get('blocked', False), 'Confirmed': user.get('email_verified', False) } } custom_attrs = intercom_user['custom_attributes'] # fill home market if available if home_market_id and home_market_id in markets: custom_attrs['Home Market ID'] = home_market_id custom_attrs['Home Market Name'] = markets[home_market_id] # fill job category if available if job_category_id and job_category_id in job_categories: custom_attrs['Job Category ID'] = job_category_id custom_attrs['Job Category Name'] = job_categories[job_category_id] return intercom_user def process_upsert_users(upsert_users: dict, logger: BoundLogger): logger.debug(f'Upserting {len(upsert_users)} users to intercom.') upsert_users_list = list(upsert_users.values()) job_id = None for i in range(0, len(upsert_users_list), config.INTERCOM_USER_UPSERT_BATCH_SIZE): users_chunk = upsert_users_list[ i: i + config.INTERCOM_USER_UPSERT_BATCH_SIZE] job_id = intercom_api.upsert_users(users_chunk, job_id) def process_user_data(user_data_list: List[dict], logger: BoundLogger): """Upsert user data (market ID and email) to MySQL. Args: user_data_list: User data (ID, market ID and email) list. logger: Logger. """ logger.debug(f'Upserting {len(user_data_list)} users to MySQL.') for i in range(0, len(user_data_list), config.MYSQL_USER_DATA_ID_UPDATE_BATCH_SIZE): user_chunk = user_data_list[i: i + config.MYSQL_USER_DATA_ID_UPDATE_BATCH_SIZE] insert_stmt = insert(UserMarket).values(user_chunk) on_duplicate_key_stmt = insert_stmt.on_duplicate_key_update( MarketId=insert_stmt.inserted.MarketId, Email=insert_stmt.inserted.Email ) with session_scope() as session: session.execute(on_duplicate_key_stmt) def deactivate_user_devices(user_list: Set[str], logger: BoundLogger): """Deactivate devices for all blocked or deleted users Args: user_list (list): list of user_id. logger: Logger. """ logger.debug('Deactivating user devices') with session_scope() as session: raw_count = session.query(UserDeviceToken).filter( UserDeviceToken.user_id.in_(user_list) ).update({'is_active': False}, synchronize_session=False) logger.debug(f'Deactivated {raw_count} tokens for {len(user_list)} users') def dump_users_data_to_redis(users_update_date: Dict[str, str], redis_client) -> None: """Dump user's last changes data to redis Args: users_update_date: dict of user_id and description. redis_client: Redis client. """ redis_pipeline = redis_client.pipeline() for user_id, description in users_update_date.items(): timestamp = int(datetime.timestamp(datetime.now())) redis_pipeline.setex(f"{DSP_APP_NAME}auth0:user_{user_id}:updated_at", config.REDIS_KEY_TTL, timestamp) redis_pipeline.setex(f"{DSP_APP_NAME}auth0:user_{user_id}:details", config.REDIS_KEY_TTL, description) redis_pipeline.execute() def handler(event: dict, redis_client: StrictRedis, logger: BoundLogger): """AWS Lambda handler. Args: event (dict): Lambda arguments. redis_client (StrictRedis): Redis client. logger: Logger. """ logger.debug('Decoding data.') data = decode_log_events(event['awslogs']['data']) log_events = [] for record in data['logEvents']: try: event_data = json.loads(record['message']) log_events.append(event_data) except json.JSONDecodeError: logger.error(f'Invalid event {record["message"]}') continue logger.debug(f'Got {len(log_events)} log events.') # do nothing if no log events found if not log_events: return upsert_users = {} user_data = {} users_update_date = {} users_updated_meta_data: Dict[str: str] = {} blocked_users = set() deleted_users = set() markets = {} job_categories = {} if config.ENABLE_INTERCOM_USER_UPDATE: logger.debug('Getting markets.') try: markets = admin_api.get_markets() except Exception as e: logger.error(e) logger.debug('Getting job categories.') try: job_categories = apollo_api.get_job_categories() except Exception as e: logger.error(e) for log_event in log_events: if log_event['type'] not in config.LOG_EVENT_TYPES: # we do not know how to process other events, please update code continue auth0_request = log_event['details']['request'] # if it is user action if auth0_request['path'].startswith('/api/v2/users'): response = log_event['details']['response'] # if it is create or update event if auth0_request['method'] in ('post', 'patch'): user = response['body'] # skip if not user change management api call if 'user_id' not in user: continue user_id = user['user_id'].replace('auth0|', '') logger.debug(f'Processing user {user_id}') update_date = datetime.strptime(user['updated_at'], '%Y-%m-%dT%H:%M:%S.%fZ') if user_id in users_update_date and users_update_date[user_id] > update_date: continue users_update_date[user_id] = update_date users_updated_meta_data[user_id] = log_event.get('description', '') home_market_id = user.get('user_metadata', {}).get('apollo', {}).get('homeMarketID') job_category = user.get('user_metadata', {}).get('apollo_go_job_category') job_category_id = job_category and isinstance(job_category, dict) and job_category.get('id') upsert_users[user_id] = transform_to_intercom_user( user, home_market_id, job_category_id, markets, job_categories) email = user.get('email') if home_market_id or email: user_data[user_id] = {'UserId': user_id, 'MarketId': home_market_id, 'Email': email} # handle is blocked user status is_blocked = user.get('blocked') if is_blocked is True: blocked_users.add(user_id) app_metadata = user.get('app_metadata', {}) # handle blocking byInactivity is_blocked_by_inactivity = app_metadata.get('blockReasons', {}).get('byInactivity') if is_blocked_by_inactivity: blocked_users.add(user_id) # handle blocking byPasswordExpired is_blocked_by_pass_expired = app_metadata.get('blockReasons', {}).get('byPasswordExpired') if is_blocked_by_pass_expired: blocked_users.add(user_id) # remove user_id from blocked_users if it was unblocked in next logs if not is_blocked and not is_blocked_by_inactivity and not is_blocked_by_pass_expired \ and user_id in blocked_users: blocked_users.remove(user_id) # if it is delete event elif auth0_request['method'] == 'delete': unquoted_path = unquote(auth0_request['path']) user_id = unquoted_path.replace('/api/v2/users/auth0|', '') deleted_users.add(user_id) dump_users_data_to_redis(users_updated_meta_data, redis_client) if config.ENABLE_INTERCOM_USER_UPDATE: process_upsert_users(upsert_users, logger) if config.ENABLE_USER_DATA_UPDATE: process_user_data(list(user_data.values()), logger) if config.ENABLE_DEACTIVATE_USER_DEVICES: deactivate_user_devices(blocked_users | deleted_users, logger)