"""Import track visits data from intercom.""" import base64 from datetime import datetime, timedelta import json import logging import os import pymysql import sys import time from typing import Any, Generator, List, Tuple from urllib import error, parse as urlparse, request from urllib.parse import urlencode import config import constants # set up logging LOGGER_LEVEL = getattr(logging, os.environ.get('LOGGER_LEVEL', 'INFO')) logger = logging.getLogger() logger.setLevel(LOGGER_LEVEL) stream_handler = logging.StreamHandler(sys.stdout) stream_handler.setLevel(logging.DEBUG) logger.addHandler(stream_handler) # set up MySQL connection sql_connection = pymysql.connect( host=config.MySQL.HOST, port=config.MySQL.PORT, user=config.MySQL.USER, passwd=config.MySQL.PASSWORD, db=config.MySQL.DATABASE ) sql_cursor = sql_connection.cursor() def handle_rate_limits(f): """Retry on rate limits.""" def wrapped(*args, **kwargs): """Decorator body.""" try_count = 0 while True: try_count = try_count + 1 # call func result = f(*args, **kwargs) # if not 'too many requests' or max retries then stop if result[0] != constants.Http.TOO_MANY_REQUESTS or try_count >= config.MAX_RETRIES: if try_count >= config.MAX_RETRIES: logger.debug('Too many retries, exit') return result logger.debug('Got {}, retry'.format(constants.Http.TOO_MANY_REQUESTS)) # wait to make another try time.sleep(try_count * config.WAIT_RATE) return wrapped @handle_rate_limits def make_request( url: str, headers: dict or None = None, basic_auth_user: str or None = None, basic_auth_password: str or None = None, bearer_token: str or None = None) -> Tuple[int, dict]: """Make HTTP request. Args: url (str): Request URL. headers (dict or None): Request headers. basic_auth_user (str or None): Request basic auth username. basic_auth_password (str or None): Request basic auth password. bearer_token (str or None): Request bearer auth token. Returns: Tuple[int, dict]: Status code and json body. """ auth_token = None if basic_auth_user and basic_auth_password: credentials = ('{}:{}'.format(basic_auth_user, basic_auth_password)) encoded_credentials = base64.b64encode(credentials.encode('ascii')) auth_token = 'Basic {}'.format(encoded_credentials.decode('ascii')) elif bearer_token: auth_token = 'Bearer {}'.format(bearer_token) if auth_token: if not headers: headers = {} headers['Authorization'] = auth_token logger.debug(url) request_obj = request.Request(url, headers=headers) try: response = request.urlopen(request_obj) response_data = response.read() except error.HTTPError as ex: logger.debug('{} {}'.format(ex.code, ex.msg)) return ex.code, {} response_str = response_data.decode("utf-8") logger.debug(response_str) return response.status, json.loads(response_str) def get_admin_account_user(user_id: str) -> Tuple[int, dict]: """Get internal user info. Args: user_id (str): User ID. Returns: Tuple[int, dict]: Status code and json body. """ return make_request( urlparse.urljoin(config.Accounts.URL, user_id), headers={'Content-Type': 'application/x-www-form-urlencoded'}, basic_auth_user=config.Accounts.USER, basic_auth_password=config.Accounts.PASSWORD ) def intercom_paginate(start_url: str) -> Generator[Any, None, None]: """Intercom pagination handling. Args: start_url (str): First page URL. Returns: Generator[Any, None, None]: Items. """ next_url = None while True: if not next_url: next_url = start_url status, data = make_request( next_url, headers={'Accept': 'application/json'}, bearer_token=config.Intercom.TOKEN ) if status != constants.Http.OK: raise ValueError('Incorrect response: {}'.format(status)) yield data next_url = None if 'next' in data['pages']: next_url = data['pages']['next'] if not next_url: break def get_intercom_users() -> Generator[dict, None, None]: """Get all intercom users. Returns: Generator[dict, None, None]: Set of users data. """ pages = intercom_paginate(urlparse.urljoin(config.Intercom.URL, 'users')) for data in pages: for user in data['users']: yield user def get_visited_tracks(user_id: str) -> Generator[Tuple[str, datetime], None, None]: """Get all visited tracks. Args: user_id (str): User ID. Returns: Generator[Tuple[str, datetime], None, None]: Visited tracks IDs and dates """ query_params = {'type': 'user', 'user_id': user_id} relative_url = '{}?{}'.format('events', urlencode(query_params)) pages = intercom_paginate(urlparse.urljoin(config.Intercom.URL, relative_url)) for data in pages: for event in data['events']: create_date = datetime.utcfromtimestamp(event['created_at']) # According to the docs: # "The event list is sorted by the created_at field and ordered descending", # So we can stop on the fist timestamp < now - 30 days. if create_date < config.DATE_FROM: return # get only track visit events if event['event_name'] != 'analyze-track' or not event['metadata']: continue # if "import up to some date" is set # and current event create date is greater # then we need to skip it if config.DATE_TO and create_date > config.DATE_TO: continue track_id = event['metadata']['uri'].replace('spotify:track:', '') yield track_id, create_date def insert_visits(visits: List[tuple]): """Batch insert track visits in MySQL. Args: visits (List[tuple]): Set of track visits data. """ sql_cursor.executemany(constants.SQL_INSERT_VISITS, visits) sql_connection.commit() def get_max_visit_id() -> int: """Get max visit ID from MySQL. Returns: int: Max visit ID. """ sql_cursor.execute(constants.SQL_MAX_VISIT_ID) row = sql_cursor.fetchone() if not row or not row[0]: return 0 return row[0] def check_user_account(user: dict, account: dict) -> bool: """Check if intercom user and internal user are the same. There are some that have the same ID, but complitely different email and name. It seems that it is not safe to load visits for such users. Args: user (dict): Intercom user data. account (dict): Internal admin accounts data. Returns: bool: Valid user or not. """ return ( user['name'] == account['name'] or user['name'] in account['name'] or account['name'] in user['name'] or user['email'] == account['email'] ) def import_data(): """Import intercom track visits data.""" logger.info('Current max visit ID: {}'.format(get_max_visit_id())) insert_queue = [] count_user = 0 for user in get_intercom_users(): count_user = count_user + 1 user_id = user['user_id'] logger.info('({}) {}: {}, {}'.format( count_user, user_id, user['name'], user['email'])) if config.CHECK_WITH_ADMIN_ACCOUNT_API: status, account = get_admin_account_user(user_id) if status == constants.Http.NOT_FOUND or not check_user_account(user, account): if status == constants.Http.NOT_FOUND: logger.info('User does not exist.') else: logger.info('Incorrect account: {} != {} and {} != {}'.format( user['name'], account['name'], user['email'], account['email'])) continue if status != constants.Http.OK: raise ValueError('Incorrect response: {}'.format(status)) for track_id, create_date in get_visited_tracks(user_id): insert_queue.append((user_id, track_id, create_date)) logger.info('{}: {} {}'.format(user_id, track_id, create_date)) if len(insert_queue) >= config.BATCH_SIZE: insert_visits(insert_queue) logger.debug('{} records saved'.format(config.BATCH_SIZE)) insert_queue.clear() if insert_queue: insert_visits(insert_queue) if __name__ == '__main__': import_data()