"""Direct Delivery alert Lambda module.""" import base64 import json import logging import urllib import boto3 import pycurl import requests import util RENEW_CYCLE = 10 DELIVERY_THRESHOLD_MIN = 90 MAX_PRIORITY_ALL_STORES = 1000 ALERT_CHANNEL = '#vector-system-alerts' OA_REPORT_URL_FORMAT = ( 'https://oa.theorchard.com/oa/warehouse/direct_delivery_detail.php?' 'encoding_order_id=&encoder=&dd_status[]=wait&dd_status[]=process' '&error_types[]=encoding&error_types[]=delivering&error_types[]=new' '&error_types[]=ready_to_encode&error_types[]=ready_to_deliver' '&error_types[]=encoded&error_types[]=queued_for_delivery&mkt_priority=' '&hd_dd_dms[]={dms_master_master_id}&from_encoding_started=&' 'to_encoding_started=&from_encoding_ended=&to_encoding_ended=&' 'meta_update=&err_log=&eo_priority=&upcs=&from_delivery_started=&' 'to_delivery_started=&from_delivery_ended=&to_delivery_ended=&' 'type=search') 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 get_shortened_url(url): """Use a URL shortener service to compact url. Args: url (str): URL to be shortened. Returns: str: a URL that redirects to original url. """ encoded_url = urllib.quote_plus(url) response = requests.get( 'http://tinyurl.com/api-create.php?url={url}'.format(url=encoded_url)) return response.text if response else 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. """ logger.info( 'Will scan {max_priority} store priorities.'.format( max_priority=MAX_PRIORITY_ALL_STORES)) with util.reportsdd_connection( db_host, name, password, db_name) as conn: with conn.cursor() as cursor: cursor.execute(''' SELECT cmm.customer_name AS store_name, cmm.priority, cmm.customer_master_master_id, a.delivery_ended AS last_delivery, b.priority AS backlog_priority, b.cnt, b.last_updated AS wait_time FROM customer_master_master cmm INNER JOIN dms_delivery_spec dds ON dds.dms_master_master_id = cmm.customer_master_master_id INNER JOIN ( SELECT MAX(eqd.delivery_ended) AS delivery_ended, eqd.dms_master_master_id, eq.encoding_order_type FROM encoding_queue_detail eqd INNER JOIN encoding_queue eq ON eq.encoding_queue_id = eqd.encoding_queue_id WHERE eqd.status = 'delivered' AND eq.encoder_id = 18 -- audio only GROUP BY eqd.dms_master_master_id ) a ON a.dms_master_master_id = cmm.customer_master_master_id AND a.encoding_order_type = dds.order_type INNER JOIN ( SELECT COUNT(eqd.encoding_queue_detail_id) AS cnt, MIN(eqd.last_updated) AS last_updated, eqd.dms_master_master_id, eq.encoding_order_type, eq.priority FROM encoding_queue_detail eqd INNER JOIN encoding_queue eq ON eq.encoding_queue_id = eqd.encoding_queue_id WHERE eqd.status IN ('new', 'ready_to_encode', 'encoding', 'ready_to_deliver', 'delivering', 'encoded', 'queued_for_delivery') AND eq.encoder_id = 18 -- audio only GROUP BY eqd.dms_master_master_id ) b ON b.dms_master_master_id = cmm.customer_master_master_id AND b.encoding_order_type = dds.order_type WHERE dds.order_type = 'release' AND dds.delivery = 'Y' AND dds.encoding = 'Y' GROUP BY cmm.customer_master_master_id HAVING TIMESTAMPDIFF(MINUTE, last_delivery, NOW()) > %s AND cnt > 0 AND TIMESTAMPDIFF(MINUTE, wait_time, NOW()) > %s ORDER BY cmm.priority LIMIT %s ''', ( DELIVERY_THRESHOLD_MIN, RENEW_CYCLE, MAX_PRIORITY_ALL_STORES)) stats = [row for row in cursor.fetchall()] logger.info('Stats found: {stats}'.format(stats=stats)) if len(stats): messages = [] for ( store_name, priority, backlog_priority, dms_master_master_id, cnt, wait_time, last_delivery) in stats: messages.append( 'Last delivery at *{last_delivery}* to *{store_name}* ' '(DMS {dms_master_master_id}, priority {priority}): ' '*{cnt}* in backlog with highest EO priority of ' '{backlog_priority} {oa_report_url}'.format( last_delivery=last_delivery, store_name=store_name, dms_master_master_id=dms_master_master_id, cnt=cnt, priority=priority, backlog_priority=backlog_priority, oa_report_url=get_shortened_url( OA_REPORT_URL_FORMAT.format( dms_master_master_id=dms_master_master_id) ))) data = { 'channel': ALERT_CHANNEL, 'text': ( '{} of top {} stores appear to have stalled *audio* ' 'deliveries for over {} minutes:\n{}').format( len(stats), MAX_PRIORITY_ALL_STORES, DELIVERY_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)