from datetime import datetime from oto import response from sqlalchemy import Column from sqlalchemy import DateTime from sqlalchemy import Enum from sqlalchemy import Integer from sqlalchemy import PickleType from sqlalchemy import String from sqlalchemy import Text from sqlalchemy.exc import DBAPIError from sqlalchemy.exc import DisconnectionError from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.exc import TimeoutError from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.sql import func from masters_registry import config from masters_registry.connectors import mysql from masters_registry.connectors import sentry from masters_registry.constant import bulk_tasks_const from masters_registry.constant import error BaseModel = declarative_base() class BigPickle(PickleType): impl = config.BLOB_TYPE def create_table(): """Creates task_status table in MySQL database """ BaseModel.metadata.create_all(mysql._bulk_statuses_engine) class TaskStatus(BaseModel): """Class representing the bulk task status.""" __tablename__ = 'task_status' id = Column(Integer, primary_key=True) correlation_id = Column(String(36)) user_id = Column(String(127)) user_name = Column(String(255), default='') type = Column(Enum(*bulk_tasks_const.BULK_TYPES)) status = Column( Enum(*bulk_tasks_const.BULK_STATUSES), default=bulk_tasks_const.PROCESSING_STATUS ) create_datetime = Column(DateTime, default=datetime.utcnow) finish_datetime = Column(DateTime, nullable=True, default=None) result = Column(BigPickle, nullable=True, default=None) count = Column(Integer, nullable=False) context = Column(Text, nullable=True) account_id = Column(Integer, nullable=True, default=None) account_type = Column(String(64), nullable=True, default=None) def as_dict(self, include_result=True): """Return object as dict Args: include_result (boolean): include or exclude `result` field in dict Returns: dict: Dictionary representation of object """ task_dict = { 'id': self.id, 'correlation_id': self.correlation_id, 'user_id': self.user_id, 'user_name': self.user_name, 'type': self.type, 'status': self.status, 'create_datetime': self.create_datetime.isoformat(), 'count': self.count, 'result': self.result if include_result else bool(self.result), 'account_type': self.account_type, 'account_id': self.account_id } if self.finish_datetime: task_dict['finish_datetime'] = self.finish_datetime.isoformat() else: task_dict['finish_datetime'] = None return task_dict def get_task_report(num_records, order_by, order_direction, page_offset=0): """Get a list of bulk processing tasks from task_status table. Args: num_records (int): number of task statuses to be returned order_by (str): create_datetime or finish_datetime order_direction (str): asc or desc page_offset (int): number of records to skip Returns: list(dict): list of tasks """ order_by_column = getattr(TaskStatus, order_by) finish_is_null = getattr( func, config.ISNULL_FUNC)(TaskStatus.finish_datetime) ordering = [getattr(order_by_column, order_direction)()] if order_by == 'finish_datetime': ordering.insert(0, getattr(finish_is_null, order_direction)()) with mysql.bulk_statuses_session_scope() as session: query = session.query(TaskStatus, finish_is_null) query = query.filter( TaskStatus.type != bulk_tasks_const.BULK_REMOVE_TERRITORIES) query = query.order_by( *ordering).offset(page_offset).limit(num_records) tasks = [task.as_dict(include_result=False) for task, _ in query.all()] return tasks def get_task_statuses_count(): """Function that returns total count of task statuses. Returns: int: count of task statuses """ with mysql.bulk_statuses_session_scope() as session: jobs_count_query = session.query(TaskStatus).filter( TaskStatus.type != bulk_tasks_const.BULK_REMOVE_TERRITORIES) jobs_count = jobs_count_query.count() return jobs_count def get_task(task_id): """Get single task from task_status table Args: task_id (int): task id Returns: TaskStatus """ try: with mysql.bulk_statuses_session_scope() as session: task = session.query(TaskStatus).filter_by(id=task_id).first() if task: session.expunge(task) return response.Response( message=task ) else: return response.create_not_found_response( message=error.BULK_TASK_DOES_NOT_EXIST ) except SQLAlchemyError as ex: return response.create_fatal_response(str(ex)) def create_task( correlation_id, user_id, task_type, count, context, user_name='', account_id=None, account_type=None): """Put new task item into task_status table. Args: correlation_id (str): correlation id user_id (int): user id task_type (str): task type count (int): amount of items in batch context (list): task context (list of ISRCs or UPCs) user_name (str): user name, optional parameter account_id (int): user account id, optional account_type (str): user account type, optional """ task_context = ','.join(context) if context else '' task = TaskStatus( correlation_id=correlation_id, user_id=user_id, user_name=user_name, type=task_type, count=count, context=task_context) if account_id and account_type: task.account_id = account_id task.account_type = account_type with mysql.bulk_statuses_session_scope() as session: session.add(task) session.flush() return task.id def update_task(task_id, status, result=None): """Updates task status Sets status, result and finished_timestamp Args: task_id (int): task id status (str): new status, should be DONE | PROCESSING | FAILED | PARTIAL_FAIL result (dict): result dict """ with mysql.bulk_statuses_session_scope() as session: try: task = session.query(TaskStatus).filter_by(id=task_id) task.update( { 'status': status, 'result': result, 'finish_datetime': datetime.utcnow() } ) return response.Response() except ( DBAPIError, DisconnectionError, SQLAlchemyError, TimeoutError ): if sentry.sentry_client: sentry.sentry_client.captureException() return response.create_error_response( code=error.INVALID_BULK_STATUS, message='Invalid bulk status', status=400 )