"""Module for fact_conflict table.""" from owsfeatures import features as owsfeatures from conflict_manager import features from conflict_manager.connectors import snowflake from conflict_manager.constants import conflict as conflict_const from conflict_manager.utils import api_utils from conflict_manager.utils import model_utils _SORT_FIELD_MAP_INNER = { 'track_name': 'track.track_name', 'track_artist_names': 'inner_artist.artist_names', 'product_name': 'releases.release_name', 'conflicting_owner_name': 'inner_conflict.conflicting_owner', 'daily_average_views': 'inner_conflict.yt_recent_daily_average', 'conflict_date': 'inner_conflict.conflict_date', 'resolved_datetime': 'inner_conflict.resolved_datetime', 'subaccount_name': 'subaccount.subaccount_name', 'tuid': 'inner_conflict.tuid', 'action_date': 'inner_action.action_date' } _SORT_FIELD_MAP_OUTER = { 'track_name': 'grouped_conflict.track_name', 'track_artist_names': 'grouped_conflict.artist_names', 'product_name': 'grouped_conflict.product_name', 'subaccount_name': 'grouped_conflict.subaccount_name', 'conflicting_owner_name': 'conflict.conflicting_owner', 'daily_average_views': 'conflict.yt_recent_daily_average', 'conflict_date': 'conflict.conflict_date', 'resolved_datetime': 'conflict.resolved_datetime', 'tuid': 'conflict.tuid', 'territory': 'conflict.territory', 'action_date': 'action.action_date' } @snowflake.db_session_wrap def get_new_youtube_conflicts_for_account( account, sort_by=None, sort_order='asc', offset=0, limit=50, session=None): """Get new YouTube conflicts for an account. Args: account (namedtuple): User account data sort_by (str): Field to sort results by sort_order (str): Order results by ascending or descending offset (int): Number of items to skip limit (int): Max results to return session (sqlalchemy.orm.session.Session): database session (optional) Returns: response.Response: dict of created action or errors """ conflict_queries = NewYouTubeConflictsQueries(session) return conflict_queries.get_conflicts_list( account, sort_by, sort_order, offset, limit) @snowflake.db_session_wrap def get_actioned_youtube_conflicts_for_account( account, sort_by=None, sort_order='asc', offset=0, limit=50, session=None): """Get actioned YouTube conflicts for an account. Args: account (namedtuple): User account data sort_by (str): Field to sort results by sort_order (str): Order results by ascending or descending offset (int): Number of items to skip limit (int): Max results to return session (sqlalchemy.orm.session.Session): database session (optional) Returns: response.Response: dict of created action or errors """ conflict_queries = ActionedYouTubeConflictsQueries(session) return conflict_queries.get_conflicts_list( account, sort_by, sort_order, offset, limit) @snowflake.db_session_wrap def get_resolved_youtube_conflicts_for_account( account, sort_by=None, sort_order='asc', offset=0, limit=50, session=None): """Get actioned YouTube conflicts for an account. Args: account (namedtuple): User account data sort_by (str): Field to sort results by sort_order (str): Order results by ascending or descending offset (int): Number of items to skip limit (int): Max results to return session (sqlalchemy.orm.session.Session): database session (optional) Returns: response.Response: dict of created action or errors """ conflict_queries = ResolvedYouTubeConflictsQueries(session) return conflict_queries.get_conflicts_list( account, sort_by, sort_order, offset, limit) class BaseConflictsFormatter(): """Base Formatter for conflicts.""" def output(self, conflicts_list): """Group conflicts by owner, date, and isrc. Args: conflicts_list (list): List of ungrouped conflicts Return: list: List of formatted conflicts """ results = [] conflict = {} last_group = () for data in conflicts_list: group = ( data['conflicting_owner'], data['conflict_date'], data['tuid']) if group != last_group: last_group = group if conflict: results.append(conflict) conflict = self._conflict_item_init(data) self.conflict_item_add_territory(conflict, data) # Make sure last conflict is added to results if conflict: results.append(conflict) return results def conflict_item_init(self, conflict, data): """Initialize conflict data. Args: conflict data """ pass def conflict_item_add_territory(self, conflict, data): """Add territory to base conflict data structure.""" pass def _conflict_item_init(self, data): """Initialize conflict data.""" track_name = data['track_name'] if data['version']: track_name = '{} ({})'.format(track_name, data['version']) if data['artist_names']: artist_names = data['artist_names'].split('\0,') else: artist_names = [] if data['subaccount_name']: subaccount_name = data['subaccount_name'] else: subaccount_name = '' if data['es_id']: es_id = data['es_id'] else: es_id = '' conflict = { 'product_id': data['product_id'], 'product_name': data['product_name'], 'display_upc': data['display_upc'], 'tuid': data['tuid'], 'track_name': track_name, 'track_artists': artist_names, 'isrc': data['isrc'], 'conflicting_owner': data['conflicting_owner'], 'conflict_date': data['conflict_date'], 'territory_standard': data['territory_standard'], 'vendor_id': data['vendor_id'], 'subaccount_id': data['subaccount_id'], 'subaccount_name': subaccount_name, 'es_id': es_id } self.conflict_item_init(conflict, data) return conflict class BaseConflictQueries(): """Base Conflict Query.""" inner_columns = [] inner_joins = [] inner_filters = [] formatter_class = BaseConflictsFormatter def __init__(self, session): """ Constructor. Args: session (sqlalchemy.orm.session.Session): database session """ self._formatter = self.formatter_class() self._session = session def get_conflicts_list( self, account, sort_by, sort_order, offset, limit): """Get list of New YouTube conflicts for an account. Args: account (namedtuple): User account data sort_by (str): Field to sort results by sort_order (str): Field to determine sort order offset (int): Number of items to skip limit (int): Max results to return Returns: list """ sql = self._get_conflicts_list_sql(sort_by, sort_order) params = {'offset': offset, 'limit': limit} results = model_utils.run_query( self._session, sql, account, params=params) formated_results = self._formatter.output(results) total_records = model_utils.get_total_records( formated_results, offset, limit) if total_records is None: total_records = self.get_total_records(account) return api_utils.make_pagination_response( formated_results, offset=offset, limit=limit, total_records=total_records) def get_total_records(self, account): """Get total count of new YouTube conflicts for an account. Args: account (namedtuple): User account data Returns: int """ sql = """ select count(distinct inner_conflict.conflicting_owner , inner_conflict.conflict_date , inner_conflict.tuid) from {fact_conflict_table} inner_conflict """ + self._inner_join_postfix + """ where inner_conflict.{account_id_field} = :account_id """ + self._inner_where_postfix query_results = model_utils.run_query(self._session, sql, account) return query_results.fetchone()[0] def _get_conflicts_list_sql(self, sort_by, sort_order): """Get SQL for conflicts list. This is just a placeholder and needs to be overridden """ raise Exception('Please implement me') # pragma: no cover @owsfeatures.load_features def _sort_columns_outer(self, sort_by, sort_order): """List of outer query sort columns by qualified name.""" sort_columns = [] if sort_by: sort_columns.append( '{} {}'.format(_SORT_FIELD_MAP_OUTER[sort_by], sort_order)) if features.is_conflict_mgr_views_in_conflict_data_enabled(): _SORT_FIELD_MAP_OUTER['views_in_conflict'] = 'views_in_conflict' if sort_by != 'views_in_conflict': sort_columns.append(_SORT_FIELD_MAP_OUTER['views_in_conflict']) # Secondary sort params are used to keep order consistent if sort_by != 'conflict_date': sort_columns.append(_SORT_FIELD_MAP_OUTER['conflict_date']) return sort_columns + [ _SORT_FIELD_MAP_OUTER['tuid'], _SORT_FIELD_MAP_OUTER['conflicting_owner_name'], _SORT_FIELD_MAP_OUTER['territory']] @owsfeatures.load_features def _sort_columns_inner(self, sort_by, sort_order): """List of inner query sort columns by qualified name.""" sort_columns = [] if sort_by: sort_columns.append( '{} {}'.format(_SORT_FIELD_MAP_INNER[sort_by], sort_order)) if features.is_conflict_mgr_views_in_conflict_data_enabled(): _SORT_FIELD_MAP_INNER['views_in_conflict'] = 'views_in_conflict' if sort_by != 'views_in_conflict': sort_columns.append(_SORT_FIELD_MAP_OUTER['views_in_conflict']) # Secondary sort params are used to keep order consistent if sort_by != 'conflict_date': sort_columns.append(_SORT_FIELD_MAP_INNER['conflict_date']) sort_columns += [ _SORT_FIELD_MAP_INNER['tuid'], _SORT_FIELD_MAP_INNER['conflicting_owner_name']] return sort_columns @owsfeatures.load_features def _inner_query_sql(self, sort_by='', sort_order='asc'): """Generate group conflicts SQL string. Args: sort_by (str): Field to sort results by sort_order (str): Order results by ascending or descending Returns: str """ sort_columns = self._sort_columns_inner(sort_by, sort_order) views_in_conflict_select = '' if features.is_conflict_mgr_views_in_conflict_data_enabled(): views_in_conflict_select = ', inner_conflict.views_in_conflict' columns_postfix = '' if self.inner_columns: columns_postfix = ', {}'.format(', '.join(self.inner_columns)) return (""" select distinct inner_conflict.conflicting_owner , inner_conflict.conflict_date , inner_conflict.isrc , inner_conflict.tuid , inner_conflict.{account_id_field} , inner_conflict.yt_recent_daily_average """ + views_in_conflict_select + """ , track.track_name , track.version , inner_artist.artist_names , track.release_id as product_id , releases.release_name as product_name , releases.display_upc , subaccount.subaccount_name """ + columns_postfix + """ from {fact_conflict_table} inner_conflict left join {track} on (inner_conflict.tuid = track.id) left join {releases} on (track.release_id = releases.release_id) left join ( select track_id, listagg(name, '\\0,') within group(order by DECODE( lower(type), 'performer', 1, 'remixer', 2, 'producer', 3 ), name) as artist_names from {track_artist} group by track_id ) inner_artist on (inner_artist.track_id = track.id) left join {subaccount} on (subaccount.subaccount_id = inner_conflict.subaccount_id) """ + self._inner_join_postfix + """ where inner_conflict.{account_id_field} = :account_id """ + self._inner_where_postfix + """ order by """ + ', '.join(sort_columns) + """ limit :limit offset :offset """) @property def _inner_join_postfix(self): """Condition to append at end of where clause.""" return ' '.join(self.inner_joins) @property def _inner_where_postfix(self): """Condition to append at end of where clause.""" if self.inner_filters: return 'and {}'.format(' and '.join(self.inner_filters)) return '' class NewYouTubeConflictsFormatter(BaseConflictsFormatter): """New YouTube Conflicts Formatter.""" def conflict_item_init(self, conflict, data): """Initialize conflict data.""" conflict['status'] = conflict_const.STATUS_NEW conflict['daily_average_views'] = data['yt_recent_daily_average'] if features.is_conflict_mgr_views_in_conflict_data_enabled(): conflict['views_in_conflict'] = data['views_in_conflict'] conflict['territories'] = [] def conflict_item_add_territory(self, conflict, data): """Add territory to conflict data.""" conflict['territories'].append({ 'conflict_id': data['conflict_id'], 'code': data['territory']}) class NewYouTubeConflictsQueries(BaseConflictQueries): """New YouTuble Conflicts Queries.""" inner_joins = [ 'left join {action_table} inner_action ' 'on (inner_conflict.conflict_id = inner_action.conflict_id)'] inner_filters = [ 'inner_action.conflict_id is NULL', 'inner_conflict.resolved_datetime is NULL'] formatter_class = NewYouTubeConflictsFormatter def _get_conflicts_list_sql(self, sort_by, sort_order): """Get Conflicts List SQL.""" sort_columns = self._sort_columns_outer(sort_by, sort_order) inner_conflict_sql = self._inner_query_sql(sort_by, sort_order) sql = """ select conflict.* , grouped_conflict.track_name , grouped_conflict.version , grouped_conflict.artist_names , grouped_conflict.product_id , grouped_conflict.product_name , grouped_conflict.display_upc , grouped_conflict.subaccount_name from {fact_conflict_table} conflict left join {action_table} action on conflict.conflict_id = action.conflict_id join (""" + inner_conflict_sql + """) grouped_conflict on ( conflict.{account_id_field} = grouped_conflict.{account_id_field} and conflict.conflicting_owner = grouped_conflict.conflicting_owner and conflict.tuid = grouped_conflict.tuid and conflict.conflict_date = grouped_conflict.conflict_date and conflict.yt_recent_daily_average = grouped_conflict.yt_recent_daily_average) where conflict.resolved_datetime is NULL and action.conflict_id is NULL """ # NOQA E501 if features.is_show_only_indexed_conflicts_enabled(): sql += ' and conflict.es_indexed = true \n' sql += 'group by all order by {};'.format(', '.join(sort_columns)) return sql class ActionedYouTubeConflictsFormatter(BaseConflictsFormatter): """Actioned YouTube Conflicts Formatter.""" ACTION_MAP = { 'assert': conflict_const.ASSERT_ACTION, 'release': conflict_const.RELEASE_ACTION } def conflict_item_init(self, conflict, data): """Initialize conflict data.""" conflict['status'] = conflict_const.STATUS_ACTIONED conflict['action_date'] = data['action_date'] conflict['response_account_id'] = data['account_id'] conflict['response_account_type'] = data['account_type'] for action_type in self.ACTION_MAP.values(): conflict[action_type] = { 'territories': [], 'reason': '', 'additional_information': '' } def conflict_item_add_territory(self, conflict, data): """Add territory to conflict data.""" action_type = self.ACTION_MAP[data['action']] conflict_action = conflict[action_type] conflict_action['territories'].append({ 'conflict_id': data['conflict_id'], 'code': data['territory'] }) if not conflict_action['reason']: conflict_action['reason'] = data['reason'] conflict_action['additional_information'] = \ data['additional_information'] class ActionedYouTubeConflictsQueries(BaseConflictQueries): """Actioned YouTuble Conflicts Queries.""" inner_columns = ['inner_action.action_date'] inner_joins = [ 'join {action_table} inner_action ' 'on (inner_conflict.conflict_id = inner_action.conflict_id)'] formatter_class = ActionedYouTubeConflictsFormatter def _get_conflicts_list_sql(self, sort_by, sort_order): """Get Conflicts List SQL.""" sort_columns = self._sort_columns_outer(sort_by, sort_order) inner_conflict_sql = self._inner_query_sql(sort_by, sort_order) sql = """ select conflict.* , grouped_conflict.track_name , IFNULL(grouped_conflict.version, '') as version , grouped_conflict.artist_names , grouped_conflict.product_id , grouped_conflict.product_name , grouped_conflict.display_upc , grouped_conflict.subaccount_name , action.* from {fact_conflict_table} conflict join (""" + inner_conflict_sql + """) grouped_conflict on ( conflict.{account_id_field} = grouped_conflict.{account_id_field} and conflict.conflicting_owner = grouped_conflict.conflicting_owner and conflict.tuid = grouped_conflict.tuid and conflict.conflict_date = grouped_conflict.conflict_date) join {action_table} action on conflict.conflict_id = action.conflict_id """ # NOQA E501 sql += 'order by {};'.format(', '.join(sort_columns)) return sql class ResolvedYouTubeConflictsFormatter(BaseConflictsFormatter): """Resolved YouTube Conflicts Formatter.""" def conflict_item_init(self, conflict, data): """Initialize conflict data.""" conflict['status'] = conflict_const.STATUS_RESOLVED conflict['resolved_datetime'] = data['resolved_datetime'] conflict['territories'] = [] def conflict_item_add_territory(self, conflict, data): """Add territory to conflict data.""" conflict['territories'].append({ 'conflict_id': data['conflict_id'], 'code': data['territory']}) class ResolvedYouTubeConflictsQueries(BaseConflictQueries): """Resolved YouTuble Conflicts Queries.""" inner_columns = ['inner_conflict.resolved_datetime'] inner_filters = ['inner_conflict.resolved_datetime is not NULL'] formatter_class = ResolvedYouTubeConflictsFormatter def _get_conflicts_list_sql(self, sort_by, sort_order): """Get Conflicts List SQL.""" sort_columns = self._sort_columns_outer(sort_by, sort_order) inner_conflict_sql = self._inner_query_sql(sort_by, sort_order) sql = """ select conflict.* , grouped_conflict.track_name , IFNULL(grouped_conflict.version, '') as version , grouped_conflict.artist_names , grouped_conflict.product_id , grouped_conflict.product_name , grouped_conflict.display_upc , grouped_conflict.subaccount_name from {fact_conflict_table} conflict join (""" + inner_conflict_sql + """) grouped_conflict on ( conflict.{account_id_field} = grouped_conflict.{account_id_field} and conflict.conflicting_owner = grouped_conflict.conflicting_owner and conflict.tuid = grouped_conflict.tuid and conflict.conflict_date = grouped_conflict.conflict_date and conflict.resolved_datetime = grouped_conflict.resolved_datetime) where conflict.resolved_datetime is not NULL """ # NOQA E501 sql += 'order by {};'.format(', '.join(sort_columns)) return sql @snowflake.db_session_wrap def get_conflict_ids_for_account_and_isrc( account, isrc, conflicting_owner, conflict_date, tuid, session=None): """Get conflict ids for given account and isrc. E.g conflict between two owners about single ISRC on 20 territories will be saved as 20 entries in fact_conflict table(one for each country). Used for validation purposes when creating action. Args: account (namedtuple): User account data isrc (str): ISRC to filter conflicts by conflicting_owner (str): Owner who claimed rights for given ISRC conflict_date (str): Date when conflict was created session (sqlalchemy.orm.session.Session): database session (optional) Returns: response.Response: list of conflict_ids """ sql = """ SELECT conflict_id FROM {fact_conflict_table} WHERE {account_id_field} = :account_id AND isrc = :isrc AND tuid = :tuid AND conflicting_owner = :conflicting_owner AND conflict_date = :conflict_date AND resolved_datetime IS NULL; """ query_results = model_utils.run_query( session, sql, account, params={ 'isrc': isrc, 'tuid': tuid, 'conflicting_owner': conflicting_owner, 'conflict_date': conflict_date }) results = [row[0] for row in query_results] return api_utils.make_pagination_response(results) @snowflake.db_session_wrap def get_conflicts_ids_for_bulk_actions(account, actions, session=None): """Get conflicts ids for each given account and action. Used for validation purposes when bulk creating actions. Args: account (namedtuple): User account data actions (list): POST /action/bulk payload session (sqlalchemy.orm.session.Session): database session (optional) Returns: response.Response: list of conflicts count """ where_sql = """ ( {{account_id_field}} = :account_id AND isrc = :isrc{i} AND tuid = :tuid{i} AND conflicting_owner = :conflicting_owner{i} AND conflict_date = :conflict_date{i} AND resolved_datetime IS NULL ) """ select_sql = """ SELECT LISTAGG(conflict_id, ',') as conflict_ids, isrc, tuid, conflicting_owner, conflict_date FROM {fact_conflict_table} WHERE """ group_by_sql = """ group by isrc, vendor_id, isrc, tuid,conflicting_owner, conflict_date; """ params = {} where_conditions = [] for i, action in enumerate(actions): where_conditions.append(where_sql.format(i=i)) params['isrc{}'.format(i)] = action['isrc'] params['conflicting_owner{}'.format(i)] = action['conflicting_owner'] params['conflict_date{}'.format(i)] = action['conflict_date'] params['tuid{}'.format(i)] = action['tuid'] sql = '{select} {where} {group_by}'.format( select=select_sql, where=' OR '.join(where_conditions), group_by=group_by_sql) query_results = model_utils.run_query( session, sql, account, params=params).fetchall() return query_results @snowflake.db_session_wrap def get_grouped_conflicts_ids(grouped_conflicts, session): """Get list of grouped conflicts ids. Used for validation purposes when bulk updating conflicts status. Args: grouped_conflicts (list): list of grouped conflicts ids form POST /conflicts/status/bulk payload session (sqlalchemy.orm.session.Session): database session Returns: set: list of grouped conflicts ids """ select_sql = """ SELECT tuid || '|' || CONFLICT_DATE || '|' || CONFLICTING_OWNER || '|' || action as grouped_id FROM {fact_conflict_table} conflict JOIN {action_table} action ON conflict.conflict_id = action.conflict_id WHERE """ where_sql = """ ( tuid = :tuid{i} AND conflicting_owner = :conflicting_owner{i} AND conflict_date = :conflict_date{i} AND action = :action{i} ) """ group_by_sql = 'GROUP BY tuid, conflict_date, conflicting_owner, action;' params = {} where_conditions = [] for i, conflict in enumerate(grouped_conflicts): where_conditions.append(where_sql.format(i=i)) params['tuid{}'.format(i)] = conflict['tuid'] params['conflicting_owner{}'.format(i)] = conflict['conflicting_owner'] params['conflict_date{}'.format(i)] = conflict['conflict_date'] params['action{}'.format(i)] = conflict['action'] sql = '{select} {where} {group_by}'.format( select=select_sql, where=' OR '.join(where_conditions), group_by=group_by_sql) query_results = model_utils.run_query( session=session, sql=sql, params=params).fetchall() return set([row[0] for row in query_results]) @snowflake.db_session_wrap def get_conflicts_territories(conflict_ids, session): """Get lists of territories for a list conflicts ids. Territories are grouped by: tuid, conflict_date, conflicting_owner and isrc. Args: conflict_ids (list): list of conflict_id session (sqlalchemy.orm.session.Session): database session Returns: list (dict): list of dicts containing territories and conflict information Example: [ { "conflict_date": "2018-12-22", "conflicting_owner": "believe_africa", "isrc": "QM6MZ1783179", "territories": ["BR", "HT"], "tuid": 25874822 }, ] """ sql = """ SELECT tuid, CONFLICT_DATE, CONFLICTING_OWNER, isrc, LISTAGG(territory, ',') as territories FROM {fact_conflict_table} conflict WHERE conflict_id in (:ids) GROUP BY tuid, conflict_date, conflicting_owner, isrc; """ query_results = model_utils.run_query( session=session, sql=sql, params={'ids': conflict_ids}).fetchall() conflicts = [dict(row) for row in query_results] for conflict in conflicts: conflict['territories'] = conflict['territories'].split(',') conflict['conflict_date'] = conflict['conflict_date'].isoformat() return conflicts