"""Multiprocessing mapper and Queues. Send requests to BulkUploader and process logging. """ import json from random import random from time import sleep from timeit import default_timer from constants.data_headers import data_headers from constants.globals import ( init_lock, keepalive_file ) from multiprocessing import Lock from pebble import ProcessPool, ProcessExpired import config import requests import util.db_utils as db_utils import util.logic_utils as ingestor from integration_scripts.file_utils import write_to_file, \ create_path_with_todays_date from integration_scripts import logger as log def _send_bulk_upload_request(params, post_params): """Send a request and parse the response.""" response = None requeue = False response_data = None # TODO Mark these rows as attempted # DEBUG - Log to files if config.DEBUG_WRITE_TO_FILE: output_path = create_path_with_todays_date(config.FILE_OUTPUT_PATH) file_name = 'request_{}_{}_{}.json'.format( params['orchard_vendor_id'], params['session_id'], params['request_num']) output_file = write_to_file( file_name, json.dumps(post_params, ensure_ascii=False), output_path) log.bind(multi=True).info( 'Request File #{} Written Successfully: ' '{}\n'.format(params['request_num'], output_file)) params['request'] = output_file try: # SEND BULK INGESTION REQUEST! response = requests.post(**post_params) except Exception as e: # TODO: urllib3 error - try to resubmit here? requeue = True response = { 'error': { 'msg': 'Request failed', 'bad_response': str(response), 'exception': str(e) } } log.bind(multi=True).error('Request Error Logged. Will Retry: {}', str(e)) # REQUEST COMPLETE log.bind(multi=True).info( 'Ingestion for {}, request #{} completed.\n'.format( params['orchard_vendor_id'], params['request_num'])) try: # Process response to dict response_data = ingestor.get_response_data(response) except Exception as e: requeue = True response = { 'error': { 'msg': 'Response is not JSON', 'bad_response': str(response), 'exception': str(e) } } log.bind(multi=True).error('Response Error Logged. Will Retry: {}', str(e)) if config.DEBUG_WRITE_TO_FILE: output_path = create_path_with_todays_date(config.FILE_OUTPUT_PATH) # Debug write to file output_file = 'response_{}_{}_{}.json'.format( params['orchard_vendor_id'], params['session_id'], params['request_num']) output_file = write_to_file( output_file, json.dumps(response_data, ensure_ascii=False), output_path) log.bind(multi=True).info( 'Response File Written Successfully: {}\n'.format( output_file)) return params, response_data, response, requeue def _fake_bulk_upload_request(params): """Provide response for fake request.""" response = None requeue = False from test_data import test_real_response # DEBUG - TEST DATA response_data = test_real_response() response_data = response_data['envelope']['response'] log.bind(multi=True).info('Skipping Request.') return params, response_data, response, requeue def _get_error_responses(params, label_chunk_request): """Create a list of responses, suitable for logging.""" response_error_releases = list() # Needed to ensure 1:1 release tracking in log for code, item in label_chunk_request['data'].items(): (project_code, product_code) = code.split('_____') params['project_code'] = project_code params['product_code'] = product_code response_error_releases.append(params.copy()) return response_error_releases def find_volume_and_track_count(data): """Return the number of volumes and tracks therein.""" vol_track_dict = dict() for key, tracks in data.items(): for track in tracks: if not vol_track_dict.get(track[data_headers['volume']]): vol_track_dict[track[data_headers['volume']]] = 0 vol_track_dict[track[data_headers['volume']]] += 1 return vol_track_dict def _log_bulk_upload_response(params, resp_data, label_request, is_requeue): """Parse out the response, and write to log tables.""" requeue = False # count tracks in release track_and_vol_count = find_volume_and_track_count(label_request['data']) # Failed responses if 'error' in resp_data: log.bind(multi=True).info( 'Logging error responses for {}, request {}.'.format( params['orchard_vendor_id'], params['request_num'])) # Clean response data params['response'] = resp_data # Create list of errors (only 1 when RELEASE_PER_REQUEST == True) response_error_releases = _get_error_responses( params, label_request) error_params = response_error_releases.pop() # Get the error data error_dict = resp_data['error'][list(resp_data['error'].keys())[0]] if 'data' in error_dict: # Create lists of DB bindings for inserts # track_log_params is a list of lists. release_log_params, track_log_params = \ ingestor.create_bulk_upload_response_log( resp_data, params, is_requeue, track_and_vol_count) db_utils.write_to_mysql_log( processed=release_log_params, tracks=track_log_params) else: # TODO - parse NOTICES like LEAD VOCALS issue (misspelled field) if not is_requeue: requeue = True if 'msg' in resp_data['error']: # TODO - Process NOTICES if resp_data['error']['msg'] == 'Malformed response': error_params['action_type'] = \ 'response error - malformed response' elif resp_data['error']['msg'] == 'Request failed': error_params['action_type'] = 'request failed' else: error_params['action_type'] = 'unknown request error' if config.DEBUG_LOG_RESPONSE_CONSOLE: log.bind(multi=True).warning( f'VAPI Response:\n{json.dumps(resp_data, indent=2)}' ) elif 'type' in resp_data['error']: if resp_data['error']['type'] == 'PDOException': error_params['action_type'] = 'SQL server error' if resp_data['error']['type'] == \ 'Web_Controller_Exception_InvalidInput': error_params['action_type'] = 'PHP server error' else: error_params['action_type'] = 'unknown server error' else: error_params['action_type'] = 'unknown error' if is_requeue: error_params['action_type'] += ' - requeue' db_utils.write_to_mysql_log(failed=[error_params]) # Good responses else: # Create lists of DB bindings for inserts # track_log_params is a list of lists. release_log_params, track_log_params = \ ingestor.create_bulk_upload_response_log( resp_data, params, is_requeue, track_and_vol_count) db_utils.write_to_mysql_log( processed=release_log_params, tracks=track_log_params) return requeue def _get_request_post_params(label_id, label_chunk_request): """Formulate POST parameters for request.""" bulk_upload_url = 'http{}://{}/bulkupload/post'.format( 's' if config.USE_SSL else '', config.BULK_UPLOAD_API_URL) request_params = { 'errorMode': '1', # 'XDEBUG_SESSION_START': 'devorch', 'access_token': label_chunk_request['access_token'], 'processMode': 'mass_ingest', 'vendorId': label_id, 'isMechadmin': 'True' } request_data = { 'data': json.dumps(label_chunk_request['data'], ensure_ascii=False), 'data_format': 'json' } post_params = { 'url': bulk_upload_url, 'data': request_data, 'params': request_params } return post_params def process_bulk_upload_request(request_num, session_id, name, label_request, label_id): """Process a single bulk upload submission. Main threaded action.""" msg = None error_msg = None # name = current_process().name = name params = dict() req_requeue = False log_requeue = False requeue_params = dict() log_params = dict() # Touch a file in the tmp folder if the request_num is divisible by 25 if request_num % 25 == 0: with open('tmp/{}'.format(keepalive_file), 'w') as f: f.write('') # request delay jitter sleep_for = random() log.bind(multi=True).info('Sleep jitter: {} secs.'.format(sleep_for)) sleep(sleep_for) is_requeue = True if 'requeue' in name else False # if config.RELEASE_PER_REQUEST: release_code = [ key for key, val in label_request['data'].items()] params['release_code'] = release_code.pop() # else: # params['release_code'] = None post_params = _get_request_post_params(label_id, label_request) # Request, session and label-level fields params['orchard_vendor_id'] = label_id if config.LOG_REQUEST: params['request'] = post_params['data'] else: params['request'] = '' params['session_id'] = session_id params['request_num'] = request_num log.bind(multi=True).info('{} starting'.format(name)) log.bind(multi=True).info( 'Now processing label {} request #{}'.format( label_id, request_num)) # For Release as input instead of label # if config.RELEASE_PER_REQUEST: log.bind(multi=True).info( 'Request release code: {}'.format(params['release_code'])) # ---- SEND REQUEST AND RECEIVE RESPONSE ------------------------------------ try: request_start_time = default_timer() log.bind(multi=True).info( 'Ingesting {}, request #{}.\n'.format(label_id, request_num)) # SEND REQUEST AND GET RESPONSE if config.MOCK_REQUESTS: params, response_data, response, req_requeue = \ _fake_bulk_upload_request(params) else: params, response_data, response, req_requeue = \ _send_bulk_upload_request(params, post_params) request_end_time = default_timer() request_time = request_end_time - request_start_time # Logging if not req_requeue: response_status = 'successfully' else: response_status = 'as failed' log.bind(multi=True).info( 'Label {}, request #{}: Returned {}. Duration: {}'.format( label_id, request_num, response_status, request_time)) except Exception as e: msg = '{} failed.'.format(name) # REQUEST COMPLETE - FAILED log.bind(multi=True).warning( 'Label {} request {} - Undefined error: {}'.format( label_id, request_num, str(e))) # for key, item in label_request['data'].items(): # log.bind(multi=True).log(25, '{}: {}'.format(name, key)) params['action_type'] = 'response error - uncaught exception' if is_requeue: params['action_type'] += ' - requeue' params['response'] = str(e) # get the error responses response_error_releases = _get_error_responses(params, label_request) # Write all the rows to log db_utils.insert_multi_mysql_release_logs( release_log_params=response_error_releases) error_msg = 'Response Error: {}'.format( json.dumps(e, ensure_ascii=False)) return msg, dict(), error_msg, request_num # ---- END REQUEST ----------------------------------------------------------- # ---- BEGIN LOGGING --------------------------------------------------------- if config.DEBUG_SKIP_LOG: log.bind(multi=True).info('Skipping Logging.') else: log.bind(multi=True).info( 'Logging ingestion responses for {}, request {}.'.format( label_id, request_num)) try: log_start_time = default_timer() # Log responses req_requeue = _log_bulk_upload_response( params, response_data, label_request, is_requeue) log_stop_time = default_timer() log_time = log_stop_time - log_start_time log.bind(multi=True).info( 'Logging completed for {}, request {}. Duration: {}'.format( label_id, request_num, log_time)) # stdout.flush() msg = '{} done.'.format(name) except Exception as e: log_requeue = True msg = '{} failed.'.format(name) # REQUEST COMPLETE - FAILED log.bind(multi=True).error('\nLogging error.\n') params['action_type'] = 'log error' if is_requeue: params['action_type'] += ' - requeue' params['response'] = str(e) response_error_releases = _get_error_responses( params, label_request) db_utils.insert_multi_mysql_release_logs( release_log_params=response_error_releases) error_msg = '\nLogging Error: {}'.format(str(e)) if req_requeue: requeue_params = { 'request_num': request_num, 'session_id': session_id, 'name': name, 'label_request': label_request, 'label_id': label_id } if log_requeue: log_params = { 'params': params, 'response_data': response_data, 'label_request': label_request } return msg, requeue_params, log_params, error_msg, request_num, label_id def map_processes(param_list): """Populate queues, and dequeue.""" results_again = list() failed_again = None # Run least threads necessary threads = min(len(param_list), config.NUM_THREADS) results = list() try: with ProcessPool( max_workers=threads, max_tasks=100, initializer=init_lock, initargs=(Lock(),)) as pool: for params in param_list: results.append( pool.schedule(process_bulk_upload_request, args=params)) except Exception as e: log.bind(multi=True).error(str(e)) return None, None, None requeue_param_list = list() log_param_list = list() failed_again = list() label_id = None success_count = 0 failure_count = 0 abnormal_count = 0 requeue_succeeded = 0 requeue_failed = 0 for res in results: try: result = res.result() except ProcessExpired as p: log.error('ExitCode {}: {}', str(p.exitcode), str(p)) abnormal_count += 1 continue except Exception as e: log.error(str(e)) continue if not label_id: try: label_id = result[5] except ValueError: label_id = param_list[0][4] if result[1]: requeue_param_list.append( (result[1]['request_num'], result[1]['session_id'], '{}-requeue'.format(result[1]['name']), result[1]['label_request'], result[1]['label_id']) ) failure_count += 1 else: success_count += 1 if result[2]: log_param_list.append(result[2].copy()) log.bind(multi=True).info('Multiprocess Report for {}:', label_id) log.bind(multi=True).info('{} results processed. ', len(results)) log.bind(multi=True).info('{} failed submissions. ', len(requeue_param_list)) log.bind(multi=True).info('{} failed logs.', len(log_param_list)) log.bind(multi=True).info('{} abnormal terminations.', abnormal_count) if config.REQUEUE_FAILED: if log_param_list: log.bind(multi=True).info( 'Attempting to submit failed logs again.') for log_params in log_param_list: try: _log_bulk_upload_response(**log_params) except Exception: failed_again.append(log_params) if failed_again: log.bind(multi=True).warning( 'The following logging failures have failed twice, and ' 'will not be retried.') for result in failed_again: log.bind(multi=True).warning(result) if requeue_param_list: log.bind(multi=True).info( 'Attempting to re-queue failed submissions.') try: with ProcessPool( max_workers=threads, max_tasks=100, initializer=init_lock, initargs=(Lock(),)) as pool: for params in requeue_param_list: results_again.append( pool.schedule(process_bulk_upload_request, params)) # Success fuse success = True for res in results_again: result = res.result() if result[1] or result[2]: success = False # Break fuse except Exception as e: success = False log.bind(multi=True).warning(str(e)) if success: requeue_succeeded += 1 log.bind(multi=True).info( 'All re-queued submissions succeeded.') else: requeue_failed += 1 log.bind(multi=True).warning('Re-queued submissions failed.') counts = { 'succeeded': success_count, 'failed': failure_count, 'requeue_succeeded': requeue_succeeded, 'requeue_failed': requeue_failed, 'abnormal': abnormal_count, } return results, results_again, failed_again, counts