"""AssetStatus Model.""" import json from sqlalchemy import and_ from sqlalchemy import Column from sqlalchemy import DateTime from sqlalchemy import ForeignKey from sqlalchemy import func from sqlalchemy import Integer from sqlalchemy import String from sqlalchemy import Text from sqlalchemy.orm import aliased from asset_transcoder.connectors import mysql from asset_transcoder.constants import error from asset_transcoder.utils.exceptions import OwsError class AssetStatus(mysql.BaseModel): """Table definition for asset_status table.""" __tablename__ = 'asset_status' asset_status_id = Column( 'id', Integer, primary_key=True, autoincrement=True) asset_upload_id = Column(Integer, ForeignKey('asset_upload.id')) status = Column(String(50), nullable=False) errors = Column(Text, nullable=True) status_time = Column(DateTime, nullable=True) created_date = Column(DateTime, default=func.now()) def to_dict(self): """Return object as dict. Returns: dict: Dictionary representation of object """ return { 'id': self.asset_status_id, 'asset_upload_id': self.asset_upload_id, 'status': self.status, 'status_time': self.status_time, 'created_date': self.created_date, 'errors': '|'.join(list(json.loads(self.errors).keys())) if self.errors else None, 'assets': [] } def create_asset_status(asset_upload_id, status, created_date, errors, session=None): """Create new asset status item in the table. Args: asset_upload_id (int): Asset upload id. status (str): Asset status. created_date (datetime): Asset status time. errors (dict): Errors if present. session (sqlalchemy.Session): db connection Returns: dict: Created status info. """ asset_status = AssetStatus( asset_upload_id=asset_upload_id, status=status, errors=json.dumps(errors) if errors else None, status_time=created_date, created_date=created_date) if session is None: with mysql.transcoder_db_session() as session: session.add(asset_status) else: session.add(asset_status) session.flush() return asset_status.to_dict() def get_last_asset_status(asset_upload_id): """Return last asset status by asset upload id. Args: asset_upload_id (int): Asset upload id. Returns: dict: Last asset status info or error """ filters = [ (AssetStatus.asset_upload_id == asset_upload_id)] with mysql.transcoder_db_session(read_only=True) as session: asset_status = session.query(AssetStatus).filter( *filters).order_by(AssetStatus.status_time.desc(), AssetStatus.asset_status_id.desc()).first() if not asset_status: raise OwsError.not_found(error.ERROR_ASSET_STATUS_NOT_FOUND) return asset_status.to_dict() def get_last_status_by_ids(asset_upload_ids, session): """Return last asset status by asset upload id. Args: asset_upload_ids (list): Asset upload id list. session (sqlalchemy.Session): live db connection. Returns: dict: Last asset status info or error """ # Subquery to get the latest asset_status_id for each asset_upload_id latest_status_subquery = session.query( AssetStatus.asset_upload_id, func.max(AssetStatus.asset_status_id).label('latest_asset_status_id') ).group_by(AssetStatus.asset_upload_id).subquery() # Aliased AssetStatus for joining with subquery latest_status = aliased(AssetStatus) # Main query to get the latest AssetStatus records asset_statuses = session.query(latest_status).join( latest_status_subquery, and_( latest_status.asset_upload_id == latest_status_subquery.c.asset_upload_id, latest_status.asset_status_id == latest_status_subquery.c.latest_asset_status_id ) ).filter( latest_status.asset_upload_id.in_(asset_upload_ids) ).order_by( latest_status.status_time.desc(), latest_status.asset_status_id.desc() ).all() return {'items': [asset_status.to_dict() for asset_status in asset_statuses]}