import boto import boto.ses from datetime import datetime import logging import newrelic.agent from celery import Celery from raven import Client from raven.contrib.celery import register_signal from raven.contrib.celery import register_logger_signal from bulkperformancerights import celeryconfig from bulkperformancerights import config from bulkperformancerights.connectors import raven from bulkperformancerights.logic import converter from bulkperformancerights.logic import export_status from bulkperformancerights.logic import exporter from bulkperformancerights.logic import job from bulkperformancerights.logic import update from bulkperformancerights.logic import validate celery = Celery('bulkperformancerights', broker=config.celery_broker) celery.config_from_object(celeryconfig) if raven.sentry_dsn: client = Client(dsn=raven.sentry_dsn) # register a custom filter to filter out duplicate logs register_logger_signal(client) # hook into the Celery error handler register_signal(client) # The register_logger_signal function can also take an optional argument # `loglevel` which is the level used for the handler created. # Defaults to `logging.ERROR` # TODO (OrCharles): have this be configurable? register_logger_signal(client, loglevel=logging.INFO) @newrelic.agent.function_trace() @celery.task(queue=config.broker_queue_name) def convert_validate_update(job_status, user_s3_path, xlsx_user_s3_path, email): """Sends request to convert XLSX to JSON Args: job_status (JobStatus): job object with filename, user_type and user_id user_s3_path (str): path to xlsx file for user on s3 xlsx_user_s3_path (str): ows-xlsx service path for user email (str): email which will receive completion email """ # in the event any error is raised, catch, mark as error, rethrow try: job_dict = job_status.dict() file_name = job_dict.get('input_file_path') user_type = job_dict.get('user_type') user_id = job_dict.get('user_id') s3 = boto.connect_s3() bucket = s3.get_bucket(config.bucket_name) job.update_job(job_status, job_status='converting') convert_xlsx_response = converter.xlsx_to_json( bucket, file_name, user_type, user_id) assert convert_xlsx_response.status_code == 200, 'conversion failed' line_json_key = bucket.get_key(convert_xlsx_response.text) job.update_job(job_status, job_status='validating') valid_rows, invalid_rows = validate.validate_lines( s3_iter_lines(line_json_key), user_type, user_id) invalid_row_count = len(invalid_rows) valid_row_count = len(valid_rows) total_row_count = invalid_row_count + valid_row_count job.update_job(job_status, total_rows=total_row_count) if total_row_count == 0: job.update_job( job_status, job_status='errored', error_message='empty file') return if invalid_row_count: job.update_job(job_status, error_rows=invalid_row_count) json_error_key = boto.s3.key.Key(bucket) json_error_filename = '{}'.format(file_name.replace('.xlsx', '_error.json')) json_error_key_name = '{}json/{}{}'.format(config.xlsx_prefix, xlsx_user_s3_path, json_error_filename) json_error_key.key = json_error_key_name[1:] # Line delimited JSON sent to svc for error report json_error_key.set_contents_from_string('\n'.join(invalid_rows)) xlsx_error_key = '{}/{}{}'.format(config.prefix, user_s3_path, json_error_filename) json_error_key.copy(bucket, xlsx_error_key) # this file must live in ows-xlsx's s3 prefix exception_template = converter.exception_template # request JSON converted to XLSX, long poll for completion xlsx_error_filename = converter.json_to_xlsx( bucket, exception_template, json_error_filename, user_type, user_id) # copy file generated by ows-xlsx to /managerights prefix xlsx_error_key_name = '{}xlsx/{}{}'.format(config.xlsx_prefix, xlsx_user_s3_path, xlsx_error_filename) managerights_error_key_name = '{}{}{}'.format(config.prefix, user_s3_path, xlsx_error_filename) xlsx_error_key = bucket.get_key(xlsx_error_key_name[1:]) xlsx_error_key.copy(bucket, managerights_error_key_name) # TODO: run update and error report generation in parallel job.update_job( job_status, error_file_path=xlsx_error_filename) if valid_row_count: job.update_job(job_status, job_status='updating') updated_rows = update.update_rows(valid_rows) # Todo (OrCharles): rename processed_rows job.update_job(job_status, processed_rows=updated_rows) ses_conn = boto.ses.connect_to_region(config.AWS_REGION) job.send_complete_mail(ses_conn, user_id, email, user_type) job.update_job(job_status, job_status='processed', end_time=datetime.now()) except Exception as e: message = str(e.args[0]) job.update_job( job_status, job_status='errored', error_message=message, end_time=datetime.now()) raise @celery.task(queue=config.broker_queue_name) def export_tracks(export_job, email=None): """Export tracks to xlsx for rights mgmt updates Args: export_job (ExportStatus): export job state """ # in the event any error is raised, catch, mark as error, rethrow user_dict = {'type': export_job.user_type, 'id': export_job.user_id} try: export_status.amend(export_job, job_status='jsonify') row_count = exporter.tracks_json_to_s3(export_status.input_file_name, user_dict) export_status.amend(export_job, job_status='xlsxify', row_count=row_count) exporter.tracks_json_to_xlsx(export_status.input_file_name, user_dict) # TODO: use file path computed by service, or is inferring okay? xlsx_path = export_job.input_file_path.replace('json', 'xlsx') if email: ses_conn = boto.ses.connect_to_region(config.AWS_REGION) export_status.send_complete_mail(ses_conn, export_job.user_id, email, export_job.user_type) export_status.amend(export_job, job_status='exported', end_time=datetime.now(), output_file_path=xlsx_path) except Exception as e: # exception may not always have args so e.args[0] won't always work message = str(e) export_status.amend( export_job, job_status='errored', error_message=message, end_time=datetime.now()) raise def s3_iter_lines(key): """Stream data from s3 key Instead of downloading the entire file to disk, we can iterate Args: key (Key): AWS s3 Key from which to read data Return: Iterator See: Credit due: https://github.com/piskvorky/smart_open Todo (OrCharles): move this, rename logic.fileupload to fileops? """ line_buf = b'' for chunk in key: line_buf += chunk start = 0 while True: end = line_buf.find(b'\n', start) + 1 if end: yield line_buf[start: end] start = end else: line_buf = line_buf[start:] break if line_buf: yield line_buf