"""Renew script alert Lambda module.""" import base64 import json import logging import boto3 import pycurl import util JOB_THRESHOLD_MAX = 5000 PROCESS_THRESHOLD_MIN = 10 ALERT_CHANNEL = '#vector-system-alerts' logger = logging.getLogger() logger.setLevel(logging.INFO) # kms boto3.set_stream_logger(name='botocore') kms = boto3.client('kms') # token = kms.encrypt(KeyId='alias/test-reportsar-ro', Plaintext='testdata') # cipherblob = base64.b64encode(token.get('CiphertextBlob')) # logger.info(cipherblob) # db settings db_host = 'reportsdd.theorchard.com' name = 'dd_alert' db_name = 'direct_delivery' port = 3306 cipherblob = ( 'AQECAHjzVuo0MC4YKXdTmnvlTqPWM7Z+rgrVs0TNIFJowpkefgAAAG0wawYJKoZIhvcNAQcGo' 'F4wXAIBADBXBgkqhkiG9w0BBwEwHgYJYIZIAWUDBAEuMBEEDKuUDGhZF7DGrpcv3wIBEIAquP' 'L2+BuCA9Iad6vwjvO/Wq3ZxRwmXZF8PWCwQMtEK2jkGE9+L2ollO7d') kms_response = kms.decrypt(CiphertextBlob=base64.b64decode(cipherblob)) password = kms_response['Plaintext'] # slack url slack_cipher = ( 'AQECAHjzVuo0MC4YKXdTmnvlTqPWM7Z+rgrVs0TNIFJowpkefgAAAK8wgawGCSqGSIb3DQEHB' 'qCBnjCBmwIBADCBlQYJKoZIhvcNAQcBMB4GCWCGSAFlAwQBLjARBAy5LeXp8a2Y++tPwBkCAR' 'CAaGVdb86Es+V62wmlTSSMcMUW5FHLUP4fvPM251FJbNqe0+7gH5Hde9/kLLWn+4DUxIDiYwW' 'PtcgbOdugGfazPfDus1JorJdMF/gDNoqtEXB1F0WimxYzSVmzkNaUEGHUirYPgpJgbRXr') slack_kms_response = kms.decrypt( CiphertextBlob=base64.b64decode(slack_cipher)) slack_url = slack_kms_response['Plaintext'] slack_url = slack_url.strip() logger.info(slack_url) def handler(event, context): """Entrypoint to Lambda function. Args: event (dict, list, str, int, float, NoneType): AWS Lambda event. context (LambdaContext): AWS Lambda context. """ with util.reportsdd_connection( db_host, name, password, db_name) as conn: with conn.cursor() as cursor: cursor.execute(''' SELECT eqd.dms_master_master_id store_id, ( SELECT customer_name FROM customer_master_master cmm WHERE eqd.dms_master_master_id = cmm.customer_master_master_id ) store_name, MIN(eqd.last_updated) oldest_date, MAX(eqd.last_updated) newest_date, COUNT(*) job_count FROM encoding_queue_detail eqd INNER JOIN encoding_queue eq ON eq.encoding_queue_id = eqd.encoding_queue_id AND eq.encoder_id=18 INNER JOIN dms_delivery_spec dds ON dds.dms_master_master_id = eqd.dms_master_master_id AND eq.encoding_order_type = dds.order_type AND dds.encoding = 'Y' WHERE eqd.status = 'new' AND TIMESTAMPDIFF(MINUTE, eqd.last_updated, NOW()) > %s GROUP BY eqd.dms_master_master_id''', PROCESS_THRESHOLD_MIN) stats = [row for row in cursor.fetchall()] logger.info('Stats found: {stats}'.format(stats=stats)) if len(stats): messages = [] for ( store_id, store_name, oldest_date, newest_date, job_count) in stats: cursor.execute(''' SELECT count(eqd.encoding_queue_detail_id) AS processing FROM encoding_queue eq INNER JOIN encoding_queue_detail eqd ON eq.encoding_queue_id = eqd.encoding_queue_id INNER JOIN dms_delivery_spec dds ON dds.dms_master_master_id = eqd.dms_master_master_id AND eq.encoding_order_type = dds.order_type AND dds.encoding = 'Y' WHERE eqd.status in ('ready_to_encode', 'encoding', 'encoded', 'ready_to_deliver', 'queued_for_delivery', 'delivering') AND eq.encoder_id = 18 AND eqd.dms_master_master_id = %s''', store_id) row_processing = cursor.fetchone() if row_processing[0] < JOB_THRESHOLD_MAX: messages.append( 'New jobs from *{oldest_date}* to *{newest_date}*' ' for *{store_name}* (DMS {store_id}): ' '*{job_count}* unprocessed'.format( store_id=store_id, store_name=store_name, oldest_date=oldest_date, newest_date=newest_date, job_count=job_count)) if messages: data = { 'channel': ALERT_CHANNEL, 'text': ( '*RENEW SCRIPT ALERT* - {} ' 'store(s) appear to have ' 'stalled new jobs for over {} minutes:\n' '{}').format( len(messages), PROCESS_THRESHOLD_MIN, '\n'.join( messages))} slack_payload = json.dumps(data) curl = pycurl.Curl() curl.setopt( curl.HTTPHEADER, ['Content-Type: application/json']) curl.setopt(curl.URL, slack_url) curl.setopt(curl.POSTFIELDS, slack_payload) curl.setopt(curl.VERBOSE, True) curl.perform() if __name__ == '__main__': handler(None, None)