# pylint: disable=no-member,too-many-lines from datetime import datetime, timedelta, timezone import pytest from db_schema.schemas.slz import ( ActivityStatusEnum, CompletenessStatusEnum, ContentMetadataStatusEnum, ContentStatus, ContentStatusEnum, ContentStatusMigrationLog, UnitOfWork, UnitOfWorkMigrationLog, UnitOfWorkPriorityEnum, ) from slz_storage.conv import UoWPGToDynamo pytestmark = [pytest.mark.integration] def test_migrated(db, now, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', reprocess_id='20191202T111111', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, created_at=now, last_updated_at=now + timedelta(minutes=5), next_run_at=now + timedelta(minutes=15), ) uow_1_migrated = UnitOfWorkMigrationLog(created_at=now) uow_1.migrated = uow_1_migrated db.connection.add(uow_1) uow_2 = UnitOfWork( unit_of_work_code='spotify-20191103-sme-sub_30_sec_streams-v2', reprocess_id='20191202T111111', licensor=db_licensors['sme'], report=db_reports['sub_30_sec_streams'], report_date='2019-11-03', version='v2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=True, priority=UnitOfWorkPriorityEnum.PRIORITY_8.value, created_at=now, last_updated_at=now + timedelta(minutes=5), next_run_at=now + timedelta(minutes=15), ) db.connection.add(uow_2) db.connection.commit() def test_get_uow_by_uow_id(db, unit_of_work_stubs): db.connection.add(unit_of_work_stubs[1]) db.connection.add(unit_of_work_stubs[2]) db.connection.commit() found = db.get_uow_by_uow_id('apple-20191124-theorchard-amEvent-v1_2-rp1') assert found == unit_of_work_stubs[1] def test_find_uow_to_cancel(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191202T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) uow_2 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191203T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_2) uow_3 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191204T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.CANCELLED.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_3) db.connection.commit() results = db.find_uow_to_cancel('apple-20191124-theorchard-amEvent-v1_2-rp20191202T111111') assert [uow_2] == results def test_update_uow_activity_status(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191202T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) db.connection.commit() db.update_uow_activity_status(uow_1, ActivityStatusEnum.IN_PROGRESS) found = db.connection.query(UnitOfWork).filter( UnitOfWork.unit_of_work_id == uow_1.unit_of_work_id ).one() assert found.activity_status == ActivityStatusEnum.IN_PROGRESS def test_cancel_uow(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191202T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) db.connection.commit() updated_time = datetime(2020, 1, 2, 3, 4, 5, 6, timezone.utc) db.cancel_uow(uow_1, updated_at=updated_time) found = db.connection.query(UnitOfWork).filter( UnitOfWork.unit_of_work_id == uow_1.unit_of_work_id ).one() assert found.activity_status == ActivityStatusEnum.NOT_IN_PROGRESS assert found.completeness_status == CompletenessStatusEnum.CANCELLED assert found.last_updated_at == updated_time def test_find_uows_by_completness_statuses(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20151124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='', report=db_reports['amEvent'], report_date='2015-11-24', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) uow_2 = UnitOfWork( unit_of_work_code='apple-20151124-sme-amEvent-v1_2', licensor=db_licensors['sme'], reprocess_id='', report=db_reports['amEvent'], report_date='2015-11-24', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) db.connection.add(uow_2) db.connection.commit() found = db.find_uows_by_completness_statuses(statuses=[ CompletenessStatusEnum.ACTIVE, ]) assert uow_2 in found assert uow_1 not in found def test_find_uows_by_uow_id(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20151125-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='', report=db_reports['amEvent'], report_date='2015-11-25', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) uow_2 = UnitOfWork( unit_of_work_code='apple-20151125-sme-amEvent-v1_2', licensor=db_licensors['sme'], reprocess_id='', report=db_reports['amEvent'], report_date='2015-11-25', version='v1_2', activity_status=ActivityStatusEnum.IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) db.connection.add(uow_2) db.connection.commit() second = db.get_uow_by_uow_id('apple-20151125-sme-amEvent-v1_2') found = db.find_uows_by_uow_id(uow_ids=[ second.unit_of_work_id, ]) assert uow_2 in found assert uow_1 not in found def test_update(db, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', licensor=db_licensors['theorchard'], reprocess_id='20191202T111111', report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.MIN_COMPLETE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) db.connection.add(uow_1) db.connection.commit() db.update(uow_1, activity_status=ActivityStatusEnum.IN_PROGRESS.value, version='v1_3') found = db.connection.query(UnitOfWork).filter( UnitOfWork.unit_of_work_id == uow_1.unit_of_work_id ).one() assert found.activity_status == ActivityStatusEnum.IN_PROGRESS assert found.version == 'v1_3' def test_find_content_status_same_day_of_week(db, now, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', reprocess_id='', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) uow_2 = UnitOfWork( unit_of_work_code='apple-20191110-theorchard-amEvent-v1_2', reprocess_id='', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-10', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) uow_3 = UnitOfWork( unit_of_work_code='apple-20191116-theorchard-amEvent-v1_2', reprocess_id='', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-16', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) uow_4 = UnitOfWork( unit_of_work_code='apple-20191201-theorchard-amEvent-v1_2', reprocess_id='', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-12-01', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) cs_1 = ContentStatus( context='US', content_name='US.txt', content_status=ContentStatusEnum.COMPLETE, content_size=100, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_2 = ContentStatus( context='MX', content_name='MX.txt', content_status=ContentStatusEnum.COMPLETE, created_at=now, content_size=200, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_3 = ContentStatus( context='US', content_name='US.txt', content_status=ContentStatusEnum.COMPLETE, content_size=300, failure_count=0, created_at=now, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_4 = ContentStatus( context='US', content_name='US.txt', content_status=ContentStatusEnum.FAILED, content_size=400, failure_count=0, created_at=now, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_5 = ContentStatus( context='US', content_name='US.txt', content_status=ContentStatusEnum.COMPLETE, content_size=500, failure_count=0, created_at=now, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.connection.add(uow_1) db.connection.add(uow_2) db.connection.add(uow_3) db.connection.add(uow_4) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) uow_3.content_statuses.append(cs_4) uow_4.content_statuses.append(cs_5) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.add(cs_4) db.connection.add(cs_5) db.connection.commit() css = db.find_content_status_same_day_of_week( cs_5, set([ContentStatusEnum.COMPLETE]), 3, ) assert [cs_1, cs_3] == css def test_find_available_content_statuses(db, now, db_licensors, db_reports): uow_1 = UnitOfWork( unit_of_work_code='apple-20191124-theorchard-amEvent-v1_2', reprocess_id='1', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-24', version='v1_2', activity_status=ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=CompletenessStatusEnum.ACTIVE.value, is_force_complete=False, priority=UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at='2019-11-07T01:10:38.377644+00:00', created_at='2019-11-07T01:01:00.599872+00:00', last_updated_at='2019-11-21T18:41:39.985786+00:00', ) cs_1 = ContentStatus( context='US', content_name='us.txt', content_status=ContentStatusEnum.ACTIVE, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_2 = ContentStatus( context='MX', content_name='mx.txt', content_status=ContentStatusEnum.MISSING, created_at=now, failure_count=0, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) cs_3 = ContentStatus( context='IT', content_name='it.txt', content_status=ContentStatusEnum.FAILED, failure_count=0, created_at=now, sub_content='{}', metadata_process_status=ContentMetadataStatusEnum.NOT_QUEUED.value, ) db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_1.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() results = db.find_available_content_statuses(uow_1) assert [cs_1, cs_2, cs_3] == results found = db.get_content_status(uow_1, context='IT') assert cs_3 == found css = db.find_content_statuses_with_statuses( uow_1, set([ContentStatusEnum.ACTIVE, ContentStatusEnum.FAILED, ContentStatusEnum.COMPLETE]) ) assert [cs_1, cs_3] == css def test_migrated_uow_with_content_statues(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) cs_1_migrated = ContentStatusMigrationLog(created_at=now) cs_1.migrated = cs_1_migrated db.connection.add(cs_1) db.connection.add(cs_2) db.connection.commit() def test_update_uow__content_statuses(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] uow_2 = unit_of_work_stubs[2] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] cs_3 = content_status_stubs[3] db.connection.add(uow_1) db.connection.add(uow_2) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() completed_at = now + timedelta(days=1, hours=2) db.update_uow__content_statuses( uow_1, failure_count=234, completed_at=completed_at, content_status=ContentStatusEnum.CANCELLED, metadata_process_status=ContentMetadataStatusEnum.QUEUED.value ) expected = set([cs_1.content_status_id, cs_2.content_status_id]) db.connection.close() css = db.connection.query(ContentStatus).filter( ContentStatus.failure_count == 234, ContentStatus.completed_at == completed_at, ContentStatus.content_status == ContentStatusEnum.CANCELLED.value, ContentStatus.metadata_process_status == ContentMetadataStatusEnum.QUEUED.value ).all() got = {cs.content_status_id for cs in css} assert got == expected def test_save_complete_volatile_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] uow_2 = unit_of_work_stubs[2] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] cs_3 = content_status_stubs[3] db.connection.add(uow_1) db.connection.add(uow_2) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() content_status_id = cs_1.content_status_id changeset = {'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context} assert db.save_active_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == content_status_id ).one() assert updated_cs.failure_count == 0 changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'content_name': 'new_content_name.gz', 'job_id': '123', 'size': 12345, 'subcontent': [1, 2, 3, 4], 'completed_at': (now + timedelta(days=1)).astimezone(timezone.utc), } assert db.save_complete_volatile_content_status(**changeset) db.connection.expire_all() db.connection.close() updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == content_status_id ).one() assert updated_cs.failure_count == 0 assert updated_cs.content_status == ContentStatusEnum.COMPLETE_VOLATILE def test_save_failed_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'job_id': '123', 'failure_code': '404', 'failure_description': 'Not Found', 'created_at': now + timedelta(days=1), } assert db.save_failed_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == cs_1.content_status_id ).one() assert updated_cs.failure_count == 1 assert updated_cs.content_status == ContentStatusEnum.FAILED created_failure_log_record = updated_cs.content_failure_logs[0] assert created_failure_log_record.job_id == changeset['job_id'] assert created_failure_log_record.failure_code == changeset['failure_code'] assert created_failure_log_record.failure_description == changeset['failure_description'] assert created_failure_log_record.created_at == changeset['created_at'] def test_save_failed_to_succeed_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] uow_2 = unit_of_work_stubs[2] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] cs_3 = content_status_stubs[3] db.connection.add(uow_1) db.connection.add(uow_2) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'job_id': '123', 'failure_code': '404', 'failure_description': 'Not Found', 'created_at': now + timedelta(days=1), } assert db.save_failed_content_status(**changeset) changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'content_name': 'new_content_name.gz', 'job_id': '123', 'size': 12345, 'subcontent': [1, 2, 3, 4], 'completed_at': (now + timedelta(days=1)).astimezone(timezone.utc), } assert db.save_succeed_content_status(**changeset) content_status_id = cs_1.content_status_id db.connection.close() updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == content_status_id ).one() assert updated_cs.content_status == ContentStatusEnum.COMPLETE assert updated_cs.latest_job_id == changeset['job_id'] assert updated_cs.content_size == changeset['size'] assert updated_cs.content_name == changeset['content_name'] assert updated_cs.sub_content == changeset['subcontent'] assert updated_cs.completed_at == changeset['completed_at'] assert updated_cs.failure_count == 1 def test_save_active_to_succeed_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] uow_2 = unit_of_work_stubs[2] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] cs_3 = content_status_stubs[3] db.connection.add(uow_1) db.connection.add(uow_2) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() content_status_id = cs_1.content_status_id changeset = {'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context} assert db.save_active_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == content_status_id ).one() assert updated_cs.failure_count == 0 changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'content_name': 'new_content_name.gz', 'job_id': '123', 'size': 12345, 'subcontent': [1, 2, 3, 4], 'completed_at': (now + timedelta(days=1)).astimezone(timezone.utc), } assert db.save_succeed_content_status(**changeset) db.connection.close() updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == content_status_id ).one() assert updated_cs.failure_count == 0 def test_save_active_content_status(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, } assert db.save_active_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == cs_1.content_status_id ).one() assert updated_cs.content_status == ContentStatusEnum.ACTIVE def test_save_on_hold_content_status(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'size': 123, 'meta_data': { 'version': '1.0', 'test': 'test' } } assert db.save_on_hold_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == cs_1.content_status_id ).one() assert updated_cs.content_status == ContentStatusEnum.ON_HOLD assert updated_cs.meta_data == changeset['meta_data'] assert updated_cs.content_size == 123 def test_save_missing_content_status(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] uow_2 = unit_of_work_stubs[2] cs_1 = content_status_stubs[1] cs_2 = content_status_stubs[2] cs_3 = content_status_stubs[3] db.connection.add(uow_1) db.connection.add(uow_2) uow_1.content_statuses.append(cs_1) uow_1.content_statuses.append(cs_2) uow_2.content_statuses.append(cs_3) db.connection.add(cs_1) db.connection.add(cs_2) db.connection.add(cs_3) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, } assert db.save_missing_content_status(**changeset) updated_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == cs_1.content_status_id ).one() assert updated_cs.content_status == ContentStatusEnum.MISSING def test_create_processing_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] db.connection.add(uow_1) assert len(uow_1.content_statuses) == 0 cs_1 = content_status_stubs[1] changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'content_name': cs_1.content_name, 'context': 'JP', 'content_status': ContentStatusEnum.ACTIVE, 'job_id': '12345', 'created_at': now, 'last_checked_at': now, } result = db.create_or_update_processing_content_status(**changeset) assert result created_cs = uow_1.content_statuses[0] assert created_cs.content_status == changeset['content_status'] assert created_cs.last_checked_at == changeset['last_checked_at'] assert created_cs.content_name == changeset['content_name'] assert created_cs.context == changeset['context'] assert created_cs.failure_count == 0 assert created_cs.created_at == changeset['created_at'] assert created_cs.latest_job_id == changeset['job_id'] assert created_cs.sub_content is None assert created_cs.metadata_process_status == ContentMetadataStatusEnum.NOT_QUEUED def test_update_processing_content_status(db, now, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] cs_1.content_status = ContentStatusEnum.MISSING db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'content_name': 'new content name', 'context': cs_1.context, 'content_status': ContentStatusEnum.ACTIVE, 'job_id': 'qwerty', 'created_at': now, 'last_checked_at': now + timedelta(minutes=45), } result = db.create_or_update_processing_content_status(**changeset) assert result cs_1_id = cs_1.content_status_id db.connection.close() created_cs = db.connection.query(ContentStatus).filter( ContentStatus.content_status_id == cs_1_id ).one() assert created_cs.latest_job_id == changeset['job_id'] assert created_cs.last_checked_at == changeset['last_checked_at'] def test_skip_completed_processing_content_status_complete( db, now, unit_of_work_stubs, content_status_stubs ): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] cs_1.content_status = ContentStatusEnum.COMPLETE db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'content_name': cs_1.content_name, 'context': 'US', 'content_status': ContentStatusEnum.ACTIVE, 'job_id': 'qwerty', 'created_at': now, 'last_checked_at': now + timedelta(minutes=45), } result = db.create_or_update_processing_content_status(**changeset) assert not result def test_skip_completed_processing_content_status_on_hold( db, now, unit_of_work_stubs, content_status_stubs ): uow_1 = unit_of_work_stubs[1] cs_1 = content_status_stubs[1] cs_1.content_status = ContentStatusEnum.ON_HOLD db.connection.add(uow_1) uow_1.content_statuses.append(cs_1) db.connection.add(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'content_name': cs_1.content_name, 'context': 'US', 'content_status': ContentStatusEnum.ACTIVE, 'job_id': 'qwerty', 'created_at': now, 'last_checked_at': now + timedelta(minutes=45), } result = db.create_or_update_processing_content_status(**changeset) assert not result def test_update_content_status_metadata_completed(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] db.connection.add(uow_1) cs_1 = content_status_stubs[1] uow_1.content_statuses.append(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'num_of_records': 2000, 'hash_of_records': 'asxcf4566', 'process_status': ContentMetadataStatusEnum.COMPLETED, } result = db.update_content_status_metadata_by_uow_id(**changeset) assert result uow_id = uow_1.unit_of_work_id db.connection.close() cs_updated = db.connection.query(ContentStatus).filter( ContentStatus.context == changeset['context'], ContentStatus.unit_of_work_id == uow_id, ).one() assert cs_updated.hash == changeset['hash_of_records'] assert cs_updated.record_count == changeset['num_of_records'] assert cs_updated.metadata_process_status == changeset['process_status'] def test_update_content_status_metadata_failed(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] db.connection.add(uow_1) cs_1 = content_status_stubs[1] uow_1.content_statuses.append(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'process_status': ContentMetadataStatusEnum.FAILED } result = db.update_content_status_metadata_by_uow_id(**changeset) assert result uow_id = uow_1.unit_of_work_id db.connection.close() cs_updated = db.connection.query(ContentStatus).filter( ContentStatus.context == changeset['context'], ContentStatus.unit_of_work_id == uow_id, ).one() assert cs_updated.metadata_process_status == changeset['process_status'] assert cs_updated.hash is None assert cs_updated.record_count is None def test_update_content_status_metadata_queued(db, unit_of_work_stubs, content_status_stubs): uow_1 = unit_of_work_stubs[1] db.connection.add(uow_1) cs_1 = content_status_stubs[1] uow_1.content_statuses.append(cs_1) db.connection.commit() changeset = { 'uow_id': UoWPGToDynamo.get_uow_id(uow_1), 'context': cs_1.context, 'process_status': ContentMetadataStatusEnum.QUEUED, } result = db.update_content_status_metadata_by_uow_id(**changeset) assert result uow_id = uow_1.unit_of_work_id db.connection.close() cs_updated = db.connection.query(ContentStatus).filter( ContentStatus.context == changeset['context'], ContentStatus.unit_of_work_id == uow_id, ).one() assert cs_updated.metadata_process_status == changeset['process_status'] assert cs_updated.metadata_process_started_at is not None assert cs_updated.metadata_process_completed_at is None