"""Jobs (encoding_queue_detail) model.""" from sqlalchemy import bindparam, text from sqlalchemy.exc import NoResultFound from vectororder.connectors.mysql import art_db_connector, dd_db_connector from vectororder.constants import META_UPDATE_TO_DELIVERY_TYPE, fields, sql from vectororder.exceptions import JobNotFoundError from vectororder.models.schemas import Job def get_jobs_duplicates(job_ids: list[int]) -> list[int]: """Get jobs which have duplicates. There are newer jobs with the same upc and store ID. Args: job_ids (list): encoding queue detail ID (int) list Returns: list: list of job IDs (int) """ with dd_db_connector.db_session() as session: result = session.execute( text(sql.DD_SELECT_EQD_DUPLICATES).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) try: rows = result.mappings().all() return [row[fields.EQ_DETAIL_ID] for row in rows] except NoResultFound: return [] def get_jobs(job_ids: list[int]) -> dict[int, Job]: """Get jobs statuses. Args: job_ids (list): encoding queue detail ID (int) list Returns: list: jobs """ with dd_db_connector.db_session() as session: result = session.execute( text(sql.DD_SELECT_JOB_STATUSES).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) try: rows = result.mappings().all() return { row[fields.EQ_DETAIL_ID]: Job( job_id=row[fields.EQ_DETAIL_ID], upc=row[fields.EQ_DETAIL_UPC], store_id=row[fields.EQ_DETAIL_STORE_ID], delivery_type=META_UPDATE_TO_DELIVERY_TYPE[ row[fields.ORDER_METADATA_UPDATE] ], encoding_order_type=row[fields.ORDER_ENCODING_TYPE], status=row[fields.EQ_DETAIL_STATUS], encoding_started=row[fields.EQ_DETAIL_ENCODING_STARTED], delivery_started=row[fields.EQ_DETAIL_DELIVERY_STARTED], ) for row in rows } except NoResultFound: return {} def set_reencode(job_ids: list[int]) -> None: """Update encoding_queue_detail table to re-encode. Args: job_ids (list): encoding queue detail ID (int) list """ with dd_db_connector.db_session(transaction=True) as session: session.execute( text(sql.DD_UPDATE_JOB_REENCODE).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) session.execute( text(sql.DD_DELETE_DELIVERY_BATCH_DETAIL).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) session.execute( text(sql.DD_DELETE_LOCATION).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) def set_redeliver(job_ids: list[int]) -> None: """Update encoding_queue_detail table to re-deliver. Args: job_ids (list): encoding queue detail ID (int) list """ with dd_db_connector.db_session(transaction=True) as session: session.execute( text(sql.DD_UPDATE_JOB_REDELIVER).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) session.execute( text(sql.DD_DELETE_DELIVERY_BATCH_DETAIL).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) session.execute( text(sql.DD_DELETE_LOCATION_IF_NEW).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) def set_cancel(job_ids: list[int]) -> None: """Update encoding_queue_detail table to cancel. Args: job_ids (list): encoding queue detail ID (int) list """ with dd_db_connector.db_session(transaction=True) as session: session.execute( text(sql.DD_UPDATE_JOB_CANCEL).bindparams( bindparam(fields.SQL_ID_LIST, expanding=True) ), {fields.SQL_ID_LIST: job_ids}, ) def get_job(job_id: int) -> Job: """Get job. Args: job_id (int) Returns: Job """ jobs = get_jobs([job_id]) if job_id in jobs: return jobs.pop(job_id) else: raise JobNotFoundError(f"Job with id {job_id} not found") def has_delivered_job(*, upc: int, store_id: int) -> bool: """Check if there is a delivered job with the given UPC. Args: upc (int) store_id (int) Returns: bool """ with dd_db_connector.db_session() as session: result = session.execute( text(sql.DD_GET_HAS_DELIVERED_JOB), { "upc": upc, "dms_master_master_id": store_id, }, ) return bool(result.scalar()) def has_delivery_history(*, upc: int, store_id: int) -> bool: """Check if there is a delivery history for the given UPC. Args: upc (int) store_id (int) Returns: bool """ with art_db_connector.db_session() as session: result = session.execute( text(sql.AR_GET_HAS_DELIVERY_HISTORY), { "upc": upc, "customer_master_master_id": store_id, }, ) return bool(result.scalar())