from utils.db_connectors.base_connector import BaseConnector from utils.db_objects.direct_delivery import EncodingQueue, \ EncodingQueueDetail, JobPriority, CustomerMasterMaster, DMSDeliverySpec class EncodingQueueConnector(BaseConnector): def delete_encoding_queue_data_if_exists(self, encoding_order_id): row = self.get_row_by_id(EncodingQueue, 'encoding_order_id', encoding_order_id) if row: self.delete_encoding_queue_data(row.encoding_queue_id) else: self.log.info(f'No row in EncodingQueue for encoding_order_id ' f'{encoding_order_id} found, no need to delete ' f'existing data') return def delete_encoding_queue_data(self, eq_id): session = self.Session() detail_rows = session.query(EncodingQueueDetail).filter_by( encoding_queue_id=eq_id) for row in detail_rows: detail_id = row.encoding_queue_detail_id session.query(JobPriority).filter_by( encoding_queue_detail_id=detail_id).delete() session.delete(row) session.query(EncodingQueue).filter_by( encoding_queue_id=eq_id).delete() session.commit() session.close() def assert_direct_delivery_order(self, encoding_order, product_id, store_id): row = self.get_row_by_id(EncodingQueue, 'encoding_order_id', encoding_order.encoding_order_id) assert row.priority == encoding_order.priority assert row.encoding_order_type == 'release' assert row.encoder_id == encoding_order.encoder_id assert row.meta_update == encoding_order.meta_update self.assert_encoding_queue_detail(row.encoding_queue_id, product_id, store_id) def assert_encoding_queue_detail(self, eqd_id, product_id, store_id): row = self.get_row_by_id(EncodingQueueDetail, 'encoding_queue_id', eqd_id) assert row.upc == product_id assert row.status == 'new' assert row.dms_master_master_id == store_id assert row.updated_oa == 'N' def update_queue_detail_status(self, detail_id, status): session = self.Session() try: session.query(EncodingQueueDetail).filter( EncodingQueueDetail.encoding_queue_detail_id == detail_id ).update( {"status": status} ) session.commit() except Exception as e: self.log.debug("Updating EncodingQueueDetail failed:", e) session.rollback() raise session.close() def get_encoding_queue_detail_for_deliver(self, dms_id, statuses, rows=10): session = self.Session() query = (session.query( EncodingQueueDetail.encoding_queue_detail_id, JobPriority.priority.label("encoding_order_priority"), EncodingQueueDetail.dms_master_master_id, EncodingQueue.encoding_order_type, EncodingQueue.encoder_id, CustomerMasterMaster.priority.label("dms_priority") ).join(EncodingQueue, EncodingQueueDetail.encoding_queue_id == EncodingQueue.encoding_queue_id ).join( CustomerMasterMaster, EncodingQueueDetail.dms_master_master_id == CustomerMasterMaster.customer_master_master_id ).join( JobPriority, EncodingQueueDetail.encoding_queue_detail_id == JobPriority.encoding_queue_detail_id ).join( DMSDeliverySpec, (DMSDeliverySpec.dms_master_master_id == CustomerMasterMaster.customer_master_master_id) & (EncodingQueue.encoding_order_type == DMSDeliverySpec.order_type) & (DMSDeliverySpec.batch_delivery == "N") & (DMSDeliverySpec.delivery == "Y") ).filter( EncodingQueueDetail.status.in_(statuses), CustomerMasterMaster.customer_master_master_id == dms_id, EncodingQueue.encoder_id.in_([18]) ).limit(rows)) results = query.all() session.close() return results def assert_encoding_queue_detail_status(self, eqd_id, status): row = self.get_row_by_id(EncodingQueueDetail, 'encoding_queue_detail_id', eqd_id) assert row.status == status, \ 'Status was {},expected {}'.format(row.status, status) def get_detail_id_for_ready_deliveries(self, dms_id, order_type, encoder_id): session = self.Session() results = session.query( EncodingQueueDetail.encoding_queue_detail_id ).join(EncodingQueue, EncodingQueue.encoding_queue_id == EncodingQueueDetail.encoding_queue_id).filter( EncodingQueue.encoder_id == encoder_id, EncodingQueue.encoding_order_type == order_type, EncodingQueueDetail.dms_master_master_id == dms_id ).limit(1).all() session.close() return results[0][0]