"""Review Queue Model.""" from collections import namedtuple from datetime import datetime, timedelta, timezone from flask import g from owsresponse import status as owsresponse_status from sqlalchemy import and_, or_ from sqlalchemy.orm import joinedload, relationship from sqlalchemy.sql.expression import desc from product_review.api import db from product_review.constants.error import ( ERROR_CODE_INVALID_STATE, ERROR_CODE_NOT_FOUND, ERROR_CODE_UNABLE_TO_LOCK, ERROR_CODE_UNABLE_TO_UNLOCK, ERROR_MESSAGE_ALREADY_COMPLETE, ERROR_MESSAGE_NOT_FOUND, ERROR_MESSAGE_UNABLE_TO_LOCK, ERROR_MESSAGE_UNABLE_TO_UNLOCK, ) from product_review.models import approval as approval_model from product_review.models import move_target as move_target_model # noqa: F401 from product_review.models import queue_move_history as queue_move_history_model from product_review.models import rejection as rejection_model from product_review.util.exception import raise_exception_for_ows_response __STATUS = ("new", "complete", "applying") STATUS = namedtuple( "Status", (s.upper() for s in __STATUS), defaults=(s.lower() for s in __STATUS) )() STATUS_CHANGE = { STATUS.NEW: {STATUS.APPLYING, STATUS.COMPLETE}, STATUS.APPLYING: {STATUS.COMPLETE}, } DEFAULT_PAGE_LIMIT = 50 DEFAULT_STATUS = STATUS.NEW DEFAULT_QUEUE_NAME = "initial" class ReviewQueue(db.Model): """Review queue model.""" __tablename__ = "review_queue" review_queue_id = db.Column( "id", db.Integer, primary_key=True, autoincrement=True, nullable=False ) product_id = db.Column(db.Integer, nullable=False) error_correction_id = db.Column(db.Integer, default=None) status = db.Column(db.Enum(*STATUS), default=DEFAULT_STATUS, nullable=False) submission_type = db.Column(db.String, default=None) submission_count = db.Column(db.Integer, default=1) queue_name = db.Column(db.String, nullable=False, default=DEFAULT_QUEUE_NAME) moved_to_queue_by_user_id = db.Column(db.String, default=None) submitted_by_user_id = db.Column(db.String, default=None) created_datetime = db.Column(db.DateTime, nullable=False, default=datetime.utcnow) locked_by_user_id = db.Column(db.String, default=None) locked_at_datetime = db.Column(db.DateTime, default=None) locked_until_datetime = db.Column(db.DateTime, default=None) moved_to_target_id = db.Column( db.Integer, db.ForeignKey("move_target.id"), default=None, ) moved_to_target_name = db.Column(db.String(80), default=None) moved_to_target_email = db.Column(db.String(225), default=None) moved_to_orchadmin_user_id = db.Column( db.Integer, default=None, comment="orchadmin_user ID of move target." ) move_note = db.Column( db.String, default=None, comment="Note explaining why the product was moved.", ) moved_at_datetime = db.Column(db.DateTime, default=None) submission_datetime = db.Column(db.DateTime, default=None) approval = relationship( "Approval", foreign_keys=[approval_model.Approval.review_queue_id], primaryjoin="ReviewQueue.review_queue_id==Approval.review_queue_id", uselist=False, lazy="select", ) rejection = relationship( "Rejection", foreign_keys=[rejection_model.Rejection.review_queue_id], primaryjoin="ReviewQueue.review_queue_id==Rejection.review_queue_id", uselist=False, lazy="select", ) queue_move_history = relationship( "QueueMoveHistory", foreign_keys=[queue_move_history_model.QueueMoveHistory.review_queue_id], primaryjoin="ReviewQueue.review_queue_id==QueueMoveHistory.review_queue_id", uselist=True, lazy="select", order_by="desc(QueueMoveHistory.moved_at_datetime)" ) score_history = relationship( "ScoreHistory", foreign_keys="[ScoreHistory.review_queue_id]", primaryjoin="ReviewQueue.review_queue_id==ScoreHistory.review_queue_id", uselist=True, lazy="select", order_by="desc(ScoreHistory.score_version)" ) def to_dict(self, include_context=False, include_history=False, include_score=False): """Return a dictionary of output properties.""" item = {} for c in self.__table__.columns: column_name = "review_queue_id" if c.name == "id" else c.name item[column_name] = getattr(self, column_name) if include_context: item["approval"] = self.approval.to_dict() if self.approval else None item["rejection"] = self.rejection.to_dict() if self.rejection else None if include_history: item["queue_move_history"] = ( [ entry.to_dict() for entry in self.queue_move_history ] if self.queue_move_history else None ) if include_score: latest_score = self.score_history[0] if self.score_history else None item["score_history"] = latest_score.to_dict() if latest_score else None return item @classmethod def latest_review_by_product(cls, product_id): """Get the current status.""" results = ( cls.query.options(joinedload(cls.approval)) .options(joinedload(cls.rejection)) .join(cls.approval, isouter=True) .join(cls.rejection, isouter=True) .filter(cls.product_id == product_id) .filter( or_( cls.status == STATUS.NEW, cls.status == STATUS.APPLYING, and_( cls.status == STATUS.COMPLETE, or_( approval_model.Approval.approval_id.is_not(None), rejection_model.Rejection.rejection_id.is_not(None), ), ), ) ) # filter out unsubmissions for status .order_by(desc(cls.created_datetime)) .limit(1) .all() ) if results: return results[0] def get_item_for_update(review_queue_id): """Return a review queue item while holding a row lock.""" return ( db.session.query(ReviewQueue) .filter(ReviewQueue.review_queue_id == review_queue_id) .with_for_update() .one_or_none() ) def get_items( status=None, page_offset=None, page_limit=None, review_queue_id=None, product_id=None, ): """Return list of items by status, sorted by created_datetime.""" # flask-sqlalchemy defaults are page = 1 and page_limit = 20 offset = int(page_offset or 0) limit = int(page_limit or DEFAULT_PAGE_LIMIT) review_queue_query = db.session.query(ReviewQueue) if status: if isinstance(status, str): review_queue_query = review_queue_query.filter(ReviewQueue.status == status) elif isinstance(status, tuple): review_queue_query = review_queue_query.filter( ReviewQueue.status.in_(status) ) else: raise Exception(f"Unsupported lookup for type: {type(status)}") if review_queue_id: review_queue_query = review_queue_query.filter( ReviewQueue.review_queue_id == review_queue_id ) if product_id: review_queue_query = review_queue_query.filter( ReviewQueue.product_id == product_id ) total_count = review_queue_query.count() results = ( review_queue_query.order_by(ReviewQueue.created_datetime.desc()) .limit(limit) .offset(offset) .all() ) return results, total_count def get_review_queue_items_by_id(review_queue_ids): """Get review queue items.""" with db.session(): review_queue_query = ( ReviewQueue.query.options(joinedload(ReviewQueue.approval)) .options(joinedload(ReviewQueue.rejection)) .join(ReviewQueue.approval, isouter=True) .join(ReviewQueue.rejection, isouter=True) .filter(ReviewQueue.review_queue_id.in_(review_queue_ids)) ) results = [] for item in review_queue_query: results.append(item.to_dict(include_context=True, include_score=True)) return results def get_target_group_by_item_id(review_queue_ids): """Get target group id off of review queue item by id.""" with db.session() as session: query = session.query( ReviewQueue.review_queue_id, ReviewQueue.moved_to_target_id ).filter(ReviewQueue.review_queue_id.in_(review_queue_ids)) results = [] for item in query: ids = {"review_queue_id": item[0], "moved_to_target_id": item[1]} results.append(ids) return results def review_history(product_id, page_limit, page_offset, current_item=False): """Return list of review queue items for a product ID, sorted by created_datetime.""" # noqa: E501 offset = int(page_offset or 0) limit = int(page_limit or DEFAULT_PAGE_LIMIT) review_queue_query = ( ReviewQueue.query.options(joinedload(ReviewQueue.approval)) .options(joinedload(ReviewQueue.rejection)) .join(ReviewQueue.approval, isouter=True) .join(ReviewQueue.rejection, isouter=True) .join(ReviewQueue.queue_move_history, isouter=True) .filter( ReviewQueue.product_id == product_id )).filter( or_( and_(current_item is True, ReviewQueue.status == STATUS.NEW), and_(ReviewQueue.status == STATUS.COMPLETE, or_( approval_model.Approval.approval_id.is_not(None), rejection_model.Rejection.rejection_id.is_not(None), )) ) ) total_count = review_queue_query.count() results = ( review_queue_query.order_by(ReviewQueue.created_datetime.desc()) .limit(limit) .offset(offset) .all() ) return results, total_count def change_status(review_queue_item, status): """Change the status of a review queue item.""" status_change = { **STATUS_CHANGE, } if getattr(g.request_context, "is_integration_test_preparation", False) or getattr( g.request_context, "is_db_resync", False ): # allow new => new transition for convenience in case record is already reset status_change[STATUS.NEW] = {*STATUS_CHANGE[STATUS.NEW], STATUS.NEW} status_change[STATUS.COMPLETE] = {STATUS.NEW} if status not in status_change.get(review_queue_item.status, {}): raise Exception( f'Status change "{review_queue_item.status} => {status}" is not valid' ) review_queue_item.status = status def lock(review_queue_id): """Lock a review queue item.""" locked_at_datetime = datetime.now(timezone.utc) row_to_lock = ( ReviewQueue.query.filter(ReviewQueue.review_queue_id == review_queue_id) .filter( and_( ReviewQueue.status == STATUS.NEW, or_( ReviewQueue.locked_by_user_id.is_(None), ReviewQueue.locked_by_user_id == g.request_context.identity_id, ReviewQueue.locked_until_datetime <= locked_at_datetime, ), ) ) .with_for_update() ).first() if not row_to_lock: raise_exception_for_ows_response( status=owsresponse_status.NOT_FOUND, code=ERROR_CODE_UNABLE_TO_LOCK, message=ERROR_MESSAGE_UNABLE_TO_LOCK, ) row_to_lock.locked_by_user_id = g.request_context.identity_id row_to_lock.locked_at_datetime = locked_at_datetime row_to_lock.locked_until_datetime = locked_at_datetime + timedelta(minutes=30) return row_to_lock.to_dict() def unlock(review_queue_id): """Unlock a review queue item.""" unlocked_at_datetime = datetime.now(timezone.utc) row_to_unlock = ( ReviewQueue.query.filter(ReviewQueue.review_queue_id == review_queue_id) .filter( and_( ReviewQueue.locked_by_user_id == g.request_context.identity_id, ReviewQueue.locked_until_datetime > unlocked_at_datetime, ) ) .with_for_update() ).first() if not row_to_unlock: raise_exception_for_ows_response( status=owsresponse_status.NOT_FOUND, code=ERROR_CODE_UNABLE_TO_UNLOCK, message=ERROR_MESSAGE_UNABLE_TO_UNLOCK, ) row_to_unlock.locked_until_datetime = unlocked_at_datetime return row_to_unlock.to_dict() def move_review_queue_item_to_queue( *, review_queue_id, queue_name, move_note, moved_to_target_id, moved_to_target_name, moved_to_target_email, moved_to_orchadmin_user_id, ): """Move a review queue item to a different queue.""" review_queue_item = db.session.get(ReviewQueue, review_queue_id) review_queue_item.queue_name = queue_name review_queue_item.move_note = move_note review_queue_item.moved_to_target_id = moved_to_target_id review_queue_item.moved_to_target_name = moved_to_target_name review_queue_item.moved_to_target_email = moved_to_target_email review_queue_item.moved_to_orchadmin_user_id = moved_to_orchadmin_user_id review_queue_item.moved_at_datetime = datetime.now(timezone.utc) review_queue_item.moved_to_queue_by_user_id = g.request_context.identity_id review_queue_item.locked_by_user_id = None review_queue_item.locked_at_datetime = None review_queue_item.locked_until_datetime = None return review_queue_item.to_dict() def approve(review_queue_item): """Approve a product.""" change_status(review_queue_item, STATUS.COMPLETE) def reject(review_queue_item): """Reject a product.""" change_status(review_queue_item, STATUS.COMPLETE) def reset_review_queue_item(review_queue_id): """For integration testing purposes only, reset a ReviewQueueItem.""" review_queue = ReviewQueue.query.filter( ReviewQueue.review_queue_id == review_queue_id ).one() review_queue.status = DEFAULT_STATUS review_queue.queue_name = DEFAULT_QUEUE_NAME review_queue.moved_to_queue_by_user_id = None review_queue.created_datetime = datetime.now(timezone.utc) review_queue.locked_by_user_id = None review_queue.locked_at_datetime = None review_queue.locked_until_datetime = None review_queue.moved_to_target_id = None review_queue.moved_to_target_name = None review_queue.moved_to_target_email = None review_queue.moved_to_orchadmin_user_id = None review_queue.moved_at_datetime = None review_queue.move_note = None def set_to_new(review_queue_id): """For integration testing purposes only, reset a product.""" item = get_item_for_update(review_queue_id) change_status(item, STATUS.NEW) def delete(review_queue_id): """For integration testing purposes only, delete an approval.""" ReviewQueue.query.filter(ReviewQueue.review_queue_id == review_queue_id).delete() def validate_for_completion(review_queue_id): """Validate that a review queue item can be completed.""" review_queue_item = get_item_for_update(review_queue_id) if not review_queue_item: raise_exception_for_ows_response( status=owsresponse_status.NOT_FOUND, code=ERROR_CODE_NOT_FOUND, message=ERROR_MESSAGE_NOT_FOUND, ) if review_queue_item.status == STATUS.COMPLETE: raise_exception_for_ows_response( status=owsresponse_status.BAD_REQUEST, code=ERROR_CODE_INVALID_STATE, message=ERROR_MESSAGE_ALREADY_COMPLETE, ) return review_queue_item