# pylint: disable=redefined-outer-name import logging import os from datetime import datetime, timedelta, timezone from time import sleep from unittest.mock import Mock import pytest from db_schema.postgres.connection import Connection from db_schema.schemas import slz from freezegun import freeze_time from slz_config.pipeline_config.service import PipelineConfigService from smelog.entities import LoggerConfig from smelog.factory import LoggerFactory from . import FIXTURES_PATH def creds_loader_writer(): return dict( active_endpoint=os.environ.get('PG_HOST', '0.0.0.0'), port=os.environ.get('PG_PORT', 5432), database=os.environ.get('PG_DB', 'slz'), username=os.environ.get('PG_USER', 'admin'), password=os.environ.get('PG_PASSWORD', 'admin'), ) def creds_loader_reader(): return dict( active_endpoint=os.environ.get('PG_HOST', '0.0.0.0'), port=os.environ.get('PG_PORT', 5432), database=os.environ.get('PG_DB', 'slz'), username=os.environ.get('PG_USER_RO', 'slz_reader'), password=os.environ.get('PG_PASSWORD_RO', 'slz_reader'), ) @pytest.fixture(scope='session') def db() -> Connection: conn = Connection(credentials_loader=creds_loader_writer, name='jm_lambda_test', version='v1') # Simple hack for waiting migrations to be completed while True: result = conn.session.execute( # pylint: disable=no-member 'SELECT id FROM databasechangelog ORDER BY dateexecuted DESC LIMIT 1' ) if result.first(): break sleep(1) yield conn @pytest.fixture(scope='function') def clean_db(db): yield db.session.rollback() for model in [slz.UnitOfWorkPipeline, slz.UnitOfWorkGroup, slz.ContentStatus, slz.UnitOfWork]: db.session.query(model).delete() db.session.commit() db.disconnect() @pytest.fixture(scope='session') def db_ro() -> Connection: conn = Connection( credentials_loader=creds_loader_reader, name='jm_lambda_test_ro', version='v1' ) yield conn @pytest.fixture(scope='function') def db_licensors(db): return {item.licensor_name: item for item in db.session.query(slz.Licensor).all()} @pytest.fixture(scope='function') def db_reports(db): return {item.report_name: item for item in db.session.query(slz.Report).all()} @pytest.fixture(scope='function') def db_data_sources(db): return {item.data_source_name: item for item in db.session.query(slz.DataSource).all()} @pytest.fixture(scope='function') def logger_test(): log_config = LoggerConfig( name='TEST', version='1', level=logging.DEBUG, environment='test', is_local=True, ) return LoggerFactory(log_config).get_logger('TEST') @pytest.fixture(scope='function') def merged_config_service_test(): dsp_config = Mock() dsp_config.path.return_value = os.path.join(FIXTURES_PATH, 'dsp_config.json') dsp_complete_criteria = Mock() dsp_complete_criteria.path.return_value = os.path.join( FIXTURES_PATH, 'dsp-complete-criteria.json' ) dsp_specific_settings = Mock() dsp_specific_settings.path.return_value = os.path.join( FIXTURES_PATH, 'dsp-specific-settings.json' ) return PipelineConfigService( logger=Mock(), dsp_config_path=dsp_config, dsp_complete_criteria_config_path=dsp_complete_criteria, dsp_specific_settings_path=dsp_specific_settings, ) @pytest.fixture(scope='function') def now_freezed(): with freeze_time(datetime(2021, 11, 24, 3, 21, 34, tzinfo=timezone.utc)): yield datetime.now().astimezone(timezone.utc) @pytest.fixture(scope='function') def unit_of_work_stubs(db_licensors, now_freezed, db_reports): return { 1: dict( 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=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at=now_freezed + timedelta(minutes=15), created_at=now_freezed, last_updated_at=now_freezed, ), 2: dict( unit_of_work_code='apple-20191125-theorchard-amEvent-v1_2', reprocess_id='1', licensor=db_licensors['theorchard'], report=db_reports['amEvent'], report_date='2019-11-25', version='v1_2', activity_status=slz.ActivityStatusEnum.NOT_IN_PROGRESS.value, completeness_status=slz.CompletenessStatusEnum.COMPLETE.value, is_force_complete=False, priority=slz.UnitOfWorkPriorityEnum.DEFAULT.value, next_run_at=now_freezed + timedelta(minutes=15), created_at=now_freezed, last_updated_at=now_freezed, ), } @pytest.fixture(scope='function') def content_status_stubs(now_freezed): return { 1: dict( context='US', content_name='us.txt', content_status=slz.ContentStatusEnum.COMPLETE, created_at=now_freezed, failure_count=0, sub_content='{}', metadata_process_status=slz.ContentMetadataStatusEnum.NOT_QUEUED.value, ), 2: dict( context='MX', content_name='mx.txt', content_status=slz.ContentStatusEnum.MISSING, created_at=now_freezed, failure_count=0, sub_content='{}', metadata_process_status=slz.ContentMetadataStatusEnum.NOT_QUEUED.value, ), 3: dict( context='IT', content_name='it.txt', content_status=slz.ContentStatusEnum.COMPLETE, failure_count=0, created_at=now_freezed, sub_content='{}', metadata_process_status=slz.ContentMetadataStatusEnum.NOT_QUEUED.value, ), }