import argparse import os import unittest from datetime import datetime, timezone, timedelta from unittest.mock import Mock import boto3 import smart_open from db_schema.schemas import slz from moto import mock_s3 from parameterized import parameterized from slz_storage.postgres import connection from ae_backfill.handler import ( queue_uows, count_almost_completed, copy_files_for_completed, count_currently_processing, complete_when_all_volatile, ) def get_uow(report_date: datetime) -> slz.UnitOfWork: return slz.UnitOfWork( unit_of_work_code=f'appreciationengine-{report_date.strftime("%Y%m%d")}-sme-activityfeed-v1', reprocess_id='', report_date=f'{report_date.strftime("%Y-%m-%d")}', report_id=106, licensor_id=1, version='v1', activity_status=slz.ActivityStatusEnum.IDLE.value, completeness_status=slz.CompletenessStatusEnum.ACTIVE.value, unit_of_work_type=slz.UnitOfWorkTypeEnum.BACKFILL, next_run_at='2021-01-28 01:00:00.715337', created_at='2021-01-28 01:00:00.715337', last_updated_at='2021-01-28 01:00:00.715337', is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.PRIORITY_6.value ) def get_cs() -> slz.ContentStatus: return slz.ContentStatus( content_name='test.json.gz', context='context', content_status=slz.ContentStatusEnum.ON_HOLD, failure_count=0, latest_job_id='backfill', created_at='2021-01-28 01:00:00.715337', last_checked_at='2021-01-28 01:00:00.715337', sub_content=None, metadata_process_status=slz.ContentMetadataStatusEnum.COMPLETED, ) class BackfillProcessCase(unittest.TestCase): def setUp(self): self.logger = Mock() self.pg = connection.get_session( host=os.environ.get('PG_HOST', '0.0.0.0'), port=os.environ.get('PG_PORT', 5432), db=os.environ.get('PG_DB', 'slz'), user=os.environ.get('PG_USER', 'admin'), password=os.environ.get('PG_PASSWORD', 'admin'), ) def tearDown(self): self.pg.rollback() self.pg.query(slz.ContentStatus).delete() self.pg.query(slz.UnitOfWork).delete() self.pg.commit() self.pg.close() def test_queue_uows(self): params = argparse.Namespace( start_date=datetime(year=2021, month=1, day=30, tzinfo=timezone.utc) ) now = datetime(year=2021, month=2, day=3, tzinfo=timezone.utc) # create backfill units for i in range(10): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=1*i) ) content_status = get_cs() unit_of_work.content_statuses.append(content_status) self.pg.add(unit_of_work) self.pg.commit() # check that initial data is correct uow_not_inprogress = self.pg.query(slz.UnitOfWork).filter( slz.UnitOfWork.activity_status == slz.ActivityStatusEnum.NOT_IN_PROGRESS.value).all() self.assertEqual(len(uow_not_inprogress), 0) cs_missing = self.pg.query(slz.ContentStatus).filter( slz.ContentStatus.content_status != slz.ContentStatusEnum.ON_HOLD).all() self.assertEqual(len(cs_missing), 0) queue_uows(count=2, pg=self.pg, now=now, logger=self.logger, params=params, ctx={}) uow_not_inprogress = self.pg.query(slz.UnitOfWork).filter( slz.UnitOfWork.activity_status == slz.ActivityStatusEnum.NOT_IN_PROGRESS.value).all() self.assertEqual(len(uow_not_inprogress), 2) self.assertEqual(uow_not_inprogress[0].created_at, now) self.assertTrue(all( cs.content_status == slz.ContentStatusEnum.MISSING for cs in uow_not_inprogress[0].content_statuses )) @parameterized.expand([ (0.5, 5), (0.75, 5), (0.8, 0), ]) def test_find_almost_completed(self, completeness_rate, expected): params = argparse.Namespace( start_date=datetime(year=2021, month=1, day=30, tzinfo=timezone.utc), completeness_rate=completeness_rate ) # create backfill uow for i in range(5): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) + timedelta(days=1 * i) ) self.pg.add(unit_of_work) # create units which are being processed(5 backfill, 5 usual) for i in range(10): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=1 * i) ) unit_of_work.activity_status = slz.ActivityStatusEnum.IN_PROGRESS.value if i >= 5: unit_of_work.unit_of_work_type = slz.UnitOfWorkTypeEnum.DAILY unit_of_work.next_run_at = '2021-01-28 02:00:00.715337' # create cs: 3-completed, 1 onhold for ii in range(3): content_status = get_cs() content_status.context = f'test{ii}' content_status.content_status = slz.ContentStatusEnum.COMPLETE unit_of_work.content_statuses.append(content_status) unit_of_work.content_statuses.append(get_cs()) self.pg.add(unit_of_work) self.pg.commit() count = count_almost_completed(params=params, pg=self.pg) self.assertEqual(count, expected) def test_complete_volatile(self): params = argparse.Namespace( start_date=datetime(year=2021, month=1, day=30, tzinfo=timezone.utc), completeness_rate=0 ) # create backfill uow for i in range(5): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) + timedelta(days=1 * i) ) self.pg.add(unit_of_work) # create units: 1 backfill unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=1) ) unit_of_work.activity_status = slz.ActivityStatusEnum.NOT_IN_PROGRESS.value # if i >= 5: # unit_of_work.unit_of_work_type = slz.UnitOfWorkTypeEnum.DAILY unit_of_work.next_run_at = '2021-01-28 02:00:00.715337' # create cs: 3-completed, 1 onhold for ii in range(3): content_status = get_cs() content_status.context = f'test{ii}' content_status.content_status = slz.ContentStatusEnum.COMPLETE unit_of_work.content_statuses.append(content_status) unit_of_work.content_statuses.append(get_cs()) # create units: 1 backfill unit all complete volatile unit_of_work_completed: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=1) ) unit_of_work_completed.activity_status = slz.ActivityStatusEnum.NOT_IN_PROGRESS.value unit_of_work_completed.next_run_at = '2021-01-28 02:00:00.715337' # create cs: 3-complete volatile for ii in range(3): content_status = get_cs() content_status.context = f'test{ii}' content_status.content_status = slz.ContentStatusEnum.COMPLETE_VOLATILE unit_of_work_completed.content_statuses.append(content_status) self.pg.add(unit_of_work_completed) self.pg.commit() # create units: 1 backfill unit with complete + complete volatile unit_of_work_completed_2: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=3) ) unit_of_work_completed_2.activity_status = slz.ActivityStatusEnum.NOT_IN_PROGRESS.value unit_of_work_completed_2.next_run_at = '2021-01-28 02:00:00.715337' # create cs: 3-completed, 1 complete volatile for ii in range(4): content_status = get_cs() content_status.context = f'test{ii}' if ii >= 3: content_status.content_status = slz.ContentStatusEnum.COMPLETE else: content_status.content_status = slz.ContentStatusEnum.COMPLETE_VOLATILE unit_of_work_completed_2.content_statuses.append(content_status) self.pg.add(unit_of_work_completed_2) self.pg.commit() # create units: 1 backfill unit just scheduled without content statuses unit_of_work_3: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=2) ) unit_of_work_3.activity_status = slz.ActivityStatusEnum.NOT_IN_PROGRESS.value unit_of_work_3.next_run_at = '2021-01-28 02:00:00.715337' self.pg.add(unit_of_work_3) self.pg.commit() complete_when_all_volatile(Mock(), params=params, pg=self.pg) assert unit_of_work.is_force_complete is False assert unit_of_work.completeness_status != slz.CompletenessStatusEnum.COMPLETE assert unit_of_work_completed.is_force_complete is True assert unit_of_work_completed.completeness_status == slz.CompletenessStatusEnum.COMPLETE assert unit_of_work_completed_2.is_force_complete is True assert unit_of_work_completed_2.completeness_status == slz.CompletenessStatusEnum.COMPLETE assert unit_of_work_3.is_force_complete is False assert unit_of_work_3.completeness_status != slz.CompletenessStatusEnum.COMPLETE def test_find_in_processing(self): params = argparse.Namespace( start_date=datetime(year=2021, month=1, day=30, tzinfo=timezone.utc), ) # create backfill uow for i in range(5): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=1, tzinfo=timezone.utc) + timedelta(days=1 * i) ) self.pg.add(unit_of_work) # create units which are being processed for i in range(10): unit_of_work: slz.UnitOfWork = get_uow( report_date=datetime(year=2021, month=1, day=26, tzinfo=timezone.utc) + timedelta(days=1 * i) ) unit_of_work.activity_status = slz.ActivityStatusEnum.IN_PROGRESS.value unit_of_work.next_run_at = '2021-01-28 02:00:00.715337' self.pg.add(unit_of_work) self.pg.commit() count = count_currently_processing(params=params, pg=self.pg) self.assertEqual(count, 5) def _create_file(self, date: str): with smart_open.open( f's3://dev-test-source-bucket/appreciationengine/activityfeed/v1/report_date={date}/report_licensor=sme/test.txt', 'wb', ignore_ext=True ) as fout: fout.write(b'Test') @mock_s3 def test_copy_already_completed(self): boto3.setup_default_session() conn = boto3.resource('s3', region_name='us-east-1') client = boto3.client('s3', region_name='us-east-1') source_bucket = 'dev-test-source-bucket' dest_bucket = 'dev-test-dest-bucket' conn.create_bucket(Bucket=source_bucket) conn.create_bucket(Bucket=dest_bucket) params = argparse.Namespace( start_date=datetime(year=2021, month=1, day=30, tzinfo=timezone.utc), source_bucket=source_bucket, dest_bucket=dest_bucket, ) now = datetime(year=2021, month=2, day=3, tzinfo=timezone.utc) # completed less than 2 days ago backfill uow unit_of_work = get_uow(datetime(year=2021, month=1, day=25, tzinfo=timezone.utc)) unit_of_work.last_updated_at = now - timedelta(days=1) unit_of_work.completeness_status = slz.CompletenessStatusEnum.COMPLETE self.pg.add(unit_of_work) self._create_file(unit_of_work.report_date) # completed more than 2 days ago backfill uow unit_of_work = get_uow(datetime(year=2021, month=1, day=27, tzinfo=timezone.utc)) unit_of_work.last_updated_at = now - timedelta(days=4) unit_of_work.completeness_status = slz.CompletenessStatusEnum.COMPLETE self.pg.add(unit_of_work) self._create_file(unit_of_work.report_date) # completed less than 2 days ago regular uow unit_of_work = get_uow(datetime(year=2021, month=2, day=2, tzinfo=timezone.utc)) unit_of_work.last_updated_at = now - timedelta(days=1) unit_of_work.completeness_status = slz.CompletenessStatusEnum.COMPLETE self.pg.add(unit_of_work) self._create_file(unit_of_work.report_date) self.pg.commit() # check that files were created successfully objects = client.list_objects_v2(Bucket=source_bucket) self.assertEqual(len(objects['Contents']), 3) objects = client.list_objects_v2(Bucket=dest_bucket) self.assertEqual(objects['KeyCount'], 0) copy_files_for_completed( logger=self.logger, params=params, pg=self.pg, now=now, ctx={}, client=client, resource=conn, ) # check that files were copied objects = client.list_objects_v2(Bucket=source_bucket) self.assertEqual(len(objects['Contents']), 2) objects = client.list_objects_v2(Bucket=dest_bucket) self.assertEqual(len(objects['Contents']), 1)