From 75952c1ff9da51d94339c6e9cbf599ea45fd3a4d Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 1 Jun 2026 17:32:57 +0100 Subject: [PATCH 01/20] Add udc auth to helm chart --- helm/daq-queuing-service/templates/deployment.yaml | 8 ++++++++ helm/daq-queuing-service/values.yaml | 6 +++++- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/helm/daq-queuing-service/templates/deployment.yaml b/helm/daq-queuing-service/templates/deployment.yaml index f96d836..5f11ea2 100644 --- a/helm/daq-queuing-service/templates/deployment.yaml +++ b/helm/daq-queuing-service/templates/deployment.yaml @@ -60,6 +60,14 @@ spec: {{- with .Values.volumeMounts }} {{- toYaml . | nindent 12 }} {{- end }} + {{- if .Values.udcSecret.enabled }} + env: + - name: UDC_SECRET + valueFrom: + secretKeyRef: + name: {{ .Values.udcSecret.name }} + key: {{ .Values.udcSecret.key }} + {{- end }} volumes: {{- with .Values.volumes }} {{- toYaml . | nindent 8 }} diff --git a/helm/daq-queuing-service/values.yaml b/helm/daq-queuing-service/values.yaml index 37a097c..5071db5 100644 --- a/helm/daq-queuing-service/values.yaml +++ b/helm/daq-queuing-service/values.yaml @@ -9,7 +9,6 @@ imagePullSecrets: [] nameOverride: "" fullnameOverride: "" - podAnnotations: {} podLabels: {} @@ -74,3 +73,8 @@ volumeMounts: [] volumes: [] nodeSelector: {} tolerations: [] + +udcSecret: + enabled: false + name: "" + key: udc-secret From edc221b44346c3755f0f6eac74bd950211626cf9 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 1 Jun 2026 18:05:40 +0100 Subject: [PATCH 02/20] Create UDC session manager --- src/daq_queuing_service/app/app.py | 9 +++++- .../blueapi_interaction/session_manager.py | 28 +++++++++++++++++++ 2 files changed, 36 insertions(+), 1 deletion(-) create mode 100644 src/daq_queuing_service/blueapi_interaction/session_manager.py diff --git a/src/daq_queuing_service/app/app.py b/src/daq_queuing_service/app/app.py index acea0c8..9355815 100644 --- a/src/daq_queuing_service/app/app.py +++ b/src/daq_queuing_service/app/app.py @@ -5,12 +5,14 @@ from blueapi.client import BlueapiClient from blueapi.client.rest import BlueapiRestClient +from blueapi.service.authentication import SessionCacheManager from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from daq_queuing_service.api.api import create_api_router from daq_queuing_service.api.errors import register_exception_handlers from daq_queuing_service.blueapi_interaction.blueapi_adapter import BlueapiClientAdapter +from daq_queuing_service.blueapi_interaction.session_manager import UDCSessionManager from daq_queuing_service.broadcaster import Broadcaster from daq_queuing_service.plugins.construct_task_request import ( construct_blueapi_call_list, @@ -63,7 +65,12 @@ def log_task_exception(task: asyncio.Task[NoReturn]): app.state.queue = TaskQueue(construct_blueapi_call_list, broadcaster) - blueapi_rest_client = BlueapiRestClient(config=config.blueapi.api) + assert config.blueapi.oidc + session_manager = UDCSessionManager(config.blueapi.oidc, SessionCacheManager(None)) + + blueapi_rest_client = BlueapiRestClient( + config=config.blueapi.api, session_manager=session_manager + ) blueapi_client = BlueapiClient.from_config(config.blueapi) blueapi_client_adapter = BlueapiClientAdapter(blueapi_client) diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py new file mode 100644 index 0000000..ebd5e59 --- /dev/null +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -0,0 +1,28 @@ +import os + +import requests +from blueapi.client.rest import SessionManager + + +class UDCSessionManager(SessionManager): + def get_valid_access_token(self) -> str: + token_url = ( + "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token" + ) + + client_id = "i15-1-udc" + client_secret = os.environ["UDC_SECRET"] + + if not (client_id and client_secret): + return "" + + response = requests.post( + token_url, + data={ + "client_id": client_id, + "client_secret": client_secret, + "grant_type": "client_credentials", + }, + ) + response.raise_for_status() + return response.json().get("access_token") From 983869f3bfeea7d317a9b83d8303307800fc1412 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Tue, 2 Jun 2026 14:00:09 +0100 Subject: [PATCH 03/20] Allow no oidc config and improve logging --- src/daq_queuing_service/api/api.py | 4 +--- src/daq_queuing_service/app/app.py | 4 +++- .../blueapi_interaction/blueapi_adapter.py | 3 +-- .../blueapi_interaction/session_manager.py | 11 ++++++++-- src/daq_queuing_service/broadcaster.py | 3 +-- src/daq_queuing_service/log.py | 21 +++++++++++++++++++ src/daq_queuing_service/task_queue/queue.py | 4 +--- src/daq_queuing_service/worker/worker.py | 5 ++++- tests/unit_tests/test_worker.py | 2 +- 9 files changed, 42 insertions(+), 15 deletions(-) create mode 100644 src/daq_queuing_service/log.py diff --git a/src/daq_queuing_service/api/api.py b/src/daq_queuing_service/api/api.py index d88fd19..0315e40 100644 --- a/src/daq_queuing_service/api/api.py +++ b/src/daq_queuing_service/api/api.py @@ -1,6 +1,5 @@ import asyncio import json -import logging from collections.abc import AsyncGenerator, Callable from blueapi.client.rest import ( @@ -16,6 +15,7 @@ from daq_queuing_service.app._config import AppConfig, load_config from daq_queuing_service.blueapi_interaction.blueapi_call import BlueapiCallResponse from daq_queuing_service.broadcaster import Broadcaster +from daq_queuing_service.log import LOGGER from daq_queuing_service.task import ExperimentDefinition, Status, Task from daq_queuing_service.task_queue.queue import ( QUEUE_EVENTS, @@ -26,8 +26,6 @@ # pyright: reportUnusedFunction=false -LOGGER = logging.getLogger(__name__) - class InvalidExperimentDefinitionsError(Exception): def __init__(self, errors: dict[int, InvalidParametersError | UnknownPlanError]): diff --git a/src/daq_queuing_service/app/app.py b/src/daq_queuing_service/app/app.py index 9355815..7689c65 100644 --- a/src/daq_queuing_service/app/app.py +++ b/src/daq_queuing_service/app/app.py @@ -2,6 +2,7 @@ import logging from contextlib import asynccontextmanager from typing import NoReturn +from unittest.mock import MagicMock from blueapi.client import BlueapiClient from blueapi.client.rest import BlueapiRestClient @@ -65,7 +66,8 @@ def log_task_exception(task: asyncio.Task[NoReturn]): app.state.queue = TaskQueue(construct_blueapi_call_list, broadcaster) - assert config.blueapi.oidc + if not config.blueapi.oidc: + config.blueapi.oidc = MagicMock() session_manager = UDCSessionManager(config.blueapi.oidc, SessionCacheManager(None)) blueapi_rest_client = BlueapiRestClient( diff --git a/src/daq_queuing_service/blueapi_interaction/blueapi_adapter.py b/src/daq_queuing_service/blueapi_interaction/blueapi_adapter.py index 87c4851..f312895 100644 --- a/src/daq_queuing_service/blueapi_interaction/blueapi_adapter.py +++ b/src/daq_queuing_service/blueapi_interaction/blueapi_adapter.py @@ -1,5 +1,4 @@ import asyncio -import logging from dataclasses import dataclass from typing import Generic, TypeVar @@ -14,7 +13,7 @@ from blueapi.service.model import TaskRequest from blueapi.worker import TaskStatus, WorkerState -LOGGER = logging.getLogger(__name__) +from daq_queuing_service.log import LOGGER T = TypeVar("T") E = TypeVar("E", bound=Exception) diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index ebd5e59..fbe7b1d 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -3,6 +3,8 @@ import requests from blueapi.client.rest import SessionManager +from daq_queuing_service.log import LOGGER + class UDCSessionManager(SessionManager): def get_valid_access_token(self) -> str: @@ -11,11 +13,14 @@ def get_valid_access_token(self) -> str: ) client_id = "i15-1-udc" - client_secret = os.environ["UDC_SECRET"] + client_secret = os.environ.get("UDC_SECRET") if not (client_id and client_secret): + LOGGER.debug("No UDC secret found") return "" + LOGGER.debug("Found UDC secret") + response = requests.post( token_url, data={ @@ -25,4 +30,6 @@ def get_valid_access_token(self) -> str: }, ) response.raise_for_status() - return response.json().get("access_token") + token = response.json().get("access_token") + LOGGER.debug(f"Got token: {token}") + return token diff --git a/src/daq_queuing_service/broadcaster.py b/src/daq_queuing_service/broadcaster.py index 23fca50..0faf5df 100644 --- a/src/daq_queuing_service/broadcaster.py +++ b/src/daq_queuing_service/broadcaster.py @@ -1,11 +1,10 @@ import asyncio -import logging from collections.abc import Iterable, Mapping from typing import Any, Generic, TypedDict, TypeVar from pydantic import BaseModel -LOGGER = logging.getLogger(__name__) +from daq_queuing_service.log import LOGGER T = TypeVar("T", bound=str) diff --git a/src/daq_queuing_service/log.py b/src/daq_queuing_service/log.py new file mode 100644 index 0000000..44629f7 --- /dev/null +++ b/src/daq_queuing_service/log.py @@ -0,0 +1,21 @@ +import logging + +import colorlog + +HANDLER = colorlog.StreamHandler() +HANDLER.setFormatter( + colorlog.ColoredFormatter( + "%(log_color)s%(asctime)s [%(name)s] %(levelname)s: %(message)s", + log_colors={ + "DEBUG": "cyan", + "INFO": "green", + "WARNING": "yellow", + "ERROR": "red", + "CRITICAL": "bold_red", + }, + ) +) +LOGGER = logging.getLogger("Queue") +LOGGER.addHandler(HANDLER) +LOGGER.setLevel(logging.DEBUG) +LOGGER.propagate = False diff --git a/src/daq_queuing_service/task_queue/queue.py b/src/daq_queuing_service/task_queue/queue.py index 05eb451..8a1b8cd 100644 --- a/src/daq_queuing_service/task_queue/queue.py +++ b/src/daq_queuing_service/task_queue/queue.py @@ -1,5 +1,4 @@ import asyncio -import logging from collections.abc import Callable, Sequence from types import TracebackType from typing import Any, Literal @@ -13,6 +12,7 @@ CallStatus, ) from daq_queuing_service.broadcaster import Broadcaster, Event +from daq_queuing_service.log import LOGGER from daq_queuing_service.task import Status, Task, TaskWithPosition from daq_queuing_service.task_queue.queue_utils import ( NegativePositionError, @@ -23,8 +23,6 @@ TaskNotInQueueError, ) -LOGGER = logging.getLogger(__name__) - Converter = Callable[ [list[TaskWithPosition], list[TaskWithPosition], list[BlueapiCall]], list[BlueapiCall], diff --git a/src/daq_queuing_service/worker/worker.py b/src/daq_queuing_service/worker/worker.py index 187b7e7..2908428 100644 --- a/src/daq_queuing_service/worker/worker.py +++ b/src/daq_queuing_service/worker/worker.py @@ -17,11 +17,14 @@ from daq_queuing_service.blueapi_interaction.blueapi_adapter import BlueapiClientAdapter from daq_queuing_service.blueapi_interaction.blueapi_call import BlueapiCall, CallStatus +from daq_queuing_service.log import HANDLER from daq_queuing_service.task import ExperimentDefinition from daq_queuing_service.task_queue.queue import TaskQueue -LOGGER = logging.getLogger(__name__) +LOGGER = logging.getLogger("Queue Worker") +LOGGER.addHandler(HANDLER) LOGGER.setLevel(logging.DEBUG) +LOGGER.propagate = False class QueueWorker: diff --git a/tests/unit_tests/test_worker.py b/tests/unit_tests/test_worker.py index ec43664..9e3d9b7 100644 --- a/tests/unit_tests/test_worker.py +++ b/tests/unit_tests/test_worker.py @@ -228,7 +228,7 @@ async def test_when_parameter_error_then_call_failed_and_error_added_to_call( worker_with_parameter_error._client.run_task.assert_called_once() # type: ignore -async def test_when_plan_name_error_then_call_failed_and_error_added_to_task( +async def test_when_plan_name_error_then_call_failed_and_error_added_to_call( worker_with_unknown_plan_error: QueueWorker, only_loop_once: type[Exception] ): From ac224acdcc16b01a3d86f6733feff9c96de77976 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 10:41:53 +0100 Subject: [PATCH 04/20] Fix tests --- src/daq_queuing_service/log.py | 1 - src/daq_queuing_service/worker/worker.py | 1 - 2 files changed, 2 deletions(-) diff --git a/src/daq_queuing_service/log.py b/src/daq_queuing_service/log.py index 44629f7..3bc74e2 100644 --- a/src/daq_queuing_service/log.py +++ b/src/daq_queuing_service/log.py @@ -18,4 +18,3 @@ LOGGER = logging.getLogger("Queue") LOGGER.addHandler(HANDLER) LOGGER.setLevel(logging.DEBUG) -LOGGER.propagate = False diff --git a/src/daq_queuing_service/worker/worker.py b/src/daq_queuing_service/worker/worker.py index 2908428..0b60906 100644 --- a/src/daq_queuing_service/worker/worker.py +++ b/src/daq_queuing_service/worker/worker.py @@ -24,7 +24,6 @@ LOGGER = logging.getLogger("Queue Worker") LOGGER.addHandler(HANDLER) LOGGER.setLevel(logging.DEBUG) -LOGGER.propagate = False class QueueWorker: From 92337177b1de9dabbde0e09a402c5e78e26663fc Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 13:32:45 +0100 Subject: [PATCH 05/20] Instantiate blueapi client with rest client --- src/daq_queuing_service/app/app.py | 15 +------ .../blueapi_interaction/clients.py | 39 +++++++++++++++++++ .../blueapi_interaction/session_manager.py | 2 +- 3 files changed, 42 insertions(+), 14 deletions(-) create mode 100644 src/daq_queuing_service/blueapi_interaction/clients.py diff --git a/src/daq_queuing_service/app/app.py b/src/daq_queuing_service/app/app.py index 7689c65..292accb 100644 --- a/src/daq_queuing_service/app/app.py +++ b/src/daq_queuing_service/app/app.py @@ -2,18 +2,14 @@ import logging from contextlib import asynccontextmanager from typing import NoReturn -from unittest.mock import MagicMock -from blueapi.client import BlueapiClient -from blueapi.client.rest import BlueapiRestClient -from blueapi.service.authentication import SessionCacheManager from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from daq_queuing_service.api.api import create_api_router from daq_queuing_service.api.errors import register_exception_handlers from daq_queuing_service.blueapi_interaction.blueapi_adapter import BlueapiClientAdapter -from daq_queuing_service.blueapi_interaction.session_manager import UDCSessionManager +from daq_queuing_service.blueapi_interaction.clients import get_blueapi_clients from daq_queuing_service.broadcaster import Broadcaster from daq_queuing_service.plugins.construct_task_request import ( construct_blueapi_call_list, @@ -66,14 +62,7 @@ def log_task_exception(task: asyncio.Task[NoReturn]): app.state.queue = TaskQueue(construct_blueapi_call_list, broadcaster) - if not config.blueapi.oidc: - config.blueapi.oidc = MagicMock() - session_manager = UDCSessionManager(config.blueapi.oidc, SessionCacheManager(None)) - - blueapi_rest_client = BlueapiRestClient( - config=config.blueapi.api, session_manager=session_manager - ) - blueapi_client = BlueapiClient.from_config(config.blueapi) + blueapi_rest_client, blueapi_client = get_blueapi_clients(config.blueapi) blueapi_client_adapter = BlueapiClientAdapter(blueapi_client) app.state.worker = QueueWorker( diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py new file mode 100644 index 0000000..b322011 --- /dev/null +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -0,0 +1,39 @@ +from unittest.mock import MagicMock + +from blueapi.client import BlueapiClient +from blueapi.client.event_bus import EventBusClient +from blueapi.client.rest import BlueapiRestClient +from blueapi.config import ApplicationConfig +from blueapi.service.authentication import SessionCacheManager +from bluesky_stomp.messaging import Broker, StompClient + +from daq_queuing_service.blueapi_interaction.session_manager import UDCSessionManager + + +def get_blueapi_clients( + blueapi_config: ApplicationConfig, +) -> tuple[BlueapiRestClient, BlueapiClient]: + if not blueapi_config.oidc: + blueapi_config.oidc = MagicMock() + + session_manager = UDCSessionManager(blueapi_config.oidc, SessionCacheManager(None)) + blueapi_rest_client = BlueapiRestClient( + config=blueapi_config.api, session_manager=session_manager + ) + + if blueapi_config.stomp.enabled: + assert blueapi_config.stomp.url.host is not None, "Stomp URL missing host" + assert blueapi_config.stomp.url.port is not None, "Stomp URL missing port" + stomp_client = StompClient.for_broker( + broker=Broker( + host=blueapi_config.stomp.url.host, + port=blueapi_config.stomp.url.port, + auth=blueapi_config.stomp.auth, + ) + ) + events = EventBusClient(stomp_client) + blueapi_client = BlueapiClient(blueapi_rest_client, events) + else: + blueapi_client = BlueapiClient(blueapi_rest_client) + + return blueapi_rest_client, blueapi_client diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index fbe7b1d..e74e877 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -31,5 +31,5 @@ def get_valid_access_token(self) -> str: ) response.raise_for_status() token = response.json().get("access_token") - LOGGER.debug(f"Got token: {token}") + LOGGER.debug("Returning token") return token From 2fe630ae68059ff650f089edde9cde8b41465eeb Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 15:06:57 +0100 Subject: [PATCH 06/20] Upgrade blueapi --- pyproject.toml | 2 +- uv.lock | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 233f095..57ace2c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -13,7 +13,7 @@ classifiers = [ ] description = "A service to queue tasks and chain BlueAPI calls" dependencies = [ - "blueapi>=1.13.0", + "blueapi>=1.14.0", "fastapi>=0.136.0", "pydantic>=2.13.2", ] diff --git a/uv.lock b/uv.lock index dec3a65..2dff840 100644 --- a/uv.lock +++ b/uv.lock @@ -381,7 +381,7 @@ wheels = [ [[package]] name = "blueapi" -version = "1.13.0" +version = "1.14.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "aioca" }, @@ -409,9 +409,9 @@ dependencies = [ { name = "tomlkit" }, { name = "uvicorn" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/89/19/9ee222e1efef179c6975d070ef99b04752cecd70eac9479b670c6cb5ae3a/blueapi-1.13.0.tar.gz", hash = "sha256:3bf77f75e6bef65326786cc01be74fcf1eded37d463d9f24d8ad446cbbaccaa7", size = 1838838, upload-time = "2026-04-17T12:31:25.597Z" } +sdist = { url = "https://files.pythonhosted.org/packages/18/16/bb4e52cbbc59ba94d2ded1778206e6ce40664fca178f66aa3fe5b8bd7d5c/blueapi-1.14.0.tar.gz", hash = "sha256:d3cb6975fa8826fd9ec9adf47e75719427e9b992002956c0433c4abb887d86d4", size = 1840934, upload-time = "2026-04-24T10:10:35.367Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/d0/a8/bb3205aca03ddc4305130b92fda5dde7c69ef91b403f1a8aada324c50256/blueapi-1.13.0-py3-none-any.whl", hash = "sha256:99a17cd87d59a7dfebeac94b68fea701d0d9d1fc0f799cfd3b710a6259610a2f", size = 83960, upload-time = "2026-04-17T12:31:24.239Z" }, + { url = "https://files.pythonhosted.org/packages/f4/0a/9499d4ee173cab1f891c86f8332e8c2d4844ed5d76a72f84ced0f134c604/blueapi-1.14.0-py3-none-any.whl", hash = "sha256:e232529e795fad0f9dec36d1bd6dd9f5bcac21f15967a6a609b01b9306a80ee2", size = 84560, upload-time = "2026-04-24T10:10:33.59Z" }, ] [[package]] @@ -1030,7 +1030,7 @@ dev = [ [package.metadata] requires-dist = [ - { name = "blueapi", specifier = ">=1.13.0" }, + { name = "blueapi", specifier = ">=1.14.0" }, { name = "fastapi", specifier = ">=0.136.0" }, { name = "pydantic", specifier = ">=2.13.2" }, ] From 3ee8f4582271f43df083b624d3c7f549d664d78b Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 15:18:47 +0100 Subject: [PATCH 07/20] Avoid duplicate logs --- .../blueapi_interaction/clients.py | 1 + .../blueapi_interaction/session_manager.py | 1 + src/daq_queuing_service/log.py | 1 + src/daq_queuing_service/worker/worker.py | 4 +++- tests/unit_tests/conftest.py | 9 +++++++++ tests/unit_tests/test_worker.py | 11 +++++++++-- 6 files changed, 24 insertions(+), 3 deletions(-) diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index b322011..9c3ef65 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -13,6 +13,7 @@ def get_blueapi_clients( blueapi_config: ApplicationConfig, ) -> tuple[BlueapiRestClient, BlueapiClient]: + # This is only needed until the blueapi client supports udc. if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index e74e877..7273c06 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -7,6 +7,7 @@ class UDCSessionManager(SessionManager): + # This is only needed until the blueapi client supports udc. def get_valid_access_token(self) -> str: token_url = ( "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token" diff --git a/src/daq_queuing_service/log.py b/src/daq_queuing_service/log.py index 3bc74e2..44629f7 100644 --- a/src/daq_queuing_service/log.py +++ b/src/daq_queuing_service/log.py @@ -18,3 +18,4 @@ LOGGER = logging.getLogger("Queue") LOGGER.addHandler(HANDLER) LOGGER.setLevel(logging.DEBUG) +LOGGER.propagate = False diff --git a/src/daq_queuing_service/worker/worker.py b/src/daq_queuing_service/worker/worker.py index 0b60906..63e2ded 100644 --- a/src/daq_queuing_service/worker/worker.py +++ b/src/daq_queuing_service/worker/worker.py @@ -24,6 +24,7 @@ LOGGER = logging.getLogger("Queue Worker") LOGGER.addHandler(HANDLER) LOGGER.setLevel(logging.DEBUG) +LOGGER.propagate = False class QueueWorker: @@ -102,7 +103,8 @@ def _on_blueapi_event(event: AnyEvent, call: BlueapiCall): assert worker_event.task_status call.blueapi_id = worker_event.task_status.task_id LOGGER.info( - f"Call {call} is in progress, blueapi ID: {call.blueapi_id}" + f"Putting call in progress, blueapi ID: {call.blueapi_id}. " + + f"Call: ({call})" ) call.put_in_progress() case ProgressEvent(): diff --git a/tests/unit_tests/conftest.py b/tests/unit_tests/conftest.py index ffeb1c9..a07b65f 100644 --- a/tests/unit_tests/conftest.py +++ b/tests/unit_tests/conftest.py @@ -1,7 +1,9 @@ import pytest from blueapi.worker.event import TaskError, TaskResult +from pytest import MonkeyPatch from daq_queuing_service.broadcaster import Broadcaster +from daq_queuing_service.log import LOGGER from daq_queuing_service.plugins.construct_task_request import ( construct_blueapi_call_list, ) @@ -9,6 +11,13 @@ from daq_queuing_service.task_queue.queue import TaskQueue +@pytest.fixture(autouse=True) +def propagate_logs(monkeypatch: MonkeyPatch): + # This is turned off in prod to avoid duplicate logs + # but needed in tests for caplog to receive logs + monkeypatch.setattr(LOGGER, "propagate", True) + + @pytest.fixture def tasks() -> list[Task]: return [ diff --git a/tests/unit_tests/test_worker.py b/tests/unit_tests/test_worker.py index 9e3d9b7..ea5c867 100644 --- a/tests/unit_tests/test_worker.py +++ b/tests/unit_tests/test_worker.py @@ -14,7 +14,7 @@ from blueapi.core import DataEvent from blueapi.service.model import TaskRequest from blueapi.worker import ProgressEvent, TaskStatus, WorkerEvent, WorkerState -from pytest import LogCaptureFixture +from pytest import LogCaptureFixture, MonkeyPatch from daq_queuing_service.blueapi_interaction.blueapi_adapter import ( BlueapiClientAdapter, @@ -23,7 +23,14 @@ from daq_queuing_service.blueapi_interaction.blueapi_call import CallStatus from daq_queuing_service.task import ExperimentDefinition, Status from daq_queuing_service.task_queue.queue import TaskError, TaskQueue, TaskResult -from daq_queuing_service.worker.worker import QueueWorker +from daq_queuing_service.worker.worker import LOGGER, QueueWorker + + +@pytest.fixture(autouse=True) +def propagate_logs(monkeypatch: MonkeyPatch): + # This is turned off in prod to avoid duplicate logs + # but needed in tests for caplog to receive logs + monkeypatch.setattr(LOGGER, "propagate", True) def _get_mock_blueapi_client( From b91d40c409d06e2460e5790d35d264000a84eb5b Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 15:37:28 +0100 Subject: [PATCH 08/20] Update docstrings + comments --- src/daq_queuing_service/blueapi_interaction/clients.py | 2 +- .../blueapi_interaction/session_manager.py | 6 +++++- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index 9c3ef65..b619c90 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -13,7 +13,7 @@ def get_blueapi_clients( blueapi_config: ApplicationConfig, ) -> tuple[BlueapiRestClient, BlueapiClient]: - # This is only needed until the blueapi client supports udc. + # This should be able to be simplified once the blueapi client supports UDC. if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index 7273c06..107b113 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -7,12 +7,16 @@ class UDCSessionManager(SessionManager): - # This is only needed until the blueapi client supports udc. + """Session manager for a UDC session. Overrides `get_valid_access_token` to get + token using a sealed secret instead of from a file. + """ + def get_valid_access_token(self) -> str: token_url = ( "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token" ) + # Need to get this from secret client_id = "i15-1-udc" client_secret = os.environ.get("UDC_SECRET") From e4a9eaa9eb61c85aeba0e02244f9331c862ece62 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 15:40:40 +0100 Subject: [PATCH 09/20] Update import --- src/daq_queuing_service/blueapi_interaction/session_manager.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index 107b113..1147e7b 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -1,7 +1,7 @@ import os import requests -from blueapi.client.rest import SessionManager +from blueapi.service.authentication import SessionManager from daq_queuing_service.log import LOGGER From 80d3d9ddb5ce4a6e2f42fd4996698172ec1932e7 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 3 Jun 2026 15:42:42 +0100 Subject: [PATCH 10/20] Comment --- src/daq_queuing_service/blueapi_interaction/session_manager.py | 1 + 1 file changed, 1 insertion(+) diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index 1147e7b..1aa7bef 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -9,6 +9,7 @@ class UDCSessionManager(SessionManager): """Session manager for a UDC session. Overrides `get_valid_access_token` to get token using a sealed secret instead of from a file. + It can probably be deleted once the blueapi client supports UDC. """ def get_valid_access_token(self) -> str: From 0da03405365870c6a4a03b2a74328d13f47deed5 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Tue, 9 Jun 2026 10:13:54 +0100 Subject: [PATCH 11/20] Update client_id name in UDC session manager --- src/daq_queuing_service/blueapi_interaction/session_manager.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/session_manager.py index 1aa7bef..56506df 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/session_manager.py @@ -18,7 +18,7 @@ def get_valid_access_token(self) -> str: ) # Need to get this from secret - client_id = "i15-1-udc" + client_id = "i15-1udc" client_secret = os.environ.get("UDC_SECRET") if not (client_id and client_secret): From d281ce6e6e5d9f99b860c971b05875362a9b8d4c Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 10 Jun 2026 13:48:01 +0100 Subject: [PATCH 12/20] Lint yaml files --- tests/test_data/test_blueapi_config.yaml | 2 +- tests/test_data/test_config.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/test_data/test_blueapi_config.yaml b/tests/test_data/test_blueapi_config.yaml index d4860ab..b681665 100644 --- a/tests/test_data/test_blueapi_config.yaml +++ b/tests/test_data/test_blueapi_config.yaml @@ -1,4 +1,4 @@ -api: +api: url: "http://localhost:8000" stomp: enabled: true # All other stomp settings will be ignored if this is false diff --git a/tests/test_data/test_config.yaml b/tests/test_data/test_config.yaml index 76bc1fc..16f7885 100644 --- a/tests/test_data/test_config.yaml +++ b/tests/test_data/test_config.yaml @@ -1,6 +1,6 @@ blueapi_call_constructor: "default" blueapi: - api: + api: url: "http://localhost:8000" stomp: enabled: true # All other stomp settings will be ignored if this is false From ed533781ec43076927b6d0a28ca34104e04c3999 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Wed, 10 Jun 2026 14:23:54 +0100 Subject: [PATCH 13/20] Rework token retriever --- .../blueapi_interaction/clients.py | 6 ++---- .../{session_manager.py => token_retriever.py} | 15 ++++++++------- 2 files changed, 10 insertions(+), 11 deletions(-) rename src/daq_queuing_service/blueapi_interaction/{session_manager.py => token_retriever.py} (71%) diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index b619c90..ff7114f 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -4,10 +4,9 @@ from blueapi.client.event_bus import EventBusClient from blueapi.client.rest import BlueapiRestClient from blueapi.config import ApplicationConfig -from blueapi.service.authentication import SessionCacheManager from bluesky_stomp.messaging import Broker, StompClient -from daq_queuing_service.blueapi_interaction.session_manager import UDCSessionManager +from daq_queuing_service.blueapi_interaction.token_retriever import UDCTokenRetriever def get_blueapi_clients( @@ -17,9 +16,8 @@ def get_blueapi_clients( if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() - session_manager = UDCSessionManager(blueapi_config.oidc, SessionCacheManager(None)) blueapi_rest_client = BlueapiRestClient( - config=blueapi_config.api, session_manager=session_manager + config=blueapi_config.api, token_retreiver=UDCTokenRetriever("i15-1udc") ) if blueapi_config.stomp.enabled: diff --git a/src/daq_queuing_service/blueapi_interaction/session_manager.py b/src/daq_queuing_service/blueapi_interaction/token_retriever.py similarity index 71% rename from src/daq_queuing_service/blueapi_interaction/session_manager.py rename to src/daq_queuing_service/blueapi_interaction/token_retriever.py index 56506df..caae60c 100644 --- a/src/daq_queuing_service/blueapi_interaction/session_manager.py +++ b/src/daq_queuing_service/blueapi_interaction/token_retriever.py @@ -1,27 +1,28 @@ import os import requests -from blueapi.service.authentication import SessionManager from daq_queuing_service.log import LOGGER -class UDCSessionManager(SessionManager): +class UDCTokenRetriever: """Session manager for a UDC session. Overrides `get_valid_access_token` to get token using a sealed secret instead of from a file. It can probably be deleted once the blueapi client supports UDC. """ + def __init__(self, client_id: str, secret_variable_name: str = "UDC_SECRET"): + self._client_id = (client_id,) + self._secret_variable_name = secret_variable_name + def get_valid_access_token(self) -> str: token_url = ( "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token" ) - # Need to get this from secret - client_id = "i15-1udc" - client_secret = os.environ.get("UDC_SECRET") + client_secret = os.environ.get(self._secret_variable_name) - if not (client_id and client_secret): + if not (self._client_id and client_secret): LOGGER.debug("No UDC secret found") return "" @@ -30,7 +31,7 @@ def get_valid_access_token(self) -> str: response = requests.post( token_url, data={ - "client_id": client_id, + "client_id": self._client_id, "client_secret": client_secret, "grant_type": "client_credentials", }, From 99e83e0becfcd1c4b8550ad024748a280d9c9557 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Thu, 11 Jun 2026 12:14:14 +0100 Subject: [PATCH 14/20] Work with current blueapi --- src/daq_queuing_service/blueapi_interaction/clients.py | 5 +++-- .../blueapi_interaction/token_retriever.py | 7 ++----- 2 files changed, 5 insertions(+), 7 deletions(-) diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index ff7114f..a5505b5 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -12,12 +12,13 @@ def get_blueapi_clients( blueapi_config: ApplicationConfig, ) -> tuple[BlueapiRestClient, BlueapiClient]: - # This should be able to be simplified once the blueapi client supports UDC. if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() blueapi_rest_client = BlueapiRestClient( - config=blueapi_config.api, token_retreiver=UDCTokenRetriever("i15-1udc") + config=blueapi_config.api, + # Waiting on https://github.com/DiamondLightSource/blueapi/pull/1553 + session_manager=UDCTokenRetriever("i15-1udc"), # type: ignore ) if blueapi_config.stomp.enabled: diff --git a/src/daq_queuing_service/blueapi_interaction/token_retriever.py b/src/daq_queuing_service/blueapi_interaction/token_retriever.py index caae60c..0d80f68 100644 --- a/src/daq_queuing_service/blueapi_interaction/token_retriever.py +++ b/src/daq_queuing_service/blueapi_interaction/token_retriever.py @@ -6,13 +6,10 @@ class UDCTokenRetriever: - """Session manager for a UDC session. Overrides `get_valid_access_token` to get - token using a sealed secret instead of from a file. - It can probably be deleted once the blueapi client supports UDC. - """ + """Implements `get_valid_access_token` to get a token using a sealed secret.""" def __init__(self, client_id: str, secret_variable_name: str = "UDC_SECRET"): - self._client_id = (client_id,) + self._client_id = client_id self._secret_variable_name = secret_variable_name def get_valid_access_token(self) -> str: From e6570a1145841790444a233af595d58fa5f50a67 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Thu, 11 Jun 2026 13:16:59 +0100 Subject: [PATCH 15/20] Add comment --- src/daq_queuing_service/blueapi_interaction/clients.py | 5 ++++- tests/unit_tests/test_udc_token_retriever.py | 2 ++ 2 files changed, 6 insertions(+), 1 deletion(-) create mode 100644 tests/unit_tests/test_udc_token_retriever.py diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index a5505b5..18ca4e6 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -15,10 +15,13 @@ def get_blueapi_clients( if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() + # TODO: Get this from an env variable or config. + client_id = "i15-1udc" + blueapi_rest_client = BlueapiRestClient( config=blueapi_config.api, # Waiting on https://github.com/DiamondLightSource/blueapi/pull/1553 - session_manager=UDCTokenRetriever("i15-1udc"), # type: ignore + session_manager=UDCTokenRetriever(client_id), # type: ignore ) if blueapi_config.stomp.enabled: diff --git a/tests/unit_tests/test_udc_token_retriever.py b/tests/unit_tests/test_udc_token_retriever.py new file mode 100644 index 0000000..cea2ac1 --- /dev/null +++ b/tests/unit_tests/test_udc_token_retriever.py @@ -0,0 +1,2 @@ +def test_get_valid_access_token(): + assert 0 From fabd2fe433e7b02564895fb5a8b009a9a8c768e2 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Fri, 12 Jun 2026 16:01:36 +0100 Subject: [PATCH 16/20] Fix tests wip --- tests/unit_tests/test_udc_token_retriever.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/unit_tests/test_udc_token_retriever.py b/tests/unit_tests/test_udc_token_retriever.py index cea2ac1..c09471a 100644 --- a/tests/unit_tests/test_udc_token_retriever.py +++ b/tests/unit_tests/test_udc_token_retriever.py @@ -1,2 +1,2 @@ -def test_get_valid_access_token(): - assert 0 +# def test_get_valid_access_token(): +# assert 0 From a49c68521a6d2ffef6305fe1e26e61d21c192ff9 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 22 Jun 2026 11:43:52 +0100 Subject: [PATCH 17/20] Get udc client ID from env variable --- .../blueapi_interaction/clients.py | 5 +---- .../blueapi_interaction/token_retriever.py | 16 ++++++++++++---- 2 files changed, 13 insertions(+), 8 deletions(-) diff --git a/src/daq_queuing_service/blueapi_interaction/clients.py b/src/daq_queuing_service/blueapi_interaction/clients.py index 18ca4e6..5a77357 100644 --- a/src/daq_queuing_service/blueapi_interaction/clients.py +++ b/src/daq_queuing_service/blueapi_interaction/clients.py @@ -15,13 +15,10 @@ def get_blueapi_clients( if not blueapi_config.oidc: blueapi_config.oidc = MagicMock() - # TODO: Get this from an env variable or config. - client_id = "i15-1udc" - blueapi_rest_client = BlueapiRestClient( config=blueapi_config.api, # Waiting on https://github.com/DiamondLightSource/blueapi/pull/1553 - session_manager=UDCTokenRetriever(client_id), # type: ignore + session_manager=UDCTokenRetriever(), # type: ignore ) if blueapi_config.stomp.enabled: diff --git a/src/daq_queuing_service/blueapi_interaction/token_retriever.py b/src/daq_queuing_service/blueapi_interaction/token_retriever.py index 0d80f68..8f82745 100644 --- a/src/daq_queuing_service/blueapi_interaction/token_retriever.py +++ b/src/daq_queuing_service/blueapi_interaction/token_retriever.py @@ -8,27 +8,35 @@ class UDCTokenRetriever: """Implements `get_valid_access_token` to get a token using a sealed secret.""" - def __init__(self, client_id: str, secret_variable_name: str = "UDC_SECRET"): - self._client_id = client_id + def __init__( + self, + secret_variable_name: str = "UDC_SECRET", + client_id_variable_name: str = "UDC_CLIENT_ID", + ): self._secret_variable_name = secret_variable_name + self._client_id_variable_name = client_id_variable_name def get_valid_access_token(self) -> str: token_url = ( "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token" ) + client_id = os.environ.get(self._client_id_variable_name) client_secret = os.environ.get(self._secret_variable_name) - if not (self._client_id and client_secret): + if not client_secret: LOGGER.debug("No UDC secret found") return "" + if not client_id: + LOGGER.debug("No UDC client ID found") + return "" LOGGER.debug("Found UDC secret") response = requests.post( token_url, data={ - "client_id": self._client_id, + "client_id": client_id, "client_secret": client_secret, "grant_type": "client_credentials", }, From 97533d794d1326f930481c42239cacdd6b7a2cb2 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 22 Jun 2026 12:29:30 +0100 Subject: [PATCH 18/20] Add client ID to env var in helm chart --- helm/daq-queuing-service/templates/deployment.yaml | 10 +++++++++- helm/daq-queuing-service/values.yaml | 3 +++ 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/helm/daq-queuing-service/templates/deployment.yaml b/helm/daq-queuing-service/templates/deployment.yaml index 5f11ea2..f4d59d3 100644 --- a/helm/daq-queuing-service/templates/deployment.yaml +++ b/helm/daq-queuing-service/templates/deployment.yaml @@ -60,13 +60,21 @@ spec: {{- with .Values.volumeMounts }} {{- toYaml . | nindent 12 }} {{- end }} - {{- if .Values.udcSecret.enabled }} + {{- if or .Values.udcSecret.enabled .Values.env }} env: + {{- if .Values.udcSecret.enabled }} - name: UDC_SECRET valueFrom: secretKeyRef: name: {{ .Values.udcSecret.name }} key: {{ .Values.udcSecret.key }} + + - name: UDC_CLIENT_ID + value: "{{ .Values.udcSecret.clientId }}" + {{- end }} + {{- with .Values.env }} + {{- toYaml . | nindent 12 }} + {{- end }} {{- end }} volumes: {{- with .Values.volumes }} diff --git a/helm/daq-queuing-service/values.yaml b/helm/daq-queuing-service/values.yaml index 5071db5..3af50af 100644 --- a/helm/daq-queuing-service/values.yaml +++ b/helm/daq-queuing-service/values.yaml @@ -74,7 +74,10 @@ volumes: [] nodeSelector: {} tolerations: [] +env: [] + udcSecret: enabled: false name: "" key: udc-secret + clientId: "" From 41876f1573df1ac36dbcd0e5fe510f16bd86b12b Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 22 Jun 2026 14:02:43 +0100 Subject: [PATCH 19/20] Add tests for token retriever --- tests/unit_tests/test_udc_token_retriever.py | 47 +++++++++++++++++++- 1 file changed, 45 insertions(+), 2 deletions(-) diff --git a/tests/unit_tests/test_udc_token_retriever.py b/tests/unit_tests/test_udc_token_retriever.py index c09471a..9dcef74 100644 --- a/tests/unit_tests/test_udc_token_retriever.py +++ b/tests/unit_tests/test_udc_token_retriever.py @@ -1,2 +1,45 @@ -# def test_get_valid_access_token(): -# assert 0 +from unittest.mock import MagicMock, patch + +import pytest +from pytest import MonkeyPatch +from requests import Response + +from daq_queuing_service.blueapi_interaction.token_retriever import UDCTokenRetriever + + +@pytest.fixture(autouse=True) +def set_secret_env_vars(monkeypatch: MonkeyPatch): + monkeypatch.setenv("UDC_SECRET", "secret") + monkeypatch.setenv("UDC_CLIENT_ID", "ixxudc") + + +@patch("daq_queuing_service.blueapi_interaction.token_retriever.requests.post") +def test_get_valid_access_token_makes_expected_request_and_returns_result( + mock_post: MagicMock, +): + mock_post.return_value = Response() + mock_post.return_value.json = MagicMock( + return_value={"access_token": "valid_token"} + ) + mock_post.return_value.status_code = 200 + token_retriever = UDCTokenRetriever() + + assert token_retriever.get_valid_access_token() == "valid_token" + + mock_post.assert_called_once_with( + "https://identity.diamond.ac.uk/realms/dls/protocol/openid-connect/token", + data={ + "client_id": "ixxudc", + "client_secret": "secret", + "grant_type": "client_credentials", + }, + ) + + +@pytest.mark.parametrize("variable_name", ("UDC_SECRET", "UDC_CLIENT_ID")) +def test_get_valid_access_token_if_env_vars_not_found_then_returns_empty_string( + variable_name: str, monkeypatch: MonkeyPatch +): + monkeypatch.delenv(variable_name) + token_retriever = UDCTokenRetriever() + assert token_retriever.get_valid_access_token() == "" From a41041672f390e5567dff6fc6dc5bbe6af7fa7b3 Mon Sep 17 00:00:00 2001 From: Jacob Williamson Date: Mon, 22 Jun 2026 14:56:34 +0100 Subject: [PATCH 20/20] Add tests for get_blueapi_clients --- tests/unit_tests/test_get_blueapi_clients.py | 46 ++++++++++++++++++++ 1 file changed, 46 insertions(+) create mode 100644 tests/unit_tests/test_get_blueapi_clients.py diff --git a/tests/unit_tests/test_get_blueapi_clients.py b/tests/unit_tests/test_get_blueapi_clients.py new file mode 100644 index 0000000..bf57b37 --- /dev/null +++ b/tests/unit_tests/test_get_blueapi_clients.py @@ -0,0 +1,46 @@ +from unittest.mock import MagicMock, patch + +from blueapi.config import ApplicationConfig, RestConfig, StompConfig +from pydantic import HttpUrl + +from daq_queuing_service.blueapi_interaction.clients import get_blueapi_clients + + +@patch("daq_queuing_service.blueapi_interaction.clients.UDCTokenRetriever") +@patch("daq_queuing_service.blueapi_interaction.clients.BlueapiClient") +@patch("daq_queuing_service.blueapi_interaction.clients.BlueapiRestClient") +def test_get_blueapi_clients_constructs_clients_with_expected_args_and_returns_clients( + mock_rest_client: MagicMock, + mock_blueapi_client: MagicMock, + mock_token_retriever: MagicMock, +): + rest_config = RestConfig(url=HttpUrl("http://test_url.com")) + rest_client, blueapi_client = get_blueapi_clients( + ApplicationConfig(api=rest_config) + ) + + mock_rest_client.assert_called_once_with( + config=rest_config, session_manager=mock_token_retriever.return_value + ) + mock_blueapi_client.assert_called_once_with(rest_client) + + assert rest_client is mock_rest_client.return_value + assert blueapi_client is mock_blueapi_client.return_value + + +@patch("daq_queuing_service.blueapi_interaction.clients.EventBusClient") +@patch("daq_queuing_service.blueapi_interaction.clients.BlueapiClient") +@patch("daq_queuing_service.blueapi_interaction.clients.BlueapiRestClient") +def test_get_blueapi_clients_constructs_blueapi_client_with_stomp_if_enabled_in_config( + mock_rest_client: MagicMock, + mock_blueapi_client: MagicMock, + mock_event_bus_client: MagicMock, +): + rest_config = RestConfig(url=HttpUrl("http://test_url.com")) + rest_client, _ = get_blueapi_clients( + ApplicationConfig(api=rest_config, stomp=StompConfig(enabled=True)) + ) + + mock_blueapi_client.assert_called_once_with( + rest_client, mock_event_bus_client.return_value + )