from datetime import timedelta from typing import Any, List import structlog from db_schema.common import ActivityStatusEnum from db_schema.schemas.slz import ( CompletenessStatusEnum, ContentStatus, ContentStatusEnum, UnitOfWork, UnitOfWorkPriorityEnum, ) from flask import current_app, flash, request from flask_admin.actions import action from flask_admin.contrib.sqla.filters import BaseSQLAFilter from sqlalchemy import and_, func, literal, null from sqlalchemy.orm import Query, Session from delphi_slz_admin.const import APP_NAME from delphi_slz_admin.helpers import cast_ids_to_int, utcnow from delphi_slz_admin.modelviews import ( BaseFilterByDateEarlierThan, BaseFilterByDateEquals, BaseFilterByDateLaterThan, BaseFilterById, BaseFilterByIdReversed, BaseModelView, ) from delphi_slz_admin.services.audit_log import AuditLogService from delphi_slz_admin.services.config import Config from delphi_slz_admin.services.notification import NotificationService from delphi_slz_admin.services.reprocessing import ( ExplorationReprocessingService, ReprocessingService, ) logger = structlog.get_logger(APP_NAME) class FilterUnitOfWorkId(BaseFilterById): def __init__(self, **kwargs: Any) -> None: self.column = ContentStatus.unit_of_work_id self.name = 'Unit of Work ID' super().__init__(self.column, self.name, **kwargs) class FilterUnitOfWorkIdReversed(BaseFilterByIdReversed): def __init__(self, **kwargs: Any) -> None: self.column = ContentStatus.unit_of_work_id self.name = 'Unit of Work ID' super().__init__(self.column, self.name, **kwargs) class FilterContentStatusId(BaseFilterById): def __init__(self, **kwargs: Any) -> None: self.column = ContentStatus.content_status_id self.name = 'Content Status ID' super().__init__(self.column, self.name, **kwargs) class FilterContentStatusIdReversed(BaseFilterByIdReversed): def __init__(self, **kwargs: Any) -> None: self.column = ContentStatus.content_status_id self.name = 'Content Status ID' super().__init__(self.column, self.name, **kwargs) class FilterContentStatus(BaseSQLAFilter): def __init__(self, **kwargs: Any) -> None: self.column = ContentStatus.content_status self.name = 'Content Status' self.options = [(item.value, item.value) for item in ContentStatusEnum] super().__init__(self.column, self.name, options=self.options, **kwargs) def apply(self, query: Query, value: str, alias: Any = None) -> Query: return query.filter(ContentStatus.content_status == value) def operation(self) -> str: return 'equals' class FilterUoWCodeLike(BaseSQLAFilter): def __init__(self, **kwargs: Any) -> None: self.column = UnitOfWork.unit_of_work_code self.name = 'UoW Code' super().__init__(self.column, self.name, **kwargs) def apply(self, query: Query, value: str, alias: Any = None) -> Query: return query.filter(self.column.like(f'%{value}%')) def operation(self) -> str: return 'like' class FilterReportDateMixin(BaseSQLAFilter): def __init__(self, **kwargs: Any) -> None: self.column = UnitOfWork.report_date self.name = 'Report date' super().__init__(self.column, self.name, data_type='datepicker', **kwargs) def operation(self) -> str: pass class FilterReportDateLaterThan(BaseFilterByDateLaterThan, FilterReportDateMixin): pass class FilterReportDateEarlierThan(BaseFilterByDateEarlierThan, FilterReportDateMixin): pass class FilterReportDateEquals(BaseFilterByDateEquals, FilterReportDateMixin): pass class FilterSameReportLicensorContext(BaseSQLAFilter): def __init__(self, **kwargs: Any) -> None: self.column = UnitOfWork.report_id self.name = 'Same Report/Licensor/Context as' super().__init__(self.column, self.name, **kwargs) def clean(self, value: str) -> int: return int(value) def apply(self, query: Query, value: str, alias: Any = None) -> Query: subquery = ( query.session.query( UnitOfWork.report_id, UnitOfWork.licensor_id, ContentStatus.context ).join(UnitOfWork).filter(ContentStatus.content_status_id == value) ).subquery() return query.join( subquery, and_( subquery.c.report_id == UnitOfWork.report_id, subquery.c.licensor_id == UnitOfWork.licensor_id, subquery.c.context == ContentStatus.context, ) ) def operation(self) -> str: return 'Content Status ID' class FilterSameDayOfWeek(BaseSQLAFilter): def __init__(self, **kwargs: Any) -> None: self.column = UnitOfWork.report_id self.name = 'Same Report/Licensor/Context/Day of week as' super().__init__(self.column, self.name, **kwargs) def clean(self, value: str) -> int: return int(value) def apply(self, query: Query, value: str, alias: Any = None) -> Query: subquery = ( query.session.query( UnitOfWork.report_id, UnitOfWork.licensor_id, UnitOfWork.report_date, ContentStatus.context ).join(UnitOfWork).filter(ContentStatus.content_status_id == value) ).subquery() return query.join( subquery, and_( subquery.c.report_id == UnitOfWork.report_id, subquery.c.licensor_id == UnitOfWork.licensor_id, subquery.c.context == ContentStatus.context, func.extract('isodow', subquery.c.report_date) == func.extract( 'isodow', UnitOfWork.report_date ), ) ) def operation(self) -> str: return 'Content Status ID' class ContentStatusModelView(BaseModelView): column_list = [ 'content_status.unit_of_work_id', 'content_status.content_status_id', 'unit_of_work.unit_of_work_code', 'unit_of_work.report_date', 'unit_of_work.timeslot', 'unit_of_work.report_id', 'unit_of_work.licensor_id', 'unit_of_work.activity_status', 'unit_of_work.completeness_status', 'content_status.context', 'content_status.content_name', 'content_status.content_size', 'content_status.content_status', 'content_status.hash', 'content_status.record_count', 'content_status.failure_count', 'content_status.latest_job_id', 'content_status.created_at', 'content_status.completed_at', 'content_status.last_checked_at', 'content_status.sub_content', 'content_status.metadata_process_status', 'content_status.metadata_process_started_at', 'content_status.metadata_process_completed_at', ] column_labels = { 'unit_of_work.report_date': '[UoW] Report Date', 'unit_of_work.timeslot': '[UoW] Timeslot', 'unit_of_work.unit_of_work_code': '[UoW] UoW Code', 'unit_of_work.report_id': '[UoW] Report ID', 'unit_of_work.licensor_id': '[UoW] Licensor ID', 'unit_of_work.activity_status': '[UoW] Activity Status', 'unit_of_work.completeness_status': '[UoW] Completeness Status', 'content_status_id': '[CS] Content Status ID', 'context': '[CS] Context', 'content_name': '[CS] Content Name', 'content_size': '[CS] Content Size', 'content_status': '[CS] Content Status', 'hash': '[CS] Hash', 'record_count': '[CS] Record Count', 'failure_count': '[CS] Failure Count', 'latest_job_id': '[CS] Latest Job Id', 'created_at': '[CS] Created At', 'completed_at': '[CS] Completed At', 'last_checked_at': '[CS] Last Checked At', 'sub_content': '[CS] Sub Content', 'metadata_process_status': '[CS] Metadata Process Status', 'metadata_process_started_at': '[CS] Metadata Process Started At', 'metadata_process_completed_at': '[CS] Metadata Process Completed At', } column_sortable_list = column_list column_filters = [ FilterUnitOfWorkId(), FilterUnitOfWorkIdReversed(), FilterContentStatusId(), FilterContentStatusIdReversed(), FilterContentStatus(), FilterUoWCodeLike(), FilterReportDateEarlierThan(), FilterReportDateLaterThan(), FilterReportDateEquals(), FilterSameReportLicensorContext(), FilterSameDayOfWeek(), ] column_default_sort = [('unit_of_work.report_date', True), ('content_status_id', True)] list_template = 'content_status_list.html' def __init__(self, session: Session, **kwargs: Any) -> None: self.model = ContentStatus self.audit_log_service = AuditLogService(session) self.notification_service = NotificationService(session) self.reprocessing_service = ReprocessingService(session) self.exp_reprocessing_service = ExplorationReprocessingService() super().__init__(self.model, session, **kwargs) @action( 'complete_cs', 'Complete Content Status', 'Are you sure you want to set Content Status completed?' ) def action_complete_cs(self, ids: List[str]) -> None: try: casted_ids = cast_ids_to_int(ids, logger) query = self.session.query(ContentStatus).filter( ContentStatus.content_status_id.in_(casted_ids) ) has_improper_status = self.session.query(literal(True)).filter( query.filter(ContentStatus.content_status != ContentStatusEnum.ON_HOLD).exists() ).scalar() if has_improper_status: flash( f'Some items are not in {ContentStatusEnum.ON_HOLD.value} status', category='error', ) return failed_ids = self.notification_service.push_content_statuses_metadata(casted_ids) if failed_ids: joined_failed_ids = ','.join(map(str, failed_ids)) flash( f'Failed to push metadata to SQS queue for CS ids: {joined_failed_ids}', category='error', ) return count = query.update( { ContentStatus.content_status: ContentStatusEnum.COMPLETE, ContentStatus.completed_at: utcnow(), ContentStatus.meta_data: null() }, synchronize_session='fetch', ) self.audit_log_service.log_content_status_set_completed(casted_ids) self.session.commit() flash(f'{count} Content Status(es) were successfully completed.') except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to complete Content Statuses, details = %s', ex) flash('Failed to complete Content Statues', category='error') def _action_complete_uow(self, ids: List[str], status: CompletenessStatusEnum) -> None: audit_log_method_mapping = { CompletenessStatusEnum.COMPLETE: self.audit_log_service.log_unit_of_work_set_completed, CompletenessStatusEnum.MIN_COMPLETE: self.audit_log_service.log_unit_of_work_set_min_completed, } try: query = self.session.query(ContentStatus.unit_of_work_id).filter( ContentStatus.content_status_id.in_(ids) ) uow_ids = [x[0] for x in query] query = self.session.query(UnitOfWork).filter(UnitOfWork.unit_of_work_id.in_(uow_ids)) count = query.update( { UnitOfWork.completeness_status: status, UnitOfWork.last_updated_at: utcnow(), }, synchronize_session='fetch', ) audit_log_method_mapping[status](uow_ids) self.session.commit() flash(f'{count} Unit(s) of Work were successfully completed.') except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to complete Unit of Work, details = %s', ex) flash('Failed to complete Unit of Work', category='error') @action( 'complete_uow', 'Complete Unit of Work', 'Are you sure you want to set Unit of Work completed?', ) def action_complete_uow(self, ids: List[str]) -> None: self._action_complete_uow(ids, CompletenessStatusEnum.COMPLETE) @action( 'min_complete_uow', 'Min Complete Unit of Work', 'Are you sure you want to set Unit of Work completed (min)?', ) def action_min_complete_uow(self, ids: List[str]) -> None: self._action_complete_uow(ids, CompletenessStatusEnum.MIN_COMPLETE) @action('reprocess', 'Reprocess Unit of Work(SLZ)') def action_reprocess(self, ids: List[str]) -> None: raw_priority = request.values.get('priority') try: if raw_priority == '': priority = None else: try: # Consider that the highest priority has the lowest integer value priority = UnitOfWorkPriorityEnum(int(raw_priority)) lowest = UnitOfWorkPriorityEnum.LOWEST.value highest = UnitOfWorkPriorityEnum.HIGHEST.value if not highest <= priority.value <= lowest: raise ValueError except (ValueError, TypeError): flash('Invalid value for Priority', category='error') return casted_ids = cast_ids_to_int(ids, logger) config: Config = current_app.config['CONFIG'] uow_codes = self.reprocessing_service.reprocess_slz_uows( casted_ids, priority, config.env ) self.audit_log_service.log_slz_unit_of_work_reprocessed(uow_codes) self.session.commit() flash(f'{len(uow_codes)} Unit(s) of Work were successfully marked for reprocessing') except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to perform reprocessing, details = %s', ex) flash('Failed to perform reprocessing, see logs for details', category='error') @action('schedule_from_idle', 'Schedule from Idle Unit of Work(SLZ)') def action_schedule_from_idle(self, ids: List[str]) -> None: try: subquery = self.session.query(ContentStatus.unit_of_work_id).filter( ContentStatus.content_status_id.in_(ids) ).subquery() query = self.session.query(UnitOfWork).filter( UnitOfWork.unit_of_work_id.in_(subquery), UnitOfWork.activity_status == ActivityStatusEnum.IDLE ) now = utcnow() count = query.update( { UnitOfWork.activity_status: ActivityStatusEnum.NOT_IN_PROGRESS, UnitOfWork.completeness_status: CompletenessStatusEnum.ACTIVE, UnitOfWork.created_at: now, UnitOfWork.last_updated_at: now, UnitOfWork.next_run_at: now + timedelta(minutes=1), UnitOfWork.is_force_complete: False, }, synchronize_session='fetch', ) self.session.commit() if count == 0: flash( 'No UoW in IDLE status were found', category='error', ) return flash(f'{count} Unit(s) of Work were successfully scheduled from idle') except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to perform scheduling from idle, details = %s', ex) flash('Failed to perform scheduling from idle, see logs for details', category='error') @action( 'reprocess_uow_exp', 'Reprocess Unit of Work(Exploration)', 'Are you sure you want to reprocess UOW(Exploration)?', ) def action_reprocess_uow_exp(self, ids: List[str]) -> None: try: query = self.session.query(ContentStatus.unit_of_work_id).filter( ContentStatus.content_status_id.in_(ids) ) uow_ids = [x[0] for x in query] content_statuses = self.session.query(ContentStatus).filter( ContentStatus.unit_of_work_id.in_(uow_ids), ContentStatus.content_status == ContentStatusEnum.COMPLETE ).all() if not content_statuses: flash( f'No content_statuses in {ContentStatusEnum.COMPLETE.value} ' f'status for selected UOW(s)', category='error', ) return config: Config = current_app.config['CONFIG'] errors = self.exp_reprocessing_service.reprocess_cs( content_statuses=content_statuses, config=config.reprocessing ) content_status_ids = [cs.content_status_id for cs in content_statuses] failed_ids = cast_ids_to_int(errors.keys(), logger) successfully_processed = [id_ for id_ in content_status_ids if id_ not in failed_ids] if successfully_processed: self.audit_log_service.log_exp_unit_of_work_reprocessed( uow_ids=uow_ids, cs_ids=successfully_processed ) self.session.commit() if errors: flash(f'Some content_statuses weren\'t reprocessed. Details: {errors}') else: flash( f'UOWs: {",".join(str(_id) for _id in uow_ids)} were successfully reprocessed' ) except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to perform UOW exploration reprocessing, details = %s', ex) flash( 'Failed to perform UOW exploration reprocessing, see logs for details', category='error' ) @action( 'reprocess_cs_exp', 'Reprocess Content Status(Exploration)', 'Are you sure you want to reprocess CS(Exploration)?', ) def action_reprocess_cs_exp(self, ids: List[str]) -> None: try: casted_ids = cast_ids_to_int(ids, logger) query = self.session.query(ContentStatus).filter( ContentStatus.content_status_id.in_(casted_ids) ) has_improper_status = self.session.query(literal(True)).filter( query.filter(ContentStatus.content_status != ContentStatusEnum.COMPLETE).exists() ).scalar() if has_improper_status: flash( f'Some items are not in {ContentStatusEnum.COMPLETE.value} status', category='error', ) return config: Config = current_app.config['CONFIG'] errors = self.exp_reprocessing_service.reprocess_cs( content_statuses=query.all(), config=config.reprocessing ) failed_ids = cast_ids_to_int(errors.keys(), logger) successfully_processed = [id_ for id_ in casted_ids if id_ not in failed_ids] if successfully_processed: self.audit_log_service.log_exp_content_status_reprocessed(successfully_processed) self.session.commit() if errors: flash(f'Some content_statuses weren\'t reprocessed. Details: {errors}') else: flash('Content statuses were successfully reprocessed') except Exception as ex: # pylint: disable=broad-except if not self.handle_view_exception(ex): raise logger.error('Failed to perform CS exploration reprocessing, details = %s', ex) flash( 'Failed to perform CS exploration reprocessing, see logs for details', category='error' )