# Unless explicitly stated otherwise all files in this repository are licensed # under the Apache License Version 2.0. # This product includes software developed at Datadog (https://www.datadoghq.com/). # Copyright 2017 Datadog, Inc. from __future__ import print_function import base64 import json import urllib import boto3 import time import os import pprint import socket import ssl import re import StringIO import gzip from base64 import b64decode from datadog import datadog_lambda_wrapper, lambda_metric # Parameters # ddApiKey: Datadog API Key try: ENCRYPTED = os.environ['DD_KMS_API_KEY'] ddApiKey = boto3.client('kms').decrypt(CiphertextBlob=b64decode(ENCRYPTED))['Plaintext'] except Exception: try: ddApiKey = os.environ['DD_API_KEY'] except Exception: pass # metadata: Additional metadata to send with the logs metadata = { "ddsourcecategory": "aws", } host = "lambda-intake.logs.datadoghq.com" ssl_port = 10516 cloudtrail_regex = re.compile('\d+_CloudTrail_\w{2}-\w{4,9}-\d_\d{8}T\d{4}Z.+.json.gz$', re.I) DD_SOURCE = "ddsource" DD_CUSTOM_TAGS = "ddtags" DD_SERVICE = "service" @datadog_lambda_wrapper def lambda_handler(event, context): # Check prerequisites if ddApiKey == "" or ddApiKey == "": raise Exception( "You must configure your API key before starting this lambda function (see #Parameters section)" ) # Attach Datadog's Socket s = connect_to_datadog(host, ssl_port) # Add the context to meta if "aws" not in metadata: metadata["aws"] = {} aws_meta = metadata["aws"] aws_meta["function_version"] = context.function_version aws_meta["invoked_function_arn"] = context.invoked_function_arn #Add custom tags here by adding new value with the following format "key1:value1, key2:value2" - might be subject to modifications metadata[DD_CUSTOM_TAGS] = "forwardername:" + context.function_name.lower()+ ",memorysize:"+ context.memory_limit_in_mb try: print("event") pprint.pprint(event) print("context") pprint.pprint(context) logs = generate_logs(event,context) print('sending auth0_data_trigger') lambda_metric("auth0.logtrigger",1) for log in logs: s = safe_submit_log(s, log) except Exception as e: print('Unexpected exception: {} for event {}'.format(str(e), event)) finally: s.close() def connect_to_datadog(host, port): s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) s = ssl.wrap_socket(s) s.connect((host, port)) return s def generate_logs(event,context): try: # Route to the corresponding parser event_type = parse_event_type(event) if event_type == "s3": logs = s3_handler(event) elif event_type == "awslogs": logs = awslogs_handler(event,context) elif event_type == "events": logs = cwevent_handler(event) elif event_type == "sns": logs = sns_handler(event) except Exception as e: # Logs through the socket the error err_message = 'Error parsing the object. Exception: {} for event {}'.format(str(e), event) logs = [err_message] return logs def safe_submit_log(s, log): try: send_entry(s, log) except Exception as e: # retry once s = connect_to_datadog(host, ssl_port) send_entry(s, log) return s # Utility functions def parse_event_type(event): if "Records" in event and len(event["Records"]) > 0: if "s3" in event["Records"][0]: return "s3" elif "Sns" in event["Records"][0]: return "sns" elif "awslogs" in event: return "awslogs" elif "detail" in event: return "events" raise Exception("Event type not supported (see #Event supported section)") # Handle S3 events def s3_handler(event): s3 = boto3.client('s3') # Get the object from the event and show its content type bucket = event['Records'][0]['s3']['bucket']['name'] key = urllib.unquote_plus(event['Records'][0]['s3']['object']['key']).decode('utf8') metadata[DD_SOURCE] = parse_event_source(event, key) ##default service to source value metadata[DD_SERVICE] = metadata[DD_SOURCE] # Extract the S3 object response = s3.get_object(Bucket=bucket, Key=key) body = response['Body'] data = body.read() structured_logs = [] # If the name has a .gz extension, then decompress the data if key[-3:] == '.gz': with gzip.GzipFile(fileobj=StringIO.StringIO(data)) as decompress_stream: data = decompress_stream.read() if is_cloudtrail(str(key)): cloud_trail = json.loads(data) for event in cloud_trail['Records']: # Create structured object and send it structured_line = merge_dicts(event, {"aws": {"s3": {"bucket": bucket, "key": key}}}) structured_logs.append(structured_line) else: # Send lines to Datadog for line in data.splitlines(): # Create structured object and send it structured_line = {"aws": {"s3": {"bucket": bucket, "key": key}}, "message": line} structured_logs.append(structured_line) return structured_logs def scrub_legacy_login(type, description): if type == 'fu' and '|' in description: return 's_orch' return type # Handle CloudWatch logs def awslogs_handler(event,context): # Get logs with gzip.GzipFile(fileobj=StringIO.StringIO(base64.b64decode(event["awslogs"]["data"]))) as decompress_stream: data = decompress_stream.read() logs = json.loads(str(data)) pprint.pprint(logs) #Set the source on the logs source = logs.get("logGroup", "cloudwatch") metadata[DD_SOURCE] = parse_event_source(event, source) ##default service to source value metadata[DD_SERVICE] = metadata[DD_SOURCE] structured_logs = [] # Send lines to Datadog i = 0 for log in logs["logEvents"]: print('sending auth0_log_processed') lambda_metric("auth0_log_processed",1) i = i + 1 print("log event") print(type(log)) pprint.pprint(log) message_dict = json.loads(str(log["message"])) print(type(message_dict)) pprint.pprint(message_dict) message_dict['type'] = scrub_legacy_login(message_dict.get('type'), message_dict.get('description')) type_long = get_event_details(message_dict.get('type')) print('log type {}'.format(type_long)) send_auth0_login_metric(logs["logGroup"], message_dict.get('type')) message_dict['type_long'] = type_long log["message"] = message_dict # Create structured object and send it structured_line = merge_dicts(log, { "aws": { "awslogs": { "logGroup": logs["logGroup"], "logStream": logs["logStream"], "owner": logs["owner"] } } }) structured_logs.append(structured_line) print('count {}'.format(i)) return structured_logs #Handle Cloudwatch Events def cwevent_handler(event): data = event #Set the source on the log source = data.get("source", "cloudwatch") service = source.split(".") if len(service)>1: metadata[DD_SOURCE] = service[1] else: metadata[DD_SOURCE] = "cloudwatch" ##default service to source value metadata[DD_SERVICE] = metadata[DD_SOURCE] structured_logs = [] structured_logs.append(data) return structured_logs # Handle Sns events def sns_handler(event): data = event # Set the source on the log metadata[DD_SOURCE] = parse_event_source(event, "sns") structured_logs = [] for ev in data['Records']: # Create structured object and send it structured_line = ev structured_logs.append(structured_line) return structured_logs def send_entry(s, log_entry): # The log_entry can only be a string or a dict if isinstance(log_entry, str): log_entry = {"message": log_entry} elif not isinstance(log_entry, dict): raise Exception( "Cannot send the entry as it must be either a string or a dict. Provided entry: " + str(log_entry) ) # Merge with metadata log_entry = merge_dicts(log_entry, metadata) # Send to Datadog str_entry = json.dumps(log_entry) #For debugging purpose uncomment the following line print(str_entry) prefix = "%s " % ddApiKey return s.send((prefix + str_entry + "\n").encode("UTF-8")) def merge_dicts(a, b, path=None): if path is None: path = [] for key in b: if key in a: if isinstance(a[key], dict) and isinstance(b[key], dict): merge_dicts(a[key], b[key], path + [str(key)]) elif a[key] == b[key]: pass # same leaf value else: raise Exception( 'Conflict while merging metadatas and the log entry at %s' % '.'.join(path + [str(key)]) ) else: a[key] = b[key] return a def is_cloudtrail(key): match = cloudtrail_regex.search(key) return bool(match) def parse_event_source(event, key): for source in ["lambda", "redshift", "cloudfront", "kinesis", "mariadb", "mysql", "apigateway", "route53", "vpc", "rds", "sns"]: if source in key: return source if "elasticloadbalancing" in key: return "elb" if is_cloudtrail(str(key)): return "cloudtrail" if "awslogs" in event: return "cloudwatch" if "Records" in event and len(event["Records"]) > 0: if "s3" in event["Records"][0]: return "s3" return "aws" def get_event_details(event_short): event_mapping = dict( api_limit='Rate Limit On API', cls='Code/Link Sent', coff='Connector Offline', con='Connector Online', cs='Code Sent', du='Deleted User', f='Failed Login', fapi='Failed API Operation', fc='Failed by Connector', fce='Failed Change Email', fco='Failed by CORS', fcoa='Failed cross-origin authentication', fcp='Failed Change Password', fcph='Failed Post Change Password Hook', fcpn='Failed Change Phone Number', fcpr='Failed Change Password Request', fcpro='Failed Connector Provisioning', fcu='Failed Change Username', fd='Failed Delegation', fdu='Failed User Deletion', feacft='Failed Exchange of authorization code for Access Token', feccft='Failed Exchange of Access Token for a Client Credentials', feoobft='Failed Exchange of Password and OOB Challenge for Access Token', feotpft='Failed Exchange of Password and OTP Challenge for Access Token', fepft='Failed Exchange of Password for Access Token', fercft='Failed Exchange of Password and MFA Recovery code for Access Token', fertft='Failed Exchange of Refresh Token for Access Token', flo='Failed Logout', fn='Failed Sending Notification', fp='Failed Login (Incorrect Password)', fs='Failed Signup', fsa='Failed Silent Auth', fu='Failed Login (Invalid Email/Username)', fui='Failed users import', fv='Failed Sending Verification Email', fvr='Failed Processing Verification Email Request', gd_auth_failed='OTP Auth failed', gd_auth_rejected='OTP Auth rejected', gd_auth_succeed='OTP Auth success', gd_enrollment_complete='Guardian enrollment complete', gd_module_switch='Module switch', gd_otp_rate_limit_exceed='Too many failures', gd_recovery_failed='Multi-factor recovery code failed', gd_recovery_rate_limit_exceed='Multi-factor recovery code has failed too many times', gd_recovery_succeed='Multi-factor recovery code succeeded authorization', gd_send_pn='Push notification for MFA sent successfully sent', gd_send_sms='SMS for MFA sent successfully sent', gd_start_auth='MFA Second factor started', gd_start_enroll='MFA Enroll started', gd_tenant_update='Guardian tenant update', gd_unenroll='MFA Device unenrolled', gd_update_device_account='MFA device updated', gd_user_delete='Deleted MFA user account.', limit_delegation='Too Many Delegation API Calls', limit_mu='Blocked IP address 100 failed logins attempts Anomaly Detection', limit_ui='Too Many Userinfo API Calls', limit_wc='Blocked Account IP address 10 failed logins Anomaly Detection', pwd_leak='Breached leaked password', s='Success Login', s_orch='Success Login Legacy', sapi='Success API Operation', sce='Success Change Email', scoa='Success cross-origin authentication', scp='Success Change Password', scph='Success Post Change Password Hook', scpn='Success Change Phone Number', scpr='Success Change Password Request', scu='Success Change Username', sd='Success Delegation', sdu='Success User Deletion', seacft='Success Exchange of authorization code for Access Token', seccft='Success Exchange of Access Token for a Client Credentials', seoobft='Success Exchange of Password and OOB Challenge for Access Token', seotpft='Success Exchange of Password and OTP Challenge for Access Token', sepft='Success Exchange of Password for Access Token', sercft='Success Exchange of Password and MFA Recovery code for Access Token', sertft='Success Exchange of Refresh Token for Access Token', slo='Success Logout', ss='Success Signup', ssa='Success Silent Auth', sui='Success users import', sv='Success Verification Email', svr='Success Verification Email Request', sys_os_update_end='Auth0 OS Update Ended', sys_os_update_start='Auth0 OS Update Started', sys_update_end='Auth0 Update Ended', sys_update_start='Auth0 Update Started', ublkdu='User login block released', w='Warnings During Login') if event_short not in event_mapping: return 'unknown' return(event_mapping[event_short]) def send_auth0_login_metric(log_group, type): if log_group == 'qa-auth0-logs': environment = 'qa' else: environment = 'prod' failed_login_types = [ 'f', #Failed Login 'feacft', #Failed Exchange of authorization code for Access Token 'feccft', #Failed Exchange of Access Token for a Client Credentials 'feoobft', #Failed Exchange of Password and OOB Challenge for Access Token 'feotpft', #Failed Exchange of Password and OTP Challenge for Access Token 'fepft', #Failed Exchange of Password for Access Token 'fercft', #Failed Exchange of Password and MFA Recovery code for Access Token 'fertft', #Failed Exchange of Refresh Token for Access Token 'fp', #Failed Login (Incorrect Password) 'fu', #Failed Login (Invalid Email/Username) 'fvr', #Failed Processing Verification Email Request 'gd_auth_failed', #OTP Auth failed 'gd_auth_rejected'] #OTP Auth rejected success_login_types = [ 's', # Success Login 's_orch'] #Success Login Legacy' if type in success_login_types: metric = 'auth0.{}.login_success'.format(environment) print(metric) lambda_metric(metric, 1) if type in failed_login_types: metric = 'auth0.{}.login_failed'.format(environment) print(metric) lambda_metric(metric, 1)