import pytest import zmq from dapd_transformation_service.entities.worker_control import ( ConsumerStateEnum, WorkerControlEnum, WorkerHealthResponseEnum, ) from dapd_transformation_service.exceptions import ConsumerNotAvailableError from dapd_transformation_service.helpers.zmq_client import ZmqClient def test_request_worker_health_ok(mocker, config, logger): zmq_request = {'header': WorkerControlEnum.GET_HEALTH, 'payload': ''} zmq_response = {'header': WorkerControlEnum.GET_HEALTH, 'payload': WorkerHealthResponseEnum.OK} mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(return_value=zmq_response) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.request_worker_health() assert res == WorkerHealthResponseEnum.OK assert mocked_socket.send_json.call_count == 1 assert mocked_socket.send_json.call_args == mocker.call(zmq_request, ) def test_request_worker_health_timeout(mocker, config, logger): mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(side_effect=zmq.error.Again) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.request_worker_health() assert res == WorkerHealthResponseEnum.TIMEOUT def test_request_worker_health_general_exception(mocker, config, logger): mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(side_effect=Exception) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.request_worker_health() assert res == WorkerHealthResponseEnum.INVALID_RESPONSE def test_get_consumer_state_ok(mocker, config, logger): zmq_request = {'header': WorkerControlEnum.GET_CONSUMER_STATE, 'payload': ''} zmq_response = { 'header': WorkerControlEnum.GET_CONSUMER_STATE, 'payload': ConsumerStateEnum.ENABLED } mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(return_value=zmq_response) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.get_consumer_state() assert res == ConsumerStateEnum.ENABLED assert mocked_socket.send_json.call_count == 1 assert mocked_socket.send_json.call_args == mocker.call(zmq_request, ) def test_get_consumer_state_general_exception(mocker, config, logger): mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(side_effect=Exception) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.get_consumer_state() assert res == ConsumerStateEnum.UNKNOWN def test_get_consumer_metrics_ok(mocker, config, logger): zmq_request = {'header': WorkerControlEnum.GET_METRICS, 'payload': ''} zmq_response = {'header': WorkerControlEnum.GET_METRICS, 'payload': 'metrics_mocked_payload'} mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(return_value=zmq_response) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) client = ZmqClient(logger, config) res = client.get_consumer_metrics() assert res == zmq_response['payload'] assert mocked_socket.send_json.call_count == 1 assert mocked_socket.send_json.call_args == mocker.call(zmq_request, ) def test_get_consumer_metrics_zmq_exception(mocker, config, logger): mocked_socket = mocker.MagicMock() mocked_socket.recv_json = mocker.MagicMock(side_effect=zmq.error.Again) mocker.patch( 'dapd_transformation_service.helpers.zmq_client.ZmqClient.initialize_socket', return_value=mocked_socket ) with pytest.raises(ConsumerNotAvailableError): client = ZmqClient(logger, config) client.get_consumer_metrics()