"""Function Logic.""" from datetime import datetime import json from content_utils.exceptions import IneligibleEventError from src import config from src.constants import ADD_INDEX_OP from src.exceptions import ReviewQueueItemExplicitSkip def _extract_update_data(db_msg, payload): """Extract update data from db message and update payload.""" locked_by_user_id = db_msg.get_field_values('locked_by_user_id') locked_until_datetime = db_msg.get_field_values('locked_until_datetime') is_changed_locked_by_user_id = locked_by_user_id.get('before') != locked_by_user_id.get('after') is_changed_locked_until_datetime = ( locked_until_datetime.get('before') != locked_until_datetime.get('after')) if not (is_changed_locked_by_user_id or is_changed_locked_until_datetime): raise IneligibleEventError( f'skip operation: {db_msg.operation}, before: {json.dumps(db_msg.before)}, after: {json.dumps(db_msg.after)}') # noqa if is_changed_locked_by_user_id: payload['locked_by_user_id'] = locked_by_user_id.get('after') if is_changed_locked_until_datetime: payload['locked_until_datetime'] = ( locked_until_datetime.get('after') and ( f"{datetime.utcfromtimestamp(locked_until_datetime.get('after') / 1e6).isoformat()}Z") # noqa ) def parse_message_payload(db_msg): """Process a single CDC event from review queue.""" payload = _get_payload(db_msg) product_id = payload.get('product_id') review_queue_id = payload.get('id') if review_queue_id in config.SKIP_REVIEW_QUEUE_ITEMS: raise ReviewQueueItemExplicitSkip('Explicit Skip') operation_type, payload = _determine_operation_type(db_msg, payload, product_id) return { 'product_id': product_id, 'review_queue_id': review_queue_id, 'payload': payload, 'operation_type': operation_type } def _get_payload(db_msg): """Extract payload from db message.""" return db_msg.after if db_msg.after is not None else db_msg.before def _determine_operation_type(db_msg, payload, product_id): """Determine the operation type and modify payload if necessary.""" status = db_msg.get_field_values('status') is_update = db_msg.operation == 'u' and not ( status.get('before') == 'applying' and status.get('after') == 'new') is_deindex = db_msg.operation == 'd' or ( is_update and status.get('before') == 'new' and status.get('after') == 'complete') is_no_op = db_msg.operation == 'c' and status.get('after') == 'complete' if is_no_op: raise IneligibleEventError( f'skip operation: {db_msg.operation}, before: {json.dumps(db_msg.before)}, after: {json.dumps(db_msg.after)}') # noqa elif is_deindex: return 'deindex', None elif is_update: payload = {'product_id': product_id} _extract_update_data(db_msg, payload) return 'update_index', payload else: return ADD_INDEX_OP, payload