import json from datetime import datetime, timedelta, timezone import pytest from dapd_transformation_service.entities.worker_control import ConsumerMetrics, ConsumerStateEnum from dapd_transformation_service.services.consumer_state import ConsumerStateService pytestmark = pytest.mark.integration FREEZED_DT = datetime(2015, 10, 15, 12, 23, 45, tzinfo=timezone.utc) def test_set_consumer_state(config, mocker, logger): client_method_mocked = mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.set_consumer_state' ) state = ConsumerStateEnum.ENABLED service = ConsumerStateService(logger, config) service.set_consumer_state(state) client_method_mocked.assert_called_once_with(state) def test_get_consumer_state(config, mocker, logger): expected = ConsumerStateEnum.ENABLED client_method_mocked = mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.get_consumer_state', return_value=expected ) service = ConsumerStateService(logger, config) result = service.get_consumer_state() assert result == expected client_method_mocked.assert_called_once_with() @pytest.mark.freeze_time(FREEZED_DT) def test_get_consumer_metrics(config, mocker, logger): mocker.patch( 'dapd_transformation_service.services.consumer_state.utcnow', return_value=FREEZED_DT ) started_at = FREEZED_DT - timedelta(seconds=4) consumer_metrics = ConsumerMetrics(started_at=started_at, consumed_records=10) client_method_mocked = mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.get_consumer_metrics', return_value=consumer_metrics.to_json() ) expected = { 'started_at': started_at.isoformat(), 'uptime': '0:00:04', 'consumed_records': 10, } service = ConsumerStateService(logger, config) result = service.get_consumer_metrics() assert result == expected client_method_mocked.assert_called_once_with() def test_get_consumer_config(config, mocker, logger): client_method_mocked = mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.get_consumer_config', return_value=json.dumps(config.to_dict()) ) expected = config.to_dict() service = ConsumerStateService(logger, config) result = service.get_consumer_config() assert result == expected client_method_mocked.assert_called_once_with() def test_set_thread_sync_event__set(config, mocker, logger): event_mocked = mocker.patch('threading.Event', autospec=True) service = ConsumerStateService(logger, config, thread_sync_event=event_mocked) service.set_thread_sync_event(is_enabled=True) assert event_mocked.set.called assert not event_mocked.clear.called def test_set_thread_sync_event__clear(config, mocker, logger): event_mocked = mocker.patch('threading.Event', autospec=True) service = ConsumerStateService(logger, config, thread_sync_event=event_mocked) service.set_thread_sync_event(is_enabled=False) assert not event_mocked.set.called assert event_mocked.clear.called def test_set_thread_sync_event__none(config, mocker, logger): event_mocked = mocker.patch('threading.Event', autospec=True) service = ConsumerStateService(logger, config, thread_sync_event=None) service.set_thread_sync_event(is_enabled=False) assert not event_mocked.set.called assert not event_mocked.clear.called