import argparse import datetime import os from unittest import mock from unittest.mock import MagicMock, patch, call from freezegun import freeze_time import pytest from feed_ingestion.util import dates_util from feed_ingestion.util import exec_utils from feed_ingestion.util.dates_util import Weekday from feed_ingestion.util.exec_utils import FlowExecBase from tests.testing_utils import (parametrize_by_dicts, raises_optionally) @pytest.mark.parametrize( 'date_str, is_valid', [ ('2024-01-01', True), ('2024-12-31', True), ('2024-13-31', False), ('2024-01-32', False), ('202-01-01', False), ('12-12', False), ('2024/12/12', False), ('string_value', False), ] ) def test_valid_date(date_str, is_valid): try: exec_utils.argparse_type_date(date_str) if not is_valid: pytest.fail(f'Expected ValueError for {date_str}') except BaseException as e: if is_valid: pytest.fail(f'Unexpected exception for {date_str}: {e}') @parametrize_by_dicts( 'kwargs, argv, expected, raises, raises_match', [ dict( kwargs=dict(), argv=['--date', '2024-01-01'], expected=datetime.date(2024, 1, 1) ), dict( kwargs=dict(required=False), argv=[], expected=None ), dict( kwargs=dict(required=True), argv=[], raises=argparse.ArgumentError ), dict( kwargs=dict(), argv=['--date=2024-12-31'], expected=datetime.date(2024, 12, 31) ), dict( case='using allow_abbrev - 1', kwargs=dict(), argv=['--dat=2024-12-31'], expected=datetime.date(2024, 12, 31) ), dict( case='using allow_abbrev - 1', kwargs=dict(), argv=['--da', '2024-01-01'], expected=datetime.date(2024, 1, 1) ), dict( case='not valid date', kwargs=dict(), argv=['--date', '2024-13-01'], raises=argparse.ArgumentError, raises_match="Invalid date: '2024-13-01'. " 'Expected format: YYYY-MM-DD' ), ] ) def test_add_arg_date( kwargs, argv, expected, raises, raises_match, parser): with raises_optionally(raises, raises_match): exec_utils.add_arg_date( parser, name='date', **kwargs ) args = parser.parse_args(argv) assert args.date == expected @parametrize_by_dicts( 'kwargs, cli, env, expected, raises, raises_match', [ dict( case='no default', kwargs=dict( ), cli='', expected=argparse.Namespace( reload=None, ) ), dict( case='cli true', cli='--reload True', expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ) ), dict( case='cli true 2', cli='--reload true', expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ) ), dict( case='cli double entry - take the last one', cli='--reload true --reload false', expected=argparse.Namespace( reload='False', _context={'reload': 'False'}, ) ), dict( case='cli true 3', cli='--reload=TRUE', expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ) ), dict( case='cli false', cli='--reload=False', expected=argparse.Namespace( reload='False', _context={'reload': 'False'}, ) ), dict( case='cli false 2', cli='--reload false', expected=argparse.Namespace( reload='False', _context={'reload': 'False'}, ) ), dict( case='cli empty', cli='--reload=', expected=argparse.Namespace( reload=None, ), ), dict( case='cli overrides env', cli='--reload=True', env={'RELOAD': 'false'}, expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ), ), dict( case='env true', env={'RELOAD': 'true'}, expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ), ), dict( case='env true 2', env={'RELOAD': 'TRUE'}, expected=argparse.Namespace( reload='True', _context={'reload': 'True'}, ), ), dict( case='env false', env={'RELOAD': 'false'}, expected=argparse.Namespace( reload='False', _context={'reload': 'False'}, ), ), dict( case='env false 2', env={'RELOAD': 'False'}, expected=argparse.Namespace( reload='False', _context={'reload': 'False'}, ), ), dict( case='env empty', env={'RELOAD': ''}, expected=argparse.Namespace( reload=None, ), ), ] ) def test_add_arg_bool( kwargs, cli, env, expected, raises, raises_match, parser, monkeypatch): cli = cli or '' argv = cli.split() if env: monkeypatch.setattr(os, 'environ', env) init_kwargs = dict( name='reload', context_key='reload', env_var='RELOAD', ) if kwargs: init_kwargs.update(kwargs) with raises_optionally(raises, raises_match): exec_utils.add_arg_bool(parser, **init_kwargs) returned = parser.parse_args(argv) assert returned == expected def test_add_wait_until_complete_arg(mocker, parser): mock = mocker.patch.object(exec_utils, 'add_arg_bool') exec_utils.add_wait_until_complete_arg(parser) assert mock.call_args_list == [ call( parser=parser, name='wait_until_complete', context_key='wait_until_complete', env_var='WAIT_UNTIL_COMPLETE', default=False, required=False, help='If "True" then wait until finish execution.' ) ] @parametrize_by_dicts( 'cli, env, expected, raises', [ dict( case='env - empty value', env={ 'CONTEXT_DATE': '' }, expected=argparse.Namespace( context_date=None, ) ), dict( case='cli - empty value', cli='--date=', expected=argparse.Namespace( context_date=None, ) ), dict( case='cli - empty value (2)', cli='--context-date=', expected=argparse.Namespace( context_date=None, ) ), dict( case='no env no cli', expected=argparse.Namespace( context_date=None ), ), dict( case='env - valid value', env={ 'CONTEXT_DATE': '2025-12-04' }, expected=argparse.Namespace( context_date=datetime.date(2025, 12, 4), _context={'context_date': '2025-12-04'} ) ), dict( case='cli - valid value', cli='--date 2024-01-01', expected=argparse.Namespace( context_date=datetime.date(2024, 1, 1), _context={'context_date': '2024-01-01'} ) ), dict( case='cli - valid value (2)', cli='--date=2024-03-01', expected=argparse.Namespace( context_date=datetime.date(2024, 3, 1), _context={'context_date': '2024-03-01'} ) ), dict( case='cli - valid value (3)', cli='--context-date 2024-01-02', expected=argparse.Namespace( context_date=datetime.date(2024, 1, 2), _context={'context_date': '2024-01-02'} ) ), dict( case='cli - valid value (4)', cli='--context-date=2024-01-06', expected=argparse.Namespace( context_date=datetime.date(2024, 1, 6), _context={'context_date': '2024-01-06'} ) ), dict( case='cli - invalid value', cli='--context-date=2024-01-36', raises=argparse.ArgumentError, ), ] ) def test_add_context_date_arg( cli, env, expected, raises, monkeypatch, parser): argv = cli or '' argv = argv.split() if env: monkeypatch.setattr('os.environ', env) with raises_optionally(raises): exec_utils.add_context_date_arg(parser) args = parser.parse_args(argv) assert args == expected @parametrize_by_dicts( 'kwargs, cli, env, expected, raises, raises_match', [ dict( case='all defaults', kwargs=dict( ), cli=[], expected=argparse.Namespace( weekday=None, ), ), dict( case='use name', kwargs=dict( name='weekday_special' ), expected=argparse.Namespace( weekday_special=None, ), ), dict( case='use name - cli', kwargs=dict( name='weekday_special' ), cli=['--weekday-special=monday'], expected=argparse.Namespace( weekday_special=dates_util.Weekday.MONDAY, ), ), dict( case='use name - env', kwargs=dict( name='weekday_special' ), cli=[], env={'WEEKDAY_SPECIAL': 'tuesday'}, expected=argparse.Namespace( weekday_special=Weekday.TUESDAY, ), ), dict( case='use name - cli overrides env', kwargs=dict( name='weekday_special' ), cli=['--weekday-special=friday'], env={'WEEKDAY_SPECIAL': 'tuesday'}, expected=argparse.Namespace( weekday_special=Weekday.FRIDAY, ), ), dict( case='default name - cli', kwargs=dict( ), cli=['--weekday=monday'], expected=argparse.Namespace( weekday=Weekday.MONDAY, ), ), dict( case='default name - env', kwargs=dict( ), cli=[], env={'WEEKDAY': 'tuesday'}, expected=argparse.Namespace( weekday=Weekday.TUESDAY, ), ), dict( case='default name - cli overrides env', kwargs=dict( ), cli=['--weekday=friday'], env={'WEEKDAY': 'tuesday'}, expected=argparse.Namespace( weekday=Weekday.FRIDAY, ), ), dict( case='default name default value', kwargs=dict( default='monday', required=False, ), cli=[], expected=argparse.Namespace( weekday=Weekday.MONDAY, ), ), dict( case='WRONG default', kwargs=dict( default='default_value', ), cli=[], raises=argparse.ArgumentError, raises_match=r"Invalid weekday: 'default_value'. " r"Expected one of: \['Monday', ", ), dict( case='required but not provided - should fail', kwargs=dict( default=None, required=True, ), cli=[], raises=argparse.ArgumentError, raises_match='required: --weekday' ), dict( case='doubled cli arg - the last one should override', kwargs=dict( default=None, required=True, ), cli=['--weekday=Monday', '--weekday', 'friday'], expected=argparse.Namespace( weekday=Weekday.FRIDAY, ), ), ] ) def test_add_weekday_arg( kwargs, cli, env, expected, raises, raises_match, monkeypatch, parser): if env: monkeypatch.setattr(os, 'environ', env) cli = cli or [] with raises_optionally(raises, raises_match): exec_utils.add_weekday_arg(parser, **kwargs) args = parser.parse_args(cli) assert args == expected @parametrize_by_dicts( 'kwargs, cli, env, expected, raises, raises_match', [ dict( case='default arguments - required', kwargs=dict( ), raises=argparse.ArgumentError, raises_match='required: --report' ), dict( case='no input - with default', kwargs=dict( default='report_default', required=False, ), expected=argparse.Namespace( report='report_default', _context={'report_name': 'report_default'} ), ), dict( case='use name - no default', kwargs=dict( name='report_special', required=False, ), expected=argparse.Namespace( report_special=None, ), ), dict( case='use name - use default', kwargs=dict( name='report_special', required=False, default='report_default' ), expected=argparse.Namespace( report_special='report_default', _context={'report_name': 'report_default'} ), marks=pytest.mark.xfail( reason='default value should be validated with choices', raises=argparse.ArgumentError ), ), dict( case='cli overrides default value', kwargs=dict( default='report2', required=False, ), cli=['--report=report3'], expected=argparse.Namespace( report='report3', _context={'report_name': 'report3'} ), ), dict( case='cli 1', kwargs=dict( ), cli=['--report=report1'], expected=argparse.Namespace( report='report1', _context={'report_name': 'report1'} ), ), dict( case='cli 2', kwargs=dict( ), cli=['--report', 'report2'], expected=argparse.Namespace( report='report2', _context={'report_name': 'report2'} ), ), dict( case='cli 3', kwargs=dict( ), cli=['--report-name=report3'], expected=argparse.Namespace( report='report3', _context={'report_name': 'report3'} ), ), dict( case='cli 4', kwargs=dict( ), cli=['--report-name', 'report4'], expected=argparse.Namespace( report='report4', _context={'report_name': 'report4'} ), ), dict( case='use context_key', kwargs=dict( context_key='report_context_key' ), cli=['--report', 'report4'], expected=argparse.Namespace( report='report4', _context={'report_context_key': 'report4'} ), ), dict( case='env - default env_var', kwargs=dict( ), cli=[], env={'REPORT': 'report3'}, expected=argparse.Namespace( report='report3', _context={'report_name': 'report3'} ), ), dict( case='env - specific env_var', kwargs=dict( env_var='REPORT_ENV' ), cli=[], env={'REPORT': 'report2', 'REPORT_ENV': 'report3'}, expected=argparse.Namespace( report='report3', _context={'report_name': 'report3'} ), ), dict( case='cli overrides env', kwargs=dict( ), cli=['--report=report1'], env={'REPORT': 'report3'}, expected=argparse.Namespace( report='report1', _context={'report_name': 'report1'} ), ), dict( case='required but not provided - should fail', kwargs=dict( required=True, ), cli=[], raises=argparse.ArgumentError, raises_match='required: --report' ), dict( case='incorrect value', kwargs=dict( required=True, ), cli=['--report=not_valid'], raises=argparse.ArgumentError, raises_match='invalid choice', ), dict( case='doubled cli arg - the last one should override', kwargs=dict( default=None, required=True, ), cli=['--report=report1', '--report=report2'], expected=argparse.Namespace( report='report2', _context={'report_name': 'report2'} ), ), ] ) def test_add_report_arg( kwargs, cli, env, expected, raises, raises_match, monkeypatch, parser): if env: monkeypatch.setattr(os, 'environ', env) cli = cli or [] kwargs.setdefault('choices', ['report1', 'report2', 'report3', 'report4']) with raises_optionally(raises, raises_match): exec_utils.add_report_arg(parser, **kwargs) args = parser.parse_args(cli) assert args == expected @pytest.fixture() def parser(monkeypatch): """Tweak the parser to use in test cases.""" class RawFormatter(argparse.HelpFormatter): """Custom formatter class that disables help formatting.""" def _format_usage(self, usage, actions, groups, prefix): return usage def _format_action(self, action): return action.help or '' def _format_action_invocation(self, action): return action.dest parser = exec_utils.CustomArgumentParser( formatter_class=RawFormatter ) def error(message): """Raise an exception instead of exiting.""" raise argparse.ArgumentError(None, message) monkeypatch.setattr( parser, 'error', error ) yield parser @parametrize_by_dicts( 'kwargs, cli, env, expected, raises, raises_match', [ dict( case='not required - no value - use default', kwargs=dict( default='default_value', name='arg_name', required=False, ), cli=[], expected=argparse.Namespace( arg_name='default_value', ), ), dict( case='not required - no value - no default', kwargs=dict( name='arg_name', ), cli=[], expected=argparse.Namespace( arg_name=None, ), ), dict( case='required and default - should fail', kwargs=dict( default='default_value', name='arg_name', required=True, ), raises=ValueError, raises_match='Cannot have both default and required', ), dict( case='required but not provided - should fail', kwargs=dict( default=None, name='arg_name', env_var='ARG_NAME', required=True, ), cli=[], raises=argparse.ArgumentError, raises_match='required: --arg-name', ), dict( case='value by cli', kwargs=dict( default='default_value', env_var='ARG_NAME', name='arg_name', ), cli=['--arg-name', 'cli_value'], env={}, expected=argparse.Namespace( arg_name='cli_value', ), ), dict( case='value by cli (another way)', kwargs=dict( default='default_value', env_var='ARG_NAME', name='arg_name', ), cli=['--arg-name=cli_value'], env={}, expected=argparse.Namespace( arg_name='cli_value', ), ), dict( case='value by env', kwargs=dict( default='default_value', env_var='ARG_NAME', name='arg_name', ), cli=[], env={'ARG_NAME': 'env_value'}, expected=argparse.Namespace( arg_name='env_value', ), ), dict( case='value by cli overrides env', kwargs=dict( default='default_value', env_var='ARG_NAME', name='arg_name', ), cli=['--arg-name', 'cli_value'], env={'ARG_NAME': 'env_value'}, expected=argparse.Namespace( arg_name='cli_value', ), ), dict( case='choices valid value from cli', kwargs=dict( env_var='ARG_NAME', name='arg_name', choices=['choice1', 'choice2'], ), cli=['--arg-name=choice2'], env={}, expected=argparse.Namespace( arg_name='choice2', ), ), dict( case='choices valid value from env', kwargs=dict( env_var='ARG_NAME', name='arg_name', choices=['choice1', 'choice2'], ), cli=[], env={'ARG_NAME': 'choice1'}, expected=argparse.Namespace( arg_name='choice1', ), ), dict( case='choices invalid value from cli - should fail', kwargs=dict( env_var='ARG_NAME', name='arg_name', choices=['choice1', 'choice2'], ), cli=['--arg-name=choice3'], env={}, raises=argparse.ArgumentError, raises_match="invalid choice: 'choice3'", ), dict( case='choices invalid value from env - should fail', kwargs=dict( env_var='ARG_NAME', name='arg_name', choices=['choice1', 'choice2'], ), cli=[], env={'ARG_NAME': 'choice4'}, raises=argparse.ArgumentError, raises_match="invalid choice: 'choice4'", ), dict( case='choices required but not provided - should fail', kwargs=dict( env_var='ARG_NAME', name='arg_name', choices=['choice1', 'choice2'], required=True, ), raises=argparse.ArgumentError, raises_match='required: --arg-name', ), dict( case='set _context by cli', kwargs=dict( name='arg_name', env_var='ARG_NAME', context_key='arg_context_key' ), cli=['--arg-name', 'cli_value'], env={'ARG_NAME': 'env_value'}, expected=argparse.Namespace( arg_name='cli_value', _context={'arg_context_key': 'cli_value'}, ), ), dict( case='set _context by env', kwargs=dict( name='arg_name', env_var='ARG_NAME', context_key='arg_context_key' ), cli=[], env={'ARG_NAME': 'env_value'}, expected=argparse.Namespace( arg_name='env_value', _context={'arg_context_key': 'env_value'}, ), ), dict( case='set _context - use default value', kwargs=dict( name='arg_name', default='default_value', env_var='ARG_NAME', context_key='arg_context_key' ), expected=argparse.Namespace( arg_name='default_value', _context={'arg_context_key': 'default_value'}, ), ), dict( case='set _context - no default value', kwargs=dict( name='arg_name', env_var='ARG_NAME', context_key='arg_context_key' ), expected=argparse.Namespace( arg_name=None, ), ), dict( case='incorrect name - should fail', kwargs=dict( name='arg name', env_var='ARG_NAME', context_key='arg_context_key' ), raises=ValueError, ), ] ) def test_add_arg( kwargs, cli, env, expected, raises, raises_match, monkeypatch, parser): if env: monkeypatch.setattr(os, 'environ', env) cli = cli or [] assert not (raises and expected), 'Cannot have both expected and raises' with raises_optionally(raises, raises_match): parser.add_arg(**kwargs) returned = parser.parse_args(cli) assert returned == expected help_text = parser.format_help() default_value = kwargs.get('default') if default_value is not None: assert f'(Default: "{default_value}"' in help_text env_var_value = kwargs.get('env_var') if env_var_value: assert f'Can also be set by env var "{env_var_value}"' in help_text @pytest.mark.parametrize( 'command_line_args, env, expected', [ ( '', {}, argparse.Namespace( check_status=None, expected_status='INGESTED', ) ), ( '', { 'CHECK_STATUS': 'true', 'EXPECTED_STATUS': 'DOWNLOADED', }, argparse.Namespace( check_status='True', expected_status='DOWNLOADED', ) ), ( '', { 'CHECK_STATUS': 'true', }, argparse.Namespace( check_status='True', expected_status='INGESTED', ) ), ( '--check-status=False', {}, argparse.Namespace( check_status='False', expected_status='INGESTED', ) ), ( '--expected-status DOWNLOADED', {}, argparse.Namespace( check_status=None, expected_status='DOWNLOADED', ) ), ] ) @patch.dict(os.environ, {}, clear=True) def test_add_check_status_args( command_line_args, env, expected, monkeypatch, parser): argv = command_line_args.split() for key, value in env.items(): monkeypatch.setenv(key, value) exec_utils.add_check_status_args( parser, default_expected_status='INGESTED', ) returned = parser.parse_args(argv) assert returned == expected @pytest.mark.parametrize( 'command_line, env, expected', [ ( '', {'CONTEXT': ''}, argparse.Namespace( context={}, ) ), ( '', {}, argparse.Namespace( context=None, ) ), ( '--context {"v":"c12"}', {}, argparse.Namespace( context={'v': 'c12'}, ) ), ( '--context={"v1":"c1"}', {}, argparse.Namespace( context={'v1': 'c1'}, ) ), ( '', {'CONTEXT': '{"source": "env"}'}, argparse.Namespace( context={'source': 'env'}, ) ), ( '--context {"source":"cli"}', {'CONTEXT': '{"source": "env"}'}, argparse.Namespace( context={'source': 'cli'}, ) ), ( '--context {invalid_json":"cli"}', {}, argparse.ArgumentError ), ] ) def test_add_context_arg( command_line, env, expected, clear_os_environ, parser, monkeypatch): argv = command_line.split() for key, value in env.items(): monkeypatch.setenv(key, value) try: exec_utils.add_raw_context_arg( parser, ) returned = parser.parse_args(argv) assert returned == expected except BaseException as e: if not (isinstance(expected, type) and isinstance(e, expected)): raise class TestFlowExecBase: EXPECTED_ARGS = { 'context', 'help', 'context_date', 'reload', 'verbose', 'wait_until_complete'} @pytest.fixture(autouse=True) def setup(self): self.flow_exec = exec_utils.FlowExecBase() def test_prepare_parser(self): parser = self.flow_exec.prepare_parser() actions = {action.dest for action in parser._actions} assert actions == self.EXPECTED_ARGS @pytest.mark.parametrize( 'case, args, expected', [ ('single date', argparse.Namespace( context_date=datetime.date(2024, 11, 1), ), [ datetime.date(2024, 11, 1), ] ), ] ) @freeze_time('2024-10-30') def test_generate_dates( self, case, args, expected): self.flow_exec.args = args try: result = self.flow_exec.generate_dates() result = list(result) assert result == expected except BaseException as e: if not (isinstance(expected, type) and isinstance(e, expected)): raise @parametrize_by_dicts( 'argv, env, expected, raises', [ dict( case='defaults', expected=dict( reload=None, wait_until_complete=None, ) ), dict( case='set by command line', argv=['--reload', 'True', '--wait-until-complete', 'True'], expected=dict( reload='True', wait_until_complete='True', _context={ 'reload': 'True', 'wait_until_complete': 'True' } ) ), dict( case='set by command line second way', argv=['--reload=True', '--wait-until-complete=True'], expected=dict( reload='True', wait_until_complete='True', _context={ 'reload': 'True', 'wait_until_complete': 'True' } ) ), dict( case='set by ENV', env={ 'RELOAD': 'False', 'WAIT_UNTIL_COMPLETE': 'True', }, expected=dict( reload='False', wait_until_complete='True', _context={ 'reload': 'False', 'wait_until_complete': 'True' } ) ), dict( case='CLI has precedence over ENV', argv=['--reload', 'True', '--wait-until-complete', 'True'], env={ 'RELOAD': 'False', 'WAIT_UNTIL_COMPLETE': 'False', }, expected=dict( reload='True', wait_until_complete='True', _context={ 'reload': 'True', 'wait_until_complete': 'True' } ) ), ] ) def test_parse_args( self, argv, env, expected, raises, monkeypatch): if env: monkeypatch.setattr(os, 'environ', env) argv = argv or [] parser = self.flow_exec.prepare_parser() args = parser.parse_args(argv) actual_args = {key: getattr(args, key, None) for key in expected.keys()} assert actual_args == expected def test_execute_context(self, mocker): date = datetime.date(2024, 1, 1) context = {'key': 'value'} cli_mock = mocker.patch.object(exec_utils, 'cli') self.flow_exec.flow_name = 'test_flow' self.flow_exec.args = argparse.Namespace( ) get_flow_instance_mock = mocker.patch.object( self.flow_exec, '_get_flow_instance' ) self.flow_exec.execute_context(date=date, context=context) assert cli_mock.execute_flow.call_args_list == [ mock.call( get_flow_instance_mock.return_value, '{"key": "value", "context_date": "2024-01-01"}', ) ] @pytest.mark.parametrize( 'case_name, params', [ ('one date, args context, empty contexts', dict( dates=[ datetime.date(2024, 1, 1), ], args=argparse.Namespace( _context={'reload': 'True'}, ), contexts=[ {}, ], expected=[ mock.call( date=datetime.date(2024, 1, 1), context={'reload': 'True'} ), ] )), ('contexts overrides args context', dict( dates=[ datetime.date(2024, 1, 1), ], args=argparse.Namespace( _context={'reload': 'True'}, ), contexts=[ {'reload': 'False'}, ], expected=[ mock.call( date=datetime.date(2024, 1, 1), context={'reload': 'False'} ), ] )), ('no args context', dict( dates=[ datetime.date(2024, 1, 1), ], args=argparse.Namespace( ), contexts=[ {'licensor': 'sme'}, ], expected=[ mock.call( date=datetime.date(2024, 1, 1), context={'licensor': 'sme'} ), ] )), ('2 dates, args context, 2 contexts', dict( dates=[ datetime.date(2024, 1, 1), datetime.date(2024, 1, 2), ], args=argparse.Namespace( _context={'reload': 'True'}, ), contexts=[ {'licensor': 'sme'}, {'licensor': 'theorchard'}, ], expected=[ mock.call( date=datetime.date(2024, 1, 1), context={'reload': 'True', 'licensor': 'sme'} ), mock.call( date=datetime.date(2024, 1, 1), context={'reload': 'True', 'licensor': 'theorchard'} ), mock.call( date=datetime.date(2024, 1, 2), context={'reload': 'True', 'licensor': 'sme'} ), mock.call( date=datetime.date(2024, 1, 2), context={'reload': 'True', 'licensor': 'theorchard'} ), ] )), ('no contexts, should raise', dict( dates=[ datetime.date(2024, 1, 1), ], args=argparse.Namespace( ), contexts=[ ], expected=ValueError )), ('no dates, should raise', dict( dates=[ ], args=argparse.Namespace( ), contexts=[ {}, ], expected=ValueError )), ] ) def test_execute(self, mocker, case_name, params): dates, args, contexts, expected = params.values() self.flow_exec.args = args mocker.patch.object( self.flow_exec, 'generate_dates', return_value=dates ) mocker.patch.object( self.flow_exec, 'generate_contexts', return_value=contexts ) execute_context_mock = mocker.patch.object( self.flow_exec, 'execute_context' ) try: self.flow_exec.execute() assert execute_context_mock.call_args_list == expected except BaseException as e: if not (isinstance(expected, type) and isinstance(e, expected)): raise finally: assert self.flow_exec.generate_dates.called assert self.flow_exec.generate_contexts.called class TestFlowCheckStatusMixin: @pytest.fixture(autouse=True) def setup(self): class TestFlowExec( exec_utils.CheckStatusMixin, exec_utils.FlowExecBase ): pass self.flow_exec = TestFlowExec() def test_get_status_for_flow_context_and_date( self, mocker): get_item_mock = mocker.patch.object( exec_utils.garcon_feed_status, '_get_item' ) flow = 'spotify_charts' date = datetime.date(2024, 11, 11) get_item_mock.return_value = { 'feed_name': 'feed_name', 'date': '2024-11-11', 'status': 'DOWNLOADED', } context = { 'chart': 'chart_name', } status = self.flow_exec.get_status_for_flow_context_and_date( flow=flow, context=context, date=date ) assert status == 'DOWNLOADED' assert get_item_mock.call_args_list == [ mock.call('spotify_charts_chart_name', '2024-11-11'), ] @pytest.mark.parametrize( 'is_completed, expected_to_call', [ (True, False), (False, True), ] ) def test_execute_context(self, mocker, is_completed, expected_to_call): execute_parent_mock = mocker.patch.object( exec_utils.FlowExecBase, 'execute_context' ) mocker.patch.object( self.flow_exec, 'is_context_completed', return_value=is_completed ) self.flow_exec.execute_context(context=MagicMock(), date=MagicMock()) result = execute_parent_mock.called assert result == expected_to_call @pytest.mark.parametrize( 'case_name, args, current_status, expected', [ ( 'no check status', argparse.Namespace( reload=None, check_status=False, expected_status='DOWNLOADED', ), 'DOWNLOADED', False, ), ( 'check_status', argparse.Namespace( reload=None, check_status=True, expected_status='DOWNLOADED', ), 'DOWNLOADED', True, ), ( 'check_status', argparse.Namespace( reload='False', check_status=True, expected_status='DOWNLOADED', ), 'DOWNLOADED', True, ), ( 'reload true should not check_status', argparse.Namespace( reload='True', check_status=True, expected_status='DOWNLOADED', ), 'DOWNLOADED', False, ), ] ) def test_is_context_completed( self, case_name, args, current_status, expected, mocker): self.flow_exec.flow_name = 'test_flow' mocker.patch.object( self.flow_exec, 'get_status_for_flow_context_and_date', lambda *v, **kw: current_status ) self.flow_exec.args = args result = self.flow_exec.is_context_completed( context=MagicMock(), date=MagicMock(), ) assert result == expected class TestFlowExecDaily: EXPECTED_ARGS = TestFlowExecBase.EXPECTED_ARGS | { 'days', 'skip', 'expected_status', 'check_status', 'max_concurrent_flows', 'concurrency_strategy', } @pytest.fixture(autouse=True) def setup(self): self.flow_exec = exec_utils.FlowExecDaily() def test_concurrency_control_settings(self): """Test that the ConcurrencyControlMixin settings are correct.""" mro = type(self.flow_exec).mro() assert mro.index(exec_utils.CheckStatusMixin) < mro.index( exec_utils.ConcurrencyControlMixin ) assert self.flow_exec.TIMEOUT_SECONDS == 60 * 60 parser = self.flow_exec.prepare_parser() actions = {action.dest: action for action in parser._actions} assert actions['concurrency_strategy'].default == 'WAIT' assert actions['max_concurrent_flows'].default == 5 def test_prepare_parser(self): parser = self.flow_exec.prepare_parser() actions = {action.dest for action in parser._actions} assert actions == self.EXPECTED_ARGS @parametrize_by_dicts( 'args, expected, raises, raises_match', [ dict( case='single date and period', args=argparse.Namespace( context_date=datetime.date(2024, 11, 2), days=2, skip=1, ), raises=ValueError, raises_match='Cannot use both days and context_date arguments', ), dict( case='2 days 1 skip', args=argparse.Namespace( context_date=None, days=2, skip=1, ), expected=[ datetime.date(2024, 10, 29), datetime.date(2024, 10, 28), ] ), dict( case='neither date or days presented', args=argparse.Namespace( context_date=None, days=None, skip=1, ), expected=[] ), dict( case='negative days - should raise', args=argparse.Namespace( context_date=None, days=-1, skip=1, ), raises=ValueError ), ] ) @freeze_time('2024-10-30') def test_generate_dates( self, args, expected, raises, raises_match): self.flow_exec.args = args with raises_optionally(raises, raises_match): result = self.flow_exec.generate_dates() result = list(result) assert result == expected @parametrize_by_dicts( 'kwargs, argv, expected, raises, raises_match', [ dict( kwargs=dict(), argv=['--skip', '2'], expected=dict( skip=2, days=None, context_date=None, ) ), dict( kwargs=dict(), argv=['--skip', '2', '--days', '0'], expected=dict( skip=2, days=0, context_date=None, ) ), dict( kwargs=dict(), argv=['--skip', '2', '--days', ''], expected=dict( skip=2, days=None, context_date=None, ) ), ] ) def test_days_and_skip( self, kwargs, argv, expected, raises, raises_match, parser): self.flow_exec.prepare_parser() self.flow_exec.parse_args(argv) expected_namespace = argparse.Namespace( skip=None, days=None, check_status=None, expected_status='INGESTED', max_concurrent_flows=5, concurrency_strategy='WAIT', verbose=False, context_date=None, reload=None, context={}, wait_until_complete=None ) for k, v in expected.items(): setattr(expected_namespace, k, v) assert self.flow_exec.args == expected_namespace class TestFlowExecMonthly: EXPECTED_ARGS = TestFlowExecBase.EXPECTED_ARGS | { 'months', 'skip', 'expected_status', 'check_status', 'max_concurrent_flows', 'concurrency_strategy', } @pytest.fixture(autouse=True) def setup(self): self.flow_exec = exec_utils.FlowExecMonthly() def test_concurrency_control_settings(self): # check status goes before concurrency check mro = type(self.flow_exec).mro() assert mro.index(exec_utils.CheckStatusMixin) < mro.index( exec_utils.ConcurrencyControlMixin ) assert self.flow_exec.TIMEOUT_SECONDS == 60 * 60 parser = self.flow_exec.prepare_parser() actions = {action.dest: action for action in parser._actions} assert actions['concurrency_strategy'].default == 'WAIT' assert actions['max_concurrent_flows'].default == 5 def test_prepare_parser(self): parser = self.flow_exec.prepare_parser() actions = {action.dest for action in parser._actions} assert actions == self.EXPECTED_ARGS @parametrize_by_dicts( 'args, expected, raises', [ dict( case='single date and period', args=argparse.Namespace( context_date=datetime.date(2024, 11, 2), months=2, skip=1, ), raises=ValueError ), dict( case='2 months 1 skip', args=argparse.Namespace( context_date=None, months=2, skip=1, ), expected=[ datetime.date(2024, 9, 1), datetime.date(2024, 8, 1), ] ), dict( case='neither date or months presented - should raise', args=argparse.Namespace( context_date=None, months=None, skip=1, ), expected=[] ), dict( case='negative months - should raise', args=argparse.Namespace( context_date=None, months=-1, skip=1, ), raises=ValueError ), ] ) @freeze_time('2024-10-30') def test_generate_dates( self, args, expected, raises): self.flow_exec.args = args with raises_optionally(raises): result = self.flow_exec.generate_dates() result = list(result) assert result == expected class TestFlowExecWeekly: EXPECTED_ARGS = TestFlowExecBase.EXPECTED_ARGS | { 'weeks', 'skip', 'expected_status', 'check_status', 'max_concurrent_flows', 'concurrency_strategy', } @pytest.fixture(autouse=True) def setup(self, mocker): self.flow_exec = exec_utils.FlowExecWeekly() self.flow_exec.flow_name = 'test_flow' def test_prepare_parser(self): parser = self.flow_exec.prepare_parser() actions = {action.dest: action for action in parser._actions} expected_actions = self.EXPECTED_ARGS actions_set = set(actions) assert expected_actions.issubset(actions_set), ( expected_actions.difference(actions_set)) @parametrize_by_dicts( 'args, expected, raises, raises_match', [ dict( case='single date and period', args=argparse.Namespace( context_date=datetime.date(2024, 11, 2), weeks=2, skip=1, weekday=dates_util.Weekday.THURSDAY, ), raises=ValueError ), dict( case='weeks and skip respect weekday - adjust to Tuesday', args=argparse.Namespace( context_date=None, weeks=2, skip=1, weekday=dates_util.Weekday.TUESDAY, ), expected=[ datetime.date(2024, 11, 12), datetime.date(2024, 11, 5), ] ), dict( case='weeks and skip without weekday - no adjust weekday', args=argparse.Namespace( context_date=None, weeks=2, skip=1, weekday=None, ), expected=[ datetime.date(2024, 11, 13), datetime.date(2024, 11, 6), ] ), dict( case='neither date or weeks presented - should raise', args=argparse.Namespace( context_date=None, weeks=None, skip=1, weekday=None, ), raises=ValueError, marks=pytest.mark.xfail( reason='Will fail when date moved from ExecFlowBase' ), ), dict( case='negative weeks - should raise', args=argparse.Namespace( context_date=None, weeks=-1, skip=1, weekday=None, ), raises=ValueError, raises_match='weeks should be positive', ), ] ) @freeze_time('2024-11-20') def test_generate_dates( self, args, expected, raises, raises_match): self.flow_exec.args = args with raises_optionally(raises, raises_match): result = self.flow_exec.generate_dates() result = list(result) assert result == expected class TestConcurrencyControlMixin: VALID_STRATEGIES = [ e.value for e in exec_utils.ConcurrencyControlMixin.Strategy ] @pytest.fixture(autouse=True) def setup(self): class TestFlowExec( exec_utils.ConcurrencyControlMixin, exec_utils.FlowExecBase ): pass self.flow_exec = TestFlowExec() self.flow_exec.flow_name = 'spotify' def test_arguments(self): parser = self.flow_exec.prepare_parser() actions = {action.dest: action for action in parser._actions} assert 'max_concurrent_flows' in actions assert actions['max_concurrent_flows'].default == 5 assert 'concurrency_strategy' in actions assert actions['concurrency_strategy'].choices == [ 'SKIP', 'IGNORE', 'FAIL', 'WAIT'] @pytest.mark.parametrize( 'strategy, expected_exception', [ ('IGNORE', None), ('SKIP', None), ('FAIL', RuntimeError), ('WAIT', None), (None, None), ('Not a valid strategy', ValueError), ] ) @pytest.mark.parametrize( 'num_flows', [ 4, 5, 6, ] ) @pytest.mark.parametrize( 'max_concurrent_flows', [ 5, 0, None, ] ) def test_execute_context( self, max_concurrent_flows, num_flows, strategy, expected_exception, mocker): mocker.patch.object( exec_utils.swf_util, 'count_running_workflows_by_type', return_value=num_flows ) parent_execute_context_mock = mocker.patch.object( FlowExecBase, 'execute_context' ) args = argparse.Namespace( concurrency_strategy=strategy, max_concurrent_flows=max_concurrent_flows, ) self.flow_exec.args = args expected_parent_execute_called = True if not max_concurrent_flows or not strategy: # not configured, should always work expected_exception = None expected_parent_execute_called = True elif num_flows < max_concurrent_flows: if strategy in self.VALID_STRATEGIES: # when within limit all strategies should work expected_exception = None expected_parent_execute_called = True else: if expected_exception: expected_parent_execute_called = False else: if strategy == 'WAIT': pytest.skip('WAIT strategy case tested separately') if expected_exception: expected_parent_execute_called = False if strategy in ['FAIL', 'SKIP']: expected_parent_execute_called = False date_mock = MagicMock() context_mock = MagicMock() try: self.flow_exec.execute_context( context=context_mock, date=date_mock, ) assert not expected_exception assert (parent_execute_context_mock.called == expected_parent_execute_called) except BaseException as e: if not (isinstance(expected_exception, type) and isinstance(e, expected_exception)): raise assert (parent_execute_context_mock.called == expected_parent_execute_called) @parametrize_by_dicts( 'num_flows, raises, expected_sleep_calls', [ dict( case='within limit, no sleep', num_flows=[4, 6], expected_sleep_calls=0, ), dict( case='out of limit but within timeout', num_flows=[5, 6, 4, 4], expected_sleep_calls=2, ), dict( case='out of limit and timeout', num_flows=[5, 6, 5, 4], raises=TimeoutError, expected_sleep_calls=3, ), ] ) def test_execute_context_wait_strategy( self, num_flows, raises, expected_sleep_calls, mocker, frozen_time, ): self.flow_exec.TIMEOUT_SECONDS = 65 self.flow_exec.SLEEP_SECONDS = 30 freeze_time_mock, sleep_mock = frozen_time mocker.patch.object( exec_utils.swf_util, 'count_running_workflows_by_type', side_effect=num_flows, ) parent_execute_context_mock = mocker.patch.object( FlowExecBase, 'execute_context' ) args = argparse.Namespace( concurrency_strategy='WAIT', max_concurrent_flows=5, ) self.flow_exec.args = args date_mock = MagicMock() context_mock = MagicMock() with raises_optionally(raises): self.flow_exec.execute_context( context=context_mock, date=date_mock, ) assert sleep_mock.call_count == expected_sleep_calls assert parent_execute_context_mock.called == (not raises)