from datetime import datetime, timedelta import argparse import copy import gzip import json import logging import os import sys import boto3 import threading import time import re from queue import Queue import config from auth0 import Auth0 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) class Auth0RequestThread(threading.Thread): def __init__(self, user_id, queue, callback): self.user_id = user_id super().__init__(target=lambda q, arg1: q.put((user_id, callback(user_id))), args=(queue, user_id)) time.sleep(constants.SLEEP_BETWEEN_AUTH0_PERMS) self.start() def get_users_with_permission(auth0: Auth0, permission_name: str or None): # get all users users = auth0.get_users_export() logger.info('Total users : {}'.format(len(users))) users_with_permission = {} que = Queue() threads = [] for user_id in list(users): # logger.info('User : {}'.format(user_id)) # Get userinfo # time.sleep(constants.SLEEP_BETWEEN_AUTH0_PERMS) threads.append(Auth0RequestThread(user_id, que, auth0.get_user)) # Get user permissions threads.append(Auth0RequestThread(user_id, que, auth0.get_user_permissions)) # Get user roles threads.append(Auth0RequestThread(user_id, que, auth0.get_user_roles)) # Join all the threads logger.info('Wait for all threads finished') # for t in threads: # t.join() alive = len(threads) while alive > 0: # logger.info(threads) alive = len([t for t in threads if t.is_alive()]) logger.info('Threads: {} of {} running.'.format(alive, len(threads))) time.sleep(1) logger.info('Get permissions result') while not que.empty(): user_id, user = que.get() if user_id in users_with_permission: users_with_permission[user_id] = {**users_with_permission[user_id], **user} else: users_with_permission[user_id] = copy.deepcopy(user) # if any(i['permission_name'] == permission_name for i in permissions): # users_with_permission[user_id] = { # 'name': users[user_id]['name'], # 'email': users[user_id]['email'], # 'events': [] # } return users_with_permission def get_events(users, kwargs_events): client = boto3.client('logs') need_next_event_chunk = True while need_next_event_chunk: response = client.get_log_events(**kwargs_events) for log_event in response['events']: # now we are in the event level message = json.loads(log_event['message']) if 'user_id' in message and message['user_id'] and message['type'] in constants.EVENT_TYPES: if message['user_id'] in users: # User exists users[message['user_id']]['events'].append({ 'date': message['date'], 'type': message['type'], 'app': message['client_name'] if 'client_name' in message else None }) # Uncomment this in case ability to add users to list # else: # users[message['user_id']] = { # 'name': message['user_name'], # 'email': message['user_name'], # 'events': [{ # 'date': message['date'], # 'type': message['type'], # 'app': message['client_name'] if 'client_name' in message else None # }] # } logger.info(message) if 'nextToken' in response: kwargs_events['nextToken'] = response['nextToken'] else: need_next_event_chunk = False def get_logs(users, start_time): threads = [] client = boto3.client('logs') need_next_stream = True next_stream_token = None stream_count = 0 events_count = 0 kwargs_streams = dict( logGroupName=config.EVENT_LOG_GROUP_NAME, orderBy='LastEventTime', descending=True, ) kwargs_events = dict( logGroupName=config.EVENT_LOG_GROUP_NAME, ) while need_next_stream: response = client.describe_log_streams(**kwargs_streams) if 'nextToken' in response: kwargs_streams['nextToken'] = response['nextToken'] else: need_next_stream = False # get events from stream for log_stream in response['logStreams']: logger.info('Log stream: {}'.format(log_stream['logStreamName'])) # break loop in case of start_date iso_format = re.sub(r'T(\d\d)-(\d\d)-(\d\d).(\d+)Z', r'T\1:\2:\3.\4+00:00', log_stream['logStreamName']) if datetime.fromisoformat(iso_format).replace(tzinfo=None) < start_time: need_next_stream = False time.sleep(constants.SLEEP_BETWEEN_LOG_EVENTS) kwargs_events['logStreamName'] = log_stream['logStreamName'] t = threading.Thread(target=get_events, args=(users, kwargs_events)) t.start() threads.append(t) stream_count += 1 # logger.info(response) logger.info('Wait for all threads finished') # Join all the threads for t in threads: t.join() logger.info('All threads finished') return users def output_users(users: dict): logger.info('Report data:') with open('report.txt', 'w') as report_file: for user in users: for event in users[user]['events']: dt = datetime.fromisoformat(event['date'].replace('Z', '+00:00')) s = ';'.join((users[user]['name'], users[user]['email'], event['app'], constants.EVENT_TYPES[event['type']], str(dt.year), str(dt.month), str(dt.day), str(dt.strftime("%d/%m/%Y")), str(dt.strftime("%H:%M:%S")) )) logger.info(s) report_file.write(s + "\n") def output_users_to_file(file_name: str, file_type: str, users: dict): if file_type == 'gzip': json_str = json.dumps(users, indent=4, separators=(',', ': ')) json_bytes = json_str.encode('utf-8') with gzip.GzipFile(file_name, 'w') as f: f.write(json_bytes) else: with open(file_name, 'w') as f: json.dump(users, f, indent=4, separators=(',', ': ')) def main(): logger.info('Start report...') parser = argparse.ArgumentParser() parser.add_argument("--dump_users", action="store_true") parser.add_argument("--upload_s3", action="store_true") parser.add_argument('file_type', choices=['json', 'gzip'], default='json') parser.add_argument('outfile', type=argparse.FileType('w')) args = parser.parse_args() if args.dump_users: auth0 = Auth0(logger) users_with_permission = get_users_with_permission(auth0, None) output_users_to_file(args.outfile.name, args.file_type, users_with_permission) if args.upload_s3: pass logger.info('Finished.') if __name__ == '__main__': main()