"""Ready Deliveries Lambda module.""" import base64 import math import re import time import tempfile import traceback from vector_utils.connections import connection_info from vector_utils.connections import exceptions from vector_utils.connections.transporter import Transporter from vector_utils.utils import parsers from vector_utils.utils import write_file_to_disk from vector_utils import queries import config from config import vector_secrets_manager_client from logger import get_current_logger from src.datadog import publish_datadog_metric class InvocationContext: """This is for passing around the invoke context to logger.""" correlation_id = None def handler(event, context): """Entrypoint to Lambda function. Args: event (optional): AWS Lambda event dependent structure with metadata. context (LambdaContext): AWS Lambda context. """ InvocationContext.correlation_id = context.aws_request_id log = get_current_logger(InvocationContext.correlation_id) try: allowed_encoders = event.get('allowed_encoders') or [17, 18, 19, 23] dms_id = event['dms_id'] order_type = event['order_type'] batch_delivery = event['batch_delivery'] dms_delivery_spec = queries.get_dms_delivery_spec( dms_id, order_type, conn_info=config.DD_MYSQL_CONN_INFO) log.info('Running "{}" for dms_id {}'.format(config.LAMBDA_NAME, dms_id)) publish_datadog_metric(dms_id, 'ready-deliveries.attempt') encoded_jobs_info = queries.get_encoded_jobs_info( dms_id, allowed_encoders, order_type, conn_info=config.DD_MYSQL_CONN_INFO) if encoded_jobs_info: if batch_delivery == 'Y': create_batches( int(encoded_jobs_info['dms_master_master_id']), order_type, int(encoded_jobs_info['encoded_jobs']), int(encoded_jobs_info['highest_priority']), dms_delivery_spec) else: ready_non_batched_jobs(dms_id, order_type) else: log.info('No encoded jobs picked up') publish_datadog_metric(dms_id, 'ready-deliveries.run') log.info('Finished "{}" for dms_id {}'.format(config.LAMBDA_NAME, dms_id)) except Exception as e: log.error('Unexpected Error: {}'.format('\n'.join(traceback.format_exception(e)))) raise e def create_batches( dms_id, order_type, encoded_jobs, highest_priority, dms_delivery_spec): """Process batch delivery jobs. Args: dms_id (int): Store ID order_type (string): Order type encoded_jobs (int): Number of encoded jobs highest_priority (int): Highest job priority dms_delivery_spec (dict): DMS delivery spec """ log = get_current_logger(InvocationContext.correlation_id) result = queries.get_open_batch_count( dms_id, order_type, conn_info=config.DD_MYSQL_CONN_INFO) num_batches = dms_delivery_spec.get('max_num_batches') if result['open_batches'] > 0: num_batches -= result['open_batches'] possible_batches = math.ceil( encoded_jobs / dms_delivery_spec.get('num_per_batch')) num_batches = min(possible_batches, num_batches) if num_batches < 1 and result['highest_priority'] > highest_priority: num_batches += 1 if num_batches < 1: log.info( 'DMS ID: {} - No more batches can be created.'.format(dms_id)) return None secrets_path = f'connection_info/{dms_id}/{order_type}' connection_data = vector_secrets_manager_client.get_cred(secrets_path) if 'private_key' in connection_data: connection_data['priv_key'] = f'{tempfile.gettempdir()}/{dms_id}_{order_type}_key' write_file_to_disk( tempfile.gettempdir(), connection_data['priv_key'], base64.b64decode(connection_data['private_key']).decode()) if ( connection_data.get('connection_type') == 'sftp' and str(dms_id) in config.STORES_WITH_DISABLED_ALGORITHMS ): log.info('Set sftp_disabled_algorithms for: {}'.format(dms_id)) connection_data['sftp_disabled_algorithms'] = { 'pubkeys': ['rsa-sha2-256', 'rsa-sha2-512']} if connection_data.get('connection_type') == 'gcs' and 'gcs_config_file' in connection_data: gcs_config = connection_data['gcs_config_file'] if isinstance(gcs_config, dict) and 'private_key' in gcs_config: try: decoded_key = base64.b64decode(gcs_config['private_key']).decode() gcs_config['private_key'] = decoded_key connection_data['gcs_config_file'] = gcs_config except Exception as e: log.error(f'Failed to decode GCS private key for {dms_id}: {e}') raise ValueError(f'Invalid base64 private key for GCS config (dms_id={dms_id}): {e}') conn_obj = connection_info.ConnectionInfo(connection_data, order_type) try: transporter = Transporter(conn_obj, logger=log) except (TimeoutError, ConnectionResetError, exceptions.ConnectionTimeout, exceptions.SSHException) as err: publish_datadog_metric(dms_id, 'ready-deliveries.failure', ['exception:connection_timeout']) log.warning(f'Error {err} for DMS {dms_id} {traceback.format_exc()}') return remote_initial_dir = connection_data.get('remote_initial_dir') if dms_delivery_spec.get('connection_type') == 's3': remote_initial_dir = remote_initial_dir.lstrip('/') while num_batches > 0: try: remote_folder = create_remote_folder( transporter, remote_initial_dir, dms_delivery_spec.get('batch_foldername')) except (exceptions.ConnectionTimeout, exceptions.SSHException, ValueError) as err: log.warning(f'Error: {err} DMS ID: {dms_id} - remote folder creation') break log.info(f'DMS ID: {dms_id} - created batch: {remote_folder}') result = queries.create_delivery_batch( dms_id, remote_folder, order_type, dms_delivery_spec.get('num_per_batch'), conn_info=config.DD_MYSQL_CONN_INFO) if not result: log.warning(f'Error: No Result DMS ID: {dms_id} - {remote_folder} batch creation error') break num_batches -= 1 try: transporter.close_connection() except Exception: pass def ready_non_batched_jobs(dms_id, order_type): """Process non batch delivery jobs. Args: dms_id (int): Store ID order_type (string): Order Type """ queries.set_job_for_delivery( dms_id, order_type, conn_info=config.DD_MYSQL_CONN_INFO) def create_remote_folder(transporter, remote_initial_dir, batch_foldername): """Create remote folder. Args: transporter (Transporter): Transporter class remote_initial_dir (string): Home folder of the remote user batch_foldername (string): Batch folder name Returns: remote_folder (string): Remote batch folder name """ counter = 1 attempts = set() replace_vars = None try: max_counter = 100 matches = re.search(r'{counter[()1-9]{0,3}}', batch_foldername) if matches: str_counter = parsers.process_string(matches[0], {'counter': counter}) max_counter = int('9' * len(str_counter)) # retry until max counter while counter <= max_counter: # generate remote folder name (counter vs time for uniqueness) if not matches or not replace_vars: replace_vars = { 'Year': time.strftime('%Y'), 'year': time.strftime('%y'), 'month': time.strftime('%m'), 'day': time.strftime('%d'), 'hour': time.strftime('%H'), 'minute': time.strftime('%M'), 'second': time.strftime('%S'), 'unix_timestamp': int(time.time()) } replace_vars['counter'] = counter remote_folder = parsers.process_string( batch_foldername, replace_vars ) counter += 1 # check if folder creation was already attempted if remote_folder in attempts: time.sleep(1) continue attempts.add(remote_folder) # try to create directory if transporter.mkdir('{}/{}'.format(remote_initial_dir, remote_folder)): return remote_folder except Exception as e: raise ValueError(e) # all attempts exhausted raise ValueError('Folder creation failed.') if __name__ == '__main__': handler(None, None)