"""Persister File for project manager model.""" from flask import g from oto import response from sqlalchemy import bindparam, text from project_manager.connector.mysql import pm_session_scope from project_manager.constant import error_const from project_manager.constant import feature_flag from project_manager.constant import field_const from project_manager.models import ows_account from project_manager.models import ows_artist from project_manager.models import ows_product from project_manager.models.sql import project_manager from project_manager.util import date_util from project_manager.util import features as features_util def check_db_connectivity(): """Query the db for health check.""" with pm_session_scope() as session: session.execute(project_manager.QUERY_DB_HEARTBEAT) def get_product_genres(): """Get the list of genre options for products. Returns: Response: Response containing the option list. """ get_genres_query = project_manager.SELECT_PRODUCT_GENRES if features_util.is_feature_flag_enabled( feature_flag.INCFEATURES_MAINGENRE_ALTERNATIVE_WORKSTATION, use_account_context=True): get_genres_query = \ project_manager.SELECT_PRODUCT_GENRES_EXCEPT_ALTERNATIVE with pm_session_scope() as session: rows = session.execute(get_genres_query) return response.Response(message=[dict(row) for row in rows]) def get_product_subgenres(genre_id): """Get the list of subgenre options for products given a genre id. Args: genre_id (int): genre id Returns: Response: Response containing the option list. """ with pm_session_scope() as session: rows = session.execute( project_manager.SELECT_PRODUCT_SUBGENRES, {'genre_id': genre_id} ) return response.Response( message=[_format_subgenre(dict(row)) for row in rows]) def get_product_types(): """Get the list of product type options for products. Returns: Response: Response containing the option list. """ return response.Response(message=field_const.PRODUCT_TYPES) def get_project_by_id(project_id, include_soft_delete=False): """Select project by id from project. This gets the project from the project table. Args: project_id (int): project id include_soft_delete (bool): include soft delete project Returns: dict: project """ with pm_session_scope() as session: return _get_project_by_id( project_id, session, include_soft_delete) def get_project_by_project_code( project_code, account_id, subaccount_id): """Select project by project_code from project. This gets the project and some artist info from the project table. Args: project_code (str): project_code Returns: dict: project """ with pm_session_scope() as session: query_parameters = { 'project_code': project_code, 'account_id': account_id, 'subaccount_id': subaccount_id, 'include_deletions': 'N' } row = session.execute( project_manager.GET_PROJECT_BY_PROJECT_CODE, query_parameters ).fetchone() if row: return response.Response(message=dict(row)) else: missing_project_message = 'No project found for project_code={}' \ .format(project_code) g.log.warning(missing_project_message) message = 'Unable to find project_code {}'.format(project_code) return response.create_not_found_response(message=message) def get_existing_projects(project_codes, account_uuid, subaccount_uuid): """Get a list of projects that exist in a list of project_codes. Args: project_codes (str[]) account_uuid (str) subaccount_uuid (str | None) Returns: list of projects that exist """ with pm_session_scope() as session: query_parameters = { 'project_codes': project_codes, 'account_uuid': account_uuid, 'subaccount_uuid': subaccount_uuid } if subaccount_uuid is None: query = text(project_manager.GET_PROJECTS_BY_PROJECT_CODES_VENDOR_ONLY) else: query = text(project_manager.GET_PROJECTS_BY_PROJECT_CODES) query = query.bindparams(bindparam('project_codes', expanding=True)) rows = session.execute( query, query_parameters ).fetchall() existing_projects = [] for row in rows: if row: project = dict(row) existing_projects.append(project) return response.Response(message=existing_projects) def get_project_and_artist_info_by_id(project_id, with_tenant_uuids=False): """Select project ready for maxwells by id from project. This gets the project and some artist info from the project table. Args: project_id (int): project id with_tenant_uuids (bool): Include 4 tenant level UUID fields in result or not. Returns: dict: project """ query = project_manager.SELECT_PROJECT_WITH_TENANT_UUIDS \ if with_tenant_uuids \ else project_manager.SELECT_PROJECT_AND_ARTIST_INFO_BY_ID with pm_session_scope() as session: result = session.execute(query, {'project_id': project_id}) row = result.fetchone() if row: message = dict(row) return response.Response(message=message) else: missing_project_message = 'No project found for project_id={}' \ .format(project_id) g.log.warning(missing_project_message) message = 'Unable to find project_id {}'.format(project_id) return response.create_not_found_response(message=message) def _get_project_by_id( project_id, session, include_soft_delete=False): """Helper function for getting project. Helper function to avoid duplicating code between update and get project functions. Args: project_id (int): project_id session (SQLAlchemy session): session include_soft_delete (bool): include soft delete Returns: dict: project """ sql = project_manager.SELECT_PROJECT_BY_ID if include_soft_delete: sql = project_manager.SELECT_PROJECT_BY_ID_INCLUDING_DELETED result = session.execute(sql, { 'project_id': project_id, 'include_deletions': 'N'}) row = result.fetchone() if row: message = dict(row) return response.Response(message=message) else: missing_project_message = 'No project found for project_id={}'.format( project_id) g.log.warning(missing_project_message) message = 'Unable to find project_id {}'.format(project_id) return response.create_not_found_response(message=message) def get_projects_by_ids(project_ids): """Bulk-select live projects for a list of project ids. Args: project_ids (list[int]): project ids to fetch Returns: Response: response with a list of project dicts """ if not project_ids: return response.Response(message=[]) with pm_session_scope() as session: query = text(project_manager.SELECT_PROJECTS_BY_IDS) query = query.bindparams(bindparam('project_ids', expanding=True)) rows = session.execute( query, {'project_ids': project_ids, 'include_deletions': 'N'} ).fetchall() return response.Response(message=[dict(row) for row in rows]) def get_products_for_project( project_id, vendor_id, ): """Select products from project by id. This gets the products from the releases table. Args: project_id (int): project id vendor_id (int): vendor_id from project Returns: list: products """ with pm_session_scope() as session: return _get_products_for_project(project_id, session, vendor_id) def _get_products_for_project( project_id, session, vendor_id): """Helper function for getting products for project. Args: project_id (int): project_id session (SQLAlchemy session): session vendor_id (int): vendor_id from project Returns: response.Response: response object with message and status """ rows = session.execute(project_manager.SELECT_PRODUCTS_BY_PROJECT_ID, { 'project_id': project_id, 'include_deletions': 'N'}) message = [] map_is_enabled_rejections_from_cr = {} for row in rows: product = dict(row) distribution_format_id = product['distribution_format_id'] if (map_is_enabled_rejections_from_cr .get(distribution_format_id) is None): map_is_enabled_rejections_from_cr[distribution_format_id] = ( features_util.get_rejections_from_cr_ff_status( vendor_id, distribution_format_id, )) product['_is_enabled_rejections_from_cr'] = ( map_is_enabled_rejections_from_cr[distribution_format_id]) message.append(product) return response.Response(message=message) def get_product_for_project(project_id, product_id, vendor_id): """Select product for a project by id. This gets a product from the releases table. Args: project_id (int): project id product_id (int): product id vendor_id (int): vendor id Returns: response.Response: response object with message and status """ with pm_session_scope() as session: return _get_product_for_project( project_id, product_id, vendor_id, session) def _get_product_for_project(project_id, product_id, vendor_id, session): """Helper function for getting a product for a project. Args: project_id (int): project id product_id (int): product id vendor_id (int): vendor id session (SQLAlchemy session): session Returns: response.Response: response object with message and status """ result = session.execute(project_manager.SELECT_PRODUCT_BY_PRODUCT_ID, { 'project_id': project_id, 'product_id': product_id}) row = result.fetchone() if not row: g.log.warning('No product found for product id {}'.format(product_id)) message = 'Unable to find product_id {}'.format(product_id) return response.create_not_found_response(message=message) product = dict(row) distribution_format_id = product['distribution_format_id'] product['_is_enabled_rejections_from_cr'] = ( features_util.get_rejections_from_cr_ff_status( vendor_id, distribution_format_id, ) ) return response.Response(message=product) def get_artists_for_product(product_id): """Select artists by product_id. This gets the artists names from the release_artists table. Args: product_id (int): product id Returns: response.Response: response object with message and status """ with pm_session_scope() as session: return _get_artists_for_product(product_id, session) def _get_artists_for_product(product_id, session): """Helper function for getting artists for a release. Args: product_id (int): product_id session (SQLAlchemy session): session Returns: response.Response: response object with message and status """ rows = session.execute(project_manager.SELECT_ARTISTS_BY_RELEASE_ID, { 'release_id': product_id}) message = [row[0] for row in rows] or [] return response.Response(message=message) def get_product_imprints_for_subaccount(subaccount_id): """Get product imprint options for a subaccount. Get the list of imprint options for the subaccount with id subaccount_id. Args: subaccount_id (int): ID of subaccount to get imprint options for. Returns: Response: Response containing the option list. """ with pm_session_scope() as session: rows = session.execute( project_manager.SELECT_PRODUCT_IMPRINTS_BY_SUBACCOUNT_ID, {'subaccount_id': subaccount_id}) return response.Response(message=[dict(row) for row in rows]) def get_product_imprints_for_vendor(vendor_id): """Get product imprint options for a vendor. Get the list of imprint options for the vendor with id vendor_id. Args: vendor_id (int): ID of vendor to get imprint options for. Returns: Response: Response containing the option list. """ with pm_session_scope() as session: rows = session.execute( project_manager.SELECT_PRODUCT_IMPRINTS_BY_VENDOR_ID, {'vendor_id': vendor_id}) return response.Response(message=[dict(row) for row in rows]) def get_release_status_for_product(product_id): """Select release_status by product_id. This gets the release_status from the releases table, with additional inferences to determine error_correction and action_required states Args: product_id (int): product_id Returns: response.Response: response object with message and status """ with pm_session_scope() as session: return _get_release_status_for_product(product_id, session) def _get_release_status_for_product(product_id, session): """Helper function for getting release_status for a release. Args: product_id (int): product_id session (SQLAlchemy session): session Returns: response.Response: response object with message and status """ result = session.execute( project_manager.SELECT_RELEASE_STATUS_BY_RELEASE_ID, {'release_id': product_id}) row = result.fetchone() message = { 'release_status': row[0], 'review_status': row[1], 'correction_status': row[2] } return response.Response(message=message) def _check_project_code_unique( session, project_code, vendor_id, subaccount_id): select_params = {field_const.PROJECT_CODE: project_code, field_const.VENDOR_ID: vendor_id, field_const.SUBACCOUNT_ID: subaccount_id} # check if project_code already exists for this vendor/subaccount res = session.execute( project_manager.SELECT_PROJECT_BY_PROJECT_VENDOR_SUBACCOUNT, select_params) row = res.fetchone() if row: error_detail = "Project code '{}' already exists".format(project_code) error_message = {'project_code': { 'validator': 'used', 'validator_value': True, 'message': error_detail, 'deletions': row['deletions'], 'existing_project_id': row['project_id']}} return response.create_error_response( code=error_const.VALIDATION_ERROR, message=error_message, status=400) return response.Response() def insert_project(req): """Insert a project. Args: req (ProjectPostRequest object): project to insert Returns: dictionary of entire project record that was inserted """ ins_upd_date = date_util.get_timestamp_utc_iso8601() insert_params = {field_const.PROJECT_CODE: req.project_code, field_const.VENDOR_ID: req.vendor_id, field_const.SUBACCOUNT_ID: req.subaccount_id, field_const.PROJECT_NAME: req.project_name, field_const.CREATED_DATE_UTC: ins_upd_date, field_const.UPDATED_DATE_UTC: ins_upd_date, field_const.CORRELATION_ID_DB: req.correlation_id, field_const.ARTIST_ID: req.artist_id, field_const.DESCRIPTION: req.description} select_params = {field_const.PROJECT_CODE: req.project_code, field_const.VENDOR_ID: req.vendor_id, field_const.SUBACCOUNT_ID: req.subaccount_id} with pm_session_scope() as session: # check if project_code already exists for this vendor/subaccount res = _check_project_code_unique( session, req.project_code, req.vendor_id, req.subaccount_id) if not res: return res user_type, user_id = req.user.split(':') insert_params.update( {'last_modified_by': user_id, 'user_type': user_type}) session.execute( project_manager.INSERT_PROJECT_WITH_USER_ID_AND_TYPE, insert_params) res = session.execute( project_manager.SELECT_PROJECT_BY_PROJECT_VENDOR_SUBACCOUNT, select_params) row = res.fetchone() if row: message = dict(row) return response.Response(message=message, status=201) else: error_detail = "Error inserting params '{}' into database".format( insert_params) return response.create_error_response( code=error_const.DATABASE_ERROR, message=error_detail, status=500) def update_project_by_id(project_id, params, user): """Update project. Updates project by project_id. Args: project_id (int): project_id params (dict): {project_name: some_project_name} user (string): orchard_user_id Returns: response.Response: response object with project """ products_response = None with pm_session_scope() as session: project_response = _get_project_by_id(project_id, session, True) if not project_response: return project_response project = project_response.message artist_id = project['artist_id'] vendor_id = project['vendor_id'] project['updated_date_utc'] = date_util.get_timestamp_utc_iso8601() artist_id_params = params.get('artist_id') artist_params = params.get('artist') if artist_params and not artist_id_params: created_artist = ows_artist.create_artist_info(artist_params) if not created_artist: return created_artist params['artist_id'] = created_artist.message.get('id') product_params = {} update_params = { 'project_id': project_id, 'project_name': project['project_name'], 'project_code': project['project_code'], 'artist_id': project['artist_id'], 'description': project['description'], 'updated_date_utc': project['updated_date_utc'], 'deletions': project['deletions']} for key in params: if key in update_params.keys() and key != 'project_code': project[key] = params[key] update_params[key] = params[key] project_code = params.get('project_code') if project_code and project_code != project['project_code']: res = _check_project_code_unique( session, project_code, project[field_const.VENDOR_ID], project[field_const.SUBACCOUNT_ID]) if not res: return res products_response = get_products_for_project(project_id, vendor_id) product_params = {'project_code': project_code} project['project_code'] = project_code update_params['project_code'] = project_code user_type, user_id = user.split(':') update_params.update({'user_id': user_id, 'user_type': user_type}) session.execute( project_manager.UPDATE_PROJECT_BY_ID_WITH_USER_ID_AND_TYPE, update_params) if artist_id and artist_id != project['artist_id']: product_params.update({'artist_id': project['artist_id']}) if not products_response: products_response = get_products_for_project( project_id, vendor_id) if products_response and products_response.message: # There is trigger magic in DB. If we update products before changes # to project committed, awesome trigger will create new project for # product to handle new product_code. Have no idea why. update_response = ows_product.update_products_by_pool( products_response.message, product_params) if not update_response: return update_response return response.Response(message=project) def get_projects(vendor_id, subaccount_id, page_offset, page_limit): """Get projects for a vendor or subaccount. Args: vendor_id (int): vendor id subaccount_id (int): subaccount id page_offset (int): record index used to start page_limit (int): number of records to fetch Returns: response.Response: containing list of projects """ with pm_session_scope() as session: subaccount_params = { 'subaccount_id': subaccount_id, 'page_offset': page_offset, 'page_limit': page_limit, 'include_deletions': 'N'} vendor_params = { 'vendor_id': vendor_id, 'page_offset': page_offset, 'page_limit': page_limit, 'include_deletions': 'N'} if subaccount_id: result = session.execute( project_manager.SELECT_PROJECTS_BY_SUBACCOUNT_ID, subaccount_params) count = session.execute( project_manager.SELECT_PROJECTS_BY_SUBACCOUNT_ID_COUNT, subaccount_params ).scalar() else: result = session.execute( project_manager.SELECT_PROJECTS_BY_VENDOR_ID, vendor_params) count = session.execute( project_manager.SELECT_PROJECTS_BY_VENDOR_ID_COUNT, vendor_params).scalar() rows = result.fetchall() if count is None: count = 0 items = [dict(row) for row in rows] projects = _get_subaccount_data_for_project(items) projects_response = response.Response(message=projects) message = { 'pagination': { 'page_offset': page_offset, 'page_limit': page_limit, 'total_records': count} } message['items'] = projects_response.message return response.Response(message=message, status=200) def _get_subaccount_data_for_project(projects): """Helper function gettin subaccount data for projects.""" for project in projects: subaccount_id = project.get('subaccount_id') if subaccount_id: subaccount_response = ows_account.get_subaccount_for_subaccount_id( subaccount_id) if subaccount_response: project['subaccount'] = subaccount_response return projects def _format_subgenre(subgenre): """Convert a subgenre dict from the DB to the desired format. Args: subgenre (dict): unformatted subgenre Returns: subgenre (dict): formatted subgenre """ return { 'id': subgenre.get('id'), 'composer': (subgenre.get('composer') == 'Y'), 'name': subgenre.get('name')}