from datetime import datetime import json import time from smelog.factory import BoundLogger from copy_auth0_logs_to_cloudwatch import auth0_management, aws_logs, config, s3_storage from copy_auth0_logs_to_cloudwatch.constants import ( KEY_LAST_ENTRY_ID, LOG_TYPE_AUTHENTICATION_API, LOG_TYPE_MANAGEMENT_API) from copy_auth0_logs_to_cloudwatch.redis_db import get_redis_lock class Log: """Log information container. """ def __init__(self, aws_group: str): """Init log object. Args: aws_group (str): AWS CloudWatch Log group name. """ self.aws_group = aws_group self.events = [] self.min_date = None self.sequence_token = None self.aws_stream = None self.size = 0 def get_log_event_size(log_event: str) -> int: """Get log event size in bytes. Args: log_event (str): Log event JSON. Returns: int: Size in bytes. """ return len(log_event.encode('utf-8')) + 26 def save_log_events(log: Log, log_type: str, logger: BoundLogger): """Save log events to AWS CloudWatch Log. Args: log (Log): Log information container. log_type (str): Log type. logger: Logger. """ logger.info(f'Saving {len(log.events)} {log_type} events to AWS.') if not log.aws_stream: log.aws_stream = aws_logs.create_stream(log.aws_group) logger.info(f'Log stream {log.aws_stream} created.') events = sorted(log.events, key=lambda k: k['timestamp']) log.sequence_token = aws_logs.put_events( log.aws_group, log.aws_stream, events, log.sequence_token) log.events.clear() log.min_date = None log.size = 0 def handler(redis_client, logger: BoundLogger): """AWS Lambda handler. Arguments: redis_client: Redis client. logger: Logger. """ # set lock for 10 m to avoid possible race condition between lambdas lock = get_redis_lock(redis_client) if not lock.acquire(blocking=False): logger.info('Locked') return {} start_time = time.time() last_log_id = s3_storage.get(KEY_LAST_ENTRY_ID, '0') new_last_log_id = last_log_id logger.info(f'Last log ID {last_log_id}.') logs = { LOG_TYPE_AUTHENTICATION_API: Log(config.AWS_LOG_GROUP_AUTHENTICATION_API), LOG_TYPE_MANAGEMENT_API: Log(config.AWS_LOG_GROUP_MANAGEMENT_API) } count = 0 for i in range(config.LOG_SEARCH_COUNT): logger.info(f'Auth0 logs loading iteration {i} from {new_last_log_id}.') for log_event in auth0_management.get_logs(new_last_log_id): log_type = ( LOG_TYPE_MANAGEMENT_API if log_event['type'] in config.MANAGEMENT_API_LOG_EVENT_TYPES else LOG_TYPE_AUTHENTICATION_API ) log_date = datetime.strptime(log_event['date'], '%Y-%m-%dT%H:%M:%S.%fZ') log_event_str = json.dumps(log_event) log_event_size = get_log_event_size(log_event_str) if log_event_size > config.AWS_LOG_EVENT_MAX_SIZE_BYTES: log_event_str = log_event_str[:config.AWS_LOG_EVENT_TRIM_SIZE] logger.error(f'Event exceeds max size: {log_event_size} {log_event_str}') log = logs[log_type] if ( ( log.min_date and abs((log_date - log.min_date).total_seconds() / 3600) > config.AWS_LOG_MAX_HOURS_INTERVAL ) or len(log.events) >= config.AWS_LOG_BATCH_MAX_COUNT or log.size + log_event_size >= config.AWS_LOG_BATCH_MAX_SIZE_BYTES): try: save_log_events(log, log_type, logger) except Exception: logger.error(f'Size {log.size}, date {log.min_date}') raise log.min_date = log_date if not log.min_date else min(log_date, log.min_date) log_timestamp = int(log_date.timestamp() * 1000) log.events.append({'timestamp': log_timestamp, 'message': log_event_str}) log.size = log.size + log_event_size new_last_log_id = log_event['log_id'] count = count + 1 if time.time() - start_time > config.MAX_EXECUTION_TIME_SECONDS: break else: continue # it can be reached only when inner loop break is called break for log_type, log in logs.items(): if log.events: save_log_events(log, log_type, logger) if new_last_log_id != last_log_id: logger.info(f'Saving last log ID {new_last_log_id}.') s3_storage.set(KEY_LAST_ENTRY_ID, new_last_log_id) logger.info(f'{count} events processed.') lock.release()