from datetime import datetime import json from uuid import uuid4 from oto import response from sqlalchemy import Column from sqlalchemy import Index from sqlalchemy import Integer from sqlalchemy import String from sqlalchemy import TIMESTAMP from sqlalchemy import UniqueConstraint from integration_scripts.connectors import mysql from integration_scripts.connectors import sentry class BulkUploadTracking(mysql.BaseModel): """Table definition for bulk_upload_tracking table.""" __tablename__ = 'BULK_UPLOAD_TRACKING' bulk_upload_id = Column(Integer, primary_key=True, autoincrement=True) orchard_vendor_id = Column(Integer, nullable=False) source_table = Column(String(255), nullable=False) source_file_name = Column(String(255), nullable=False) output_file_names = Column(String(8383), default='') session_id = Column(String(255), nullable=False) code = Column(String(10), nullable=True) environment = Column(String(24), default='dev') end_time = Column( TIMESTAMP, nullable=False, default=datetime.now()) start_time = Column(TIMESTAMP, default=datetime.now()) rows = Column(Integer, default=0) succeeded = Column(Integer, default=0) failed = Column(Integer, default=0) abnormal = Column(Integer, default=0) requeue_succeeded = Column(Integer, default=0) requeue_failed = Column(Integer, default=0) total_releases = Column(Integer, default=0) assets_ingested = Column(Integer, default=0) notes = Column(String(8383), default='') __table_args__ = ( Index( 'source_file_name_idx', 'source_file_name', 'session_id', 'orchard_vendor_id' ), Index( 'session_id_dx', 'session_id' ), UniqueConstraint( 'source_file_name', 'session_id', 'orchard_vendor_id', name='source_file_name_uk' ) ) def as_dict(self): """Return object as dict. Returns: dict: Dictionary representation of object """ bulk_upload_tracking_dict = { 'bulk_upload_id': self.bulk_upload_id, 'orchard_vendor_id': self.orchard_vendor_id, 'source_table': self.source_table, 'source_file_name': self.source_file_name, 'output_file_names': self.output_file_names.split('::') if self.output_file_names else [], 'session_id': self.session_id, 'code': self.code, 'environment': self.environment, 'start_time': self.start_time, 'end_time': self.end_time, 'rows': self.rows, 'failed': self.failed, 'succeeded': self.succeeded, 'abnormal': self.abnormal, 'requeue_succeeded': self.requeue_succeeded, 'requeue_failed': self.requeue_failed, 'total_releases': self.total_releases, 'assets_ingested': self.assets_ingested, 'notes': json.loads(self.notes) if self.notes else [] } return bulk_upload_tracking_dict @sentry.sentry_wrap def create_bulk_upload_tracking_table(): """Create Bulk Upload Tracking Tables""" BulkUploadTracking.metadata.create_all(mysql.rds_db_engine) @mysql.rds_db_session_wrap def create_bulk_upload_tracking_row(session, filename, vendor_id, source_table, session_id=None, code=None, environment='dev'): """Put new raw asset item into asset_upload table. Args: session (Session): An SQLAlchemy session filename (str): unique filename vendor_id (int): vendor id source_table (str): the table being used for the source of the ingest session_id (str): a uuid4 code (str): a shortened uuid4 environment (str): dev, qa, or prod Returns: response.Response: BulkUploadTracking.as_dict() in message attribute or error response. """ session_id = session_id or str(uuid4().hex) bulk_upload = BulkUploadTracking( source_file_name=filename, source_table=source_table, session_id=session_id, orchard_vendor_id=vendor_id, environment=environment, code=code ) session.add(bulk_upload) session.commit() session.refresh(bulk_upload) bulk_upload_dict = bulk_upload.as_dict() return response.Response(bulk_upload_dict) @mysql.rds_db_session_wrap def get_bulk_upload_rows(session, session_id, vendor_id=None, filename=None): filters = [ (BulkUploadTracking.session_id == session_id) ] if vendor_id: filters.extend([ (BulkUploadTracking.orchard_vendor_id == vendor_id) ]) if filename: filters.extend([ (BulkUploadTracking.source_file_name == filename) ]) bulk_upload_tracking = session.query( BulkUploadTracking).filter(*filters) bulk_upload_tracking_list = [but.as_dict() for but in bulk_upload_tracking] return response.Response(bulk_upload_tracking_list) @mysql.rds_db_session_wrap def get_bulk_upload_row( session, session_id, vendor_id, filename, as_dict=True): if not session_id and vendor_id and filename: bulk_upload_tracking = None else: filters = [ (BulkUploadTracking.session_id == session_id), (BulkUploadTracking.orchard_vendor_id == vendor_id), (BulkUploadTracking.source_file_name == filename) ] bulk_upload_tracking = session.query( BulkUploadTracking).filter(*filters).first() if bulk_upload_tracking is None: return response.create_not_found_response( 'bulk_upload_tracking_not_found') if as_dict: bulk_upload_tracking = bulk_upload_tracking.as_dict() return response.Response(bulk_upload_tracking) @mysql.rds_db_session_wrap def get_bulk_upload_row_by_id(session, bulk_upload_id, as_dict=True): filters = [ (BulkUploadTracking.bulk_upload_id == bulk_upload_id) ] bulk_upload_tracking = session.query( BulkUploadTracking).filter(*filters).first() if bulk_upload_tracking is None: return response.create_not_found_response( 'bulk_upload_tracking_not_found') if as_dict: bulk_upload_tracking = bulk_upload_tracking.as_dict() return response.Response(bulk_upload_tracking) @mysql.rds_db_session_wrap def update_bulk_upload_row(session, **kwargs): if 'bulk_upload_id' in kwargs.keys(): bulk_upload_tracking = get_bulk_upload_row_by_id( session=session, bulk_upload_id=kwargs.get('bulk_upload_id'), as_dict=False ).message else: bulk_upload_tracking = get_bulk_upload_row( session=session, session_id=kwargs.get('session_id'), vendor_id=kwargs.get('orchard_vendor_id'), filename=kwargs.get('filename'), as_dict=False, ).message if not bulk_upload_tracking: return bulk_upload_tracking for attr, val in kwargs.items(): if attr in ['notes', 'output_file_names'] and type(val) is list: val = json.dumps(val) if val else '' if attr != 'bulk_upload_id': setattr(bulk_upload_tracking, attr, val) setattr(bulk_upload_tracking, 'end_time', datetime.now()) session.merge(bulk_upload_tracking) session.commit() session.refresh(bulk_upload_tracking) bulk_upload_tracking_dict = bulk_upload_tracking.as_dict() return response.Response(bulk_upload_tracking_dict)