From 206b54cb33f1dd2899c6425f3cefe2a23497ac8c Mon Sep 17 00:00:00 2001 From: Daniel Fernandes Date: Tue, 22 Sep 2026 10:53:50 +0000 Subject: [PATCH 1/5] Add message-broker.md documentation --- docs/explanations/message-brokers.md | 49 ++++++++++++++++++++++++++++ 1 file changed, 49 insertions(+) create mode 100644 docs/explanations/message-brokers.md diff --git a/docs/explanations/message-brokers.md b/docs/explanations/message-brokers.md new file mode 100644 index 000000000..0b96e6b36 --- /dev/null +++ b/docs/explanations/message-brokers.md @@ -0,0 +1,49 @@ +# Message Brokers + +Blueapi uses a message broker to communicate with downstream services such as the Nexus Filewriter, and tools such as the Blueapi CLI client. Blueapi is broker agnostic, communicating in STOMP via the [bluesky-stomp](https://github.com/DiamondLightSource/bluesky-stomp) library. However, RabbitMQ is most commonly used in DLS deployments. + +## Messages + +When a plan is run, Blueapi broadcasts all worker events, progress events and data events to the configured message broker. Some example messages can be seen below. These are the messages sent to RabbitMQ during a `sleep` plan (printed with an `[x]` at the start): +``` + [x] public.worker.event:b'{"state":"RUNNING","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' + [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' + [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":{"outcome":"success","result":null,"type":"NoneType"},"task_complete":true,"task_failed":false},"errors":[],"warnings":[]}' + [x] public.worker.event:b'{"state":"RUNNING","task_status":{"task_id":"a3c51955-943b-4099-85f5-8263a6c02c53","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' + [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"a3c51955-943b-4099-85f5-8263a6c02c53","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' + [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"a3c51955-943b-4099-85f5-8263a6c02c53","result":{"outcome":"success","result":null,"type":"NoneType"},"task_complete":true,"task_failed":false},"errors":[],"warnings":[]}' +``` + +## Topology +The topology of the message broker is entirely deployment-specific. Nonetheless, the following details how deployments at DLS are generally structured. + + +RabbitMQ receives messages from Blueapi via its default topic exchange, `amq.topic`. All messages have the routing key `public.worker.event`, so will be forwarded to all queues bound to `amq.topic` with a compatible routing key (eg. `public.worker.event`, `*.*.event` etc.). + +No default queues exist, meaning that if no downstream services have registered a queue, all messages will be dropped. + +In RabbitMQ, different queues can have the same routing key, and consequently receive a copy of the same messages without interfering with each other. In Blueapi a common pair of queues to have is `public.worker.event.nexus` and `stomp-subscription-{random GUID}`, both bound to the `amq.topic` exchange with routing key `public.worker.event`. The prior is for the Nexus Filewriter, and the latter the Blueapi CLI client. In this case both have access to all messages published by Blueapi, but due to automatically declaring seperate queues at subscription, each consumer's message consumption will not deprive the other. + + +## Blueapi Messaging Mechanism + +When Blueapi starts up, a `StompClient` is instantiated which registers a callback to each event stream (worker, progress and data). When an event occurs, it will be sent to the broker's default topic exchange with the routing key `public.worker.event`. In RabbitMQ, the default topic exchange is `amq.topic`. + + +## Subscribing to Blueapi Messages + +When Blueapi is configured to use a message broker, services can subscribe to the broker to receive events generated during plan execution. In general, having services subscribe directly to the message broker is not recommended, as there should be a more appropriate way to access the information you want. + +Due to the dynamic nature of RabbitMQ, new services can subscribe to the exchange at any point, without reconfiguring the server or interfering with other consumer's subscriptions. + +To to connect a consumer to a RabbitMQ topic, follow RabbitMQ's [Tutorial 5: Topics](https://www.rabbitmq.com/tutorials) in your preferred language. Change the exchange name to `amq.topic` and the routing key to `public.worker.event`. + +The consumer created in this tutorial will capture all messages generated during its runtime. It will not have access to messages that were generated previously. If this is a requirement for your use-case, RabbitMQ can be reconfigured to generate specific queues at startup. This can be achieved through [schema definitions](https://www.rabbitmq.com/docs/definitions). + +## Tips + +It is important to remember that while all queues are guaranteed to receive the same set of messages, there is no guarantee that each queue has consumed their messages. + +For example, a data analysis service listening for a Stop document may process events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. + +The only guarantee we make is that all queues will receive all events in the correct order. From 905001d693f5c4ba4ecd929f33f069f26c5a0caf Mon Sep 17 00:00:00 2001 From: Daniel Fernandes Date: Tue, 22 Sep 2026 10:55:40 +0000 Subject: [PATCH 2/5] Reword tips --- docs/explanations/message-brokers.md | 2 +- helm/blueapi/output.yaml | 326 +++++++++++++++++++++++++++ helm/blueapi/overriding_values.yaml | 101 +++++++++ 3 files changed, 428 insertions(+), 1 deletion(-) create mode 100755 helm/blueapi/output.yaml create mode 100755 helm/blueapi/overriding_values.yaml diff --git a/docs/explanations/message-brokers.md b/docs/explanations/message-brokers.md index 0b96e6b36..60233551b 100644 --- a/docs/explanations/message-brokers.md +++ b/docs/explanations/message-brokers.md @@ -44,6 +44,6 @@ The consumer created in this tutorial will capture all messages generated during It is important to remember that while all queues are guaranteed to receive the same set of messages, there is no guarantee that each queue has consumed their messages. -For example, a data analysis service listening for a Stop document may process events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. +For example, a service listening for a Stop document to kick off data analysis may consume events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. The only guarantee we make is that all queues will receive all events in the correct order. diff --git a/helm/blueapi/output.yaml b/helm/blueapi/output.yaml new file mode 100755 index 000000000..80a327438 --- /dev/null +++ b/helm/blueapi/output.yaml @@ -0,0 +1,326 @@ +--- +# Source: blueapi/templates/configmap.yaml +apiVersion: v1 +kind: ConfigMap +metadata: + name: release-name-blueapi-config +data: + config.yaml: |- + api: + url: http://0.0.0.0:8000/ + env: + events: + broadcast_status_events: false + metadata: + instrument: p48 + sources: + - kind: deviceManager + module: dodal.beamlines.training_rig + - kind: planFunctions + module: dodal.plans + logging: + graylog: + enabled: true + url: tcp://graylog-log-target.diamond.ac.uk:12231/ + level: INFO + numtracker: + url: https://numtracker-staging.diamond.ac.uk/graphql + oidc: + client_audience: account + client_id: blueapiCli + logout_redirect_endpoint: oauth2/sign_out + well_known_url: https://identity.diamond.ac.uk/realms/dls/.well-known/openid-configuration + scratch: + repositories: + - name: dodal + remote_url: https://github.com/DiamondLightSource/dodal.git + - name: htss-rig-bluesky + remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git + target_revision: update-copier-template + root: /workspaces + stomp: + auth: + password: guest + username: guest + enabled: true + url: tcp://p48-rabbitmq-daq.diamond.ac.uk:61613 + tiled: + authentication: + client_id: p48TiledWriter + client_secret: ${TILED_WRITER_SECRET} + enabled: true + url: https://tiled.diamond.ac.uk/api/v1/ +--- +# Source: blueapi/templates/configmap.yaml +apiVersion: v1 +kind: ConfigMap +metadata: + name: release-name-blueapi-otel-config +data: +--- +# Source: blueapi/templates/configmap.yaml +apiVersion: v1 +kind: ConfigMap +metadata: + name: release-name-blueapi-init-config +data: + init_config.yaml: |- + scratch: + repositories: + - name: dodal + remote_url: https://github.com/DiamondLightSource/dodal.git + - name: htss-rig-bluesky + remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git + target_revision: update-copier-template + root: /workspaces +--- +# Source: blueapi/templates/tests/test-connection.yaml +apiVersion: v1 +kind: ConfigMap +metadata: + name: release-name-blueapi-test-config +data: + init_config.yaml: |- + api: + url: http://release-name-blueapi:80/ + stomp: + enabled: false + auth: + username: guest + password: guest + url: http://rabbitmq:61613/ + logging: + level: "INFO" + graylog: + enabled: False + url: http://graylog-log-target.diamond.ac.uk:12232/ +--- +# Source: blueapi/templates/volumes.yaml +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: release-name-blueapi-scratch-0.1.0 + annotations: + argocd.argoproj.io/sync-options: Prune=false,Delete=false + argocd.argoproj.io/compare-options: IgnoreExtraneous +spec: + accessModes: + - ReadWriteMany + resources: + requests: + storage: 2Gi +--- +# Source: blueapi/templates/service.yaml +apiVersion: v1 +kind: Service +metadata: + name: release-name-blueapi + labels: + helm.sh/chart: blueapi-0.1.0 + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name + app.kubernetes.io/version: "0.1.0" + app.kubernetes.io/managed-by: Helm +spec: + type: ClusterIP + ports: + - port: 80 + targetPort: http + protocol: TCP + name: http + selector: + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name +--- +# Source: blueapi/templates/statefulset.yaml +apiVersion: apps/v1 +kind: StatefulSet +metadata: + name: release-name-blueapi + labels: + helm.sh/chart: blueapi-0.1.0 + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name + app.kubernetes.io/version: "0.1.0" + app.kubernetes.io/managed-by: Helm +spec: + selector: + matchLabels: + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name + template: + metadata: + annotations: + fluentd-ignore: "true" + checksum/config: e1f22a59a75c2819feb0c50ed96a46e42b9aa293571acfa9740d7642e8f22af2 + labels: + helm.sh/chart: blueapi-0.1.0 + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name + app.kubernetes.io/version: "0.1.0" + app.kubernetes.io/managed-by: Helm + spec: + serviceAccountName: default + volumes: + - name: worker-config + projected: + sources: + - configMap: + name: release-name-blueapi-config + - name: init-config + projected: + sources: + - configMap: + name: release-name-blueapi-init-config + - name: venv + emptyDir: + sizeLimit: 5Gi + - name: scratch + persistentVolumeClaim: + claimName: release-name-blueapi-scratch-0.1.0 + initContainers: + - name: setup-scratch + securityContext: + runAsNonRoot: true + runAsUser: 1000 + image: "ghcr.io/dan-fernandes/blueapi:0.1.0" + imagePullPolicy: IfNotPresent + resources: + limits: + cpu: 2000m + memory: 4000Mi + requests: + cpu: 200m + memory: 400Mi + command: ["/bin/sh", "-c"] + args: + - | + echo "Setting up scratch area" + blueapi -c /config/init_config.yaml setup-scratch + if [ $? -ne 0 ]; then echo 'Blueapi failed'; exit 1; fi; + echo "Exporting venv as artefact" + cp -r /app/.venv/* /artefacts + env: + - name: UV_CACHE_DIR + value: /app/.uv-cache + volumeMounts: + - name: init-config + mountPath: "/config" + readOnly: true + - name: venv + mountPath: /artefacts + - name: scratch + mountPath: /workspaces + containers: + - name: blueapi + securityContext: + runAsNonRoot: true + runAsUser: 1000 + image: "ghcr.io/dan-fernandes/blueapi:0.1.0" + imagePullPolicy: IfNotPresent + ports: + - name: http + containerPort: 8000 + protocol: TCP + resources: + limits: + cpu: 2000m + memory: 4000Mi + requests: + cpu: 200m + memory: 400Mi + volumeMounts: + - name: worker-config + mountPath: "/config" + readOnly: true + - name: scratch + mountPath: /workspaces + - name: venv + mountPath: /app/.venv + args: + - "-c" + - "/config/config.yaml" + - "serve" + command: + - "python" + - "-Xfrozen_modules=off" + - "-m" + - "debugpy" + - "--listen" + - "5678" + - "--configure-subProcess" + - "true" + - "-m" + - "blueapi" + envFrom: + - configMapRef: + name: release-name-blueapi-otel-config + env: + - name: UV_CACHE_DIR + value: /app/.uv-cache + - name: DEBUG_MODE + value: "ON" + - name: BEAMLINE + value: p48 + - name: INSTRUMENT + value: p48 + - name: EPICS_PVA_NAME_SERVERS + value: p48-epics-gateways:9075 + - name: EPICS_CA_NAME_SERVERS + value: p48-epics-gateways:9064 + - name: EPICS_PVA_AUTO_ADDR_LIST + value: "NO" + - name: EPICS_CA_AUTO_ADDR_LIST + value: "NO" + - name: TILED_WRITER_SECRET + valueFrom: + secretKeyRef: + key: tiled-writer-secret + name: blueapi-secret + nodeSelector: + beamline: p48 + tolerations: + - effect: NoSchedule + key: beamline + operator: Equal + value: p48 + - effect: NoSchedule + key: nodetype + operator: Equal + value: training-rig +--- +# Source: blueapi/templates/tests/test-connection.yaml +apiVersion: v1 +kind: Pod +metadata: + name: "release-name-blueapi-test-connection" + labels: + helm.sh/chart: blueapi-0.1.0 + app.kubernetes.io/name: blueapi + app.kubernetes.io/instance: release-name + app.kubernetes.io/version: "0.1.0" + app.kubernetes.io/managed-by: Helm + annotations: + "helm.sh/hook": test +spec: + volumes: + - name: test-config + projected: + sources: + - configMap: + name: release-name-blueapi-test-config + containers: + - name: ping + volumeMounts: + - name: worker-config + mountPath: "/config" + readOnly: true + image: "ghcr.io/dan-fernandes/blueapi:0.1.0" + imagePullPolicy: IfNotPresent + command: ["blueapi"] + args: + - "-c" + - "/config/config.yaml" + - "controller" + - "plans" + restartPolicy: Never diff --git a/helm/blueapi/overriding_values.yaml b/helm/blueapi/overriding_values.yaml new file mode 100755 index 000000000..58137547e --- /dev/null +++ b/helm/blueapi/overriding_values.yaml @@ -0,0 +1,101 @@ +image: + repository: ghcr.io/dan-fernandes/blueapi + # tag: "1.12.13" +debug: + enabled: true +hostNetwork: false +ingress: + enabled: false + hosts: + - host: p48-blueapi.diamond.ac.uk + paths: + - path: / + pathType: Prefix + +nodeSelector: + beamline: p48 + +tolerations: + - key: beamline + operator: Equal + value: p48 + effect: NoSchedule + - key: nodetype + operator: Equal + value: training-rig + effect: NoSchedule + +extraEnvVars: + - name: BEAMLINE + value: p48 + - name: INSTRUMENT + value: p48 + - name: EPICS_PVA_NAME_SERVERS + value: p48-epics-gateways:9075 + - name: EPICS_CA_NAME_SERVERS + value: p48-epics-gateways:9064 + - name: EPICS_PVA_AUTO_ADDR_LIST + value: "NO" + - name: EPICS_CA_AUTO_ADDR_LIST + value: "NO" + + - name: TILED_WRITER_SECRET + valueFrom: + secretKeyRef: + name: blueapi-secret + key: tiled-writer-secret + +worker: + env: + metadata: + instrument: p48 + sources: + - kind: deviceManager + module: dodal.beamlines.training_rig + - kind: planFunctions + module: dodal.plans + events: + broadcast_status_events: false + stomp: + auth: + username: guest + password: guest + url: tcp://p48-rabbitmq-daq.diamond.ac.uk:61613 + enabled: true + logging: + level: "INFO" + graylog: + enabled: true + + oidc: + well_known_url: "https://identity.diamond.ac.uk/realms/dls/.well-known/openid-configuration" + client_id: "blueapiCli" + client_audience: "account" + logout_redirect_endpoint: "oauth2/sign_out" + + numtracker: + url: https://numtracker-staging.diamond.ac.uk/graphql + + tiled: + enabled: true + url: "https://tiled.diamond.ac.uk/api/v1/" + authentication: + client_id: "p48TiledWriter" + client_secret: ${TILED_WRITER_SECRET} + + scratch: + root: /workspaces + repositories: + - name: dodal + remote_url: https://github.com/DiamondLightSource/dodal.git + - name: htss-rig-bluesky + remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git + target_revision: update-copier-template + # - name: ophyd-async + # remote_url: https://github.com/bluesky/ophyd-async.git + +initContainer: + enabled: true + persistentVolume: + enabled: true + size: "2Gi" From b12dfab8621266c024d9cb15888d4713e1258feb Mon Sep 17 00:00:00 2001 From: Daniel Fernandes Date: Tue, 22 Sep 2026 12:44:57 +0000 Subject: [PATCH 3/5] Revert "Reword tips" This reverts commit 905001d693f5c4ba4ecd929f33f069f26c5a0caf. --- docs/explanations/message-brokers.md | 2 +- helm/blueapi/output.yaml | 326 --------------------------- helm/blueapi/overriding_values.yaml | 101 --------- 3 files changed, 1 insertion(+), 428 deletions(-) delete mode 100755 helm/blueapi/output.yaml delete mode 100755 helm/blueapi/overriding_values.yaml diff --git a/docs/explanations/message-brokers.md b/docs/explanations/message-brokers.md index 60233551b..0b96e6b36 100644 --- a/docs/explanations/message-brokers.md +++ b/docs/explanations/message-brokers.md @@ -44,6 +44,6 @@ The consumer created in this tutorial will capture all messages generated during It is important to remember that while all queues are guaranteed to receive the same set of messages, there is no guarantee that each queue has consumed their messages. -For example, a service listening for a Stop document to kick off data analysis may consume events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. +For example, a data analysis service listening for a Stop document may process events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. The only guarantee we make is that all queues will receive all events in the correct order. diff --git a/helm/blueapi/output.yaml b/helm/blueapi/output.yaml deleted file mode 100755 index 80a327438..000000000 --- a/helm/blueapi/output.yaml +++ /dev/null @@ -1,326 +0,0 @@ ---- -# Source: blueapi/templates/configmap.yaml -apiVersion: v1 -kind: ConfigMap -metadata: - name: release-name-blueapi-config -data: - config.yaml: |- - api: - url: http://0.0.0.0:8000/ - env: - events: - broadcast_status_events: false - metadata: - instrument: p48 - sources: - - kind: deviceManager - module: dodal.beamlines.training_rig - - kind: planFunctions - module: dodal.plans - logging: - graylog: - enabled: true - url: tcp://graylog-log-target.diamond.ac.uk:12231/ - level: INFO - numtracker: - url: https://numtracker-staging.diamond.ac.uk/graphql - oidc: - client_audience: account - client_id: blueapiCli - logout_redirect_endpoint: oauth2/sign_out - well_known_url: https://identity.diamond.ac.uk/realms/dls/.well-known/openid-configuration - scratch: - repositories: - - name: dodal - remote_url: https://github.com/DiamondLightSource/dodal.git - - name: htss-rig-bluesky - remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git - target_revision: update-copier-template - root: /workspaces - stomp: - auth: - password: guest - username: guest - enabled: true - url: tcp://p48-rabbitmq-daq.diamond.ac.uk:61613 - tiled: - authentication: - client_id: p48TiledWriter - client_secret: ${TILED_WRITER_SECRET} - enabled: true - url: https://tiled.diamond.ac.uk/api/v1/ ---- -# Source: blueapi/templates/configmap.yaml -apiVersion: v1 -kind: ConfigMap -metadata: - name: release-name-blueapi-otel-config -data: ---- -# Source: blueapi/templates/configmap.yaml -apiVersion: v1 -kind: ConfigMap -metadata: - name: release-name-blueapi-init-config -data: - init_config.yaml: |- - scratch: - repositories: - - name: dodal - remote_url: https://github.com/DiamondLightSource/dodal.git - - name: htss-rig-bluesky - remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git - target_revision: update-copier-template - root: /workspaces ---- -# Source: blueapi/templates/tests/test-connection.yaml -apiVersion: v1 -kind: ConfigMap -metadata: - name: release-name-blueapi-test-config -data: - init_config.yaml: |- - api: - url: http://release-name-blueapi:80/ - stomp: - enabled: false - auth: - username: guest - password: guest - url: http://rabbitmq:61613/ - logging: - level: "INFO" - graylog: - enabled: False - url: http://graylog-log-target.diamond.ac.uk:12232/ ---- -# Source: blueapi/templates/volumes.yaml -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - name: release-name-blueapi-scratch-0.1.0 - annotations: - argocd.argoproj.io/sync-options: Prune=false,Delete=false - argocd.argoproj.io/compare-options: IgnoreExtraneous -spec: - accessModes: - - ReadWriteMany - resources: - requests: - storage: 2Gi ---- -# Source: blueapi/templates/service.yaml -apiVersion: v1 -kind: Service -metadata: - name: release-name-blueapi - labels: - helm.sh/chart: blueapi-0.1.0 - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name - app.kubernetes.io/version: "0.1.0" - app.kubernetes.io/managed-by: Helm -spec: - type: ClusterIP - ports: - - port: 80 - targetPort: http - protocol: TCP - name: http - selector: - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name ---- -# Source: blueapi/templates/statefulset.yaml -apiVersion: apps/v1 -kind: StatefulSet -metadata: - name: release-name-blueapi - labels: - helm.sh/chart: blueapi-0.1.0 - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name - app.kubernetes.io/version: "0.1.0" - app.kubernetes.io/managed-by: Helm -spec: - selector: - matchLabels: - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name - template: - metadata: - annotations: - fluentd-ignore: "true" - checksum/config: e1f22a59a75c2819feb0c50ed96a46e42b9aa293571acfa9740d7642e8f22af2 - labels: - helm.sh/chart: blueapi-0.1.0 - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name - app.kubernetes.io/version: "0.1.0" - app.kubernetes.io/managed-by: Helm - spec: - serviceAccountName: default - volumes: - - name: worker-config - projected: - sources: - - configMap: - name: release-name-blueapi-config - - name: init-config - projected: - sources: - - configMap: - name: release-name-blueapi-init-config - - name: venv - emptyDir: - sizeLimit: 5Gi - - name: scratch - persistentVolumeClaim: - claimName: release-name-blueapi-scratch-0.1.0 - initContainers: - - name: setup-scratch - securityContext: - runAsNonRoot: true - runAsUser: 1000 - image: "ghcr.io/dan-fernandes/blueapi:0.1.0" - imagePullPolicy: IfNotPresent - resources: - limits: - cpu: 2000m - memory: 4000Mi - requests: - cpu: 200m - memory: 400Mi - command: ["/bin/sh", "-c"] - args: - - | - echo "Setting up scratch area" - blueapi -c /config/init_config.yaml setup-scratch - if [ $? -ne 0 ]; then echo 'Blueapi failed'; exit 1; fi; - echo "Exporting venv as artefact" - cp -r /app/.venv/* /artefacts - env: - - name: UV_CACHE_DIR - value: /app/.uv-cache - volumeMounts: - - name: init-config - mountPath: "/config" - readOnly: true - - name: venv - mountPath: /artefacts - - name: scratch - mountPath: /workspaces - containers: - - name: blueapi - securityContext: - runAsNonRoot: true - runAsUser: 1000 - image: "ghcr.io/dan-fernandes/blueapi:0.1.0" - imagePullPolicy: IfNotPresent - ports: - - name: http - containerPort: 8000 - protocol: TCP - resources: - limits: - cpu: 2000m - memory: 4000Mi - requests: - cpu: 200m - memory: 400Mi - volumeMounts: - - name: worker-config - mountPath: "/config" - readOnly: true - - name: scratch - mountPath: /workspaces - - name: venv - mountPath: /app/.venv - args: - - "-c" - - "/config/config.yaml" - - "serve" - command: - - "python" - - "-Xfrozen_modules=off" - - "-m" - - "debugpy" - - "--listen" - - "5678" - - "--configure-subProcess" - - "true" - - "-m" - - "blueapi" - envFrom: - - configMapRef: - name: release-name-blueapi-otel-config - env: - - name: UV_CACHE_DIR - value: /app/.uv-cache - - name: DEBUG_MODE - value: "ON" - - name: BEAMLINE - value: p48 - - name: INSTRUMENT - value: p48 - - name: EPICS_PVA_NAME_SERVERS - value: p48-epics-gateways:9075 - - name: EPICS_CA_NAME_SERVERS - value: p48-epics-gateways:9064 - - name: EPICS_PVA_AUTO_ADDR_LIST - value: "NO" - - name: EPICS_CA_AUTO_ADDR_LIST - value: "NO" - - name: TILED_WRITER_SECRET - valueFrom: - secretKeyRef: - key: tiled-writer-secret - name: blueapi-secret - nodeSelector: - beamline: p48 - tolerations: - - effect: NoSchedule - key: beamline - operator: Equal - value: p48 - - effect: NoSchedule - key: nodetype - operator: Equal - value: training-rig ---- -# Source: blueapi/templates/tests/test-connection.yaml -apiVersion: v1 -kind: Pod -metadata: - name: "release-name-blueapi-test-connection" - labels: - helm.sh/chart: blueapi-0.1.0 - app.kubernetes.io/name: blueapi - app.kubernetes.io/instance: release-name - app.kubernetes.io/version: "0.1.0" - app.kubernetes.io/managed-by: Helm - annotations: - "helm.sh/hook": test -spec: - volumes: - - name: test-config - projected: - sources: - - configMap: - name: release-name-blueapi-test-config - containers: - - name: ping - volumeMounts: - - name: worker-config - mountPath: "/config" - readOnly: true - image: "ghcr.io/dan-fernandes/blueapi:0.1.0" - imagePullPolicy: IfNotPresent - command: ["blueapi"] - args: - - "-c" - - "/config/config.yaml" - - "controller" - - "plans" - restartPolicy: Never diff --git a/helm/blueapi/overriding_values.yaml b/helm/blueapi/overriding_values.yaml deleted file mode 100755 index 58137547e..000000000 --- a/helm/blueapi/overriding_values.yaml +++ /dev/null @@ -1,101 +0,0 @@ -image: - repository: ghcr.io/dan-fernandes/blueapi - # tag: "1.12.13" -debug: - enabled: true -hostNetwork: false -ingress: - enabled: false - hosts: - - host: p48-blueapi.diamond.ac.uk - paths: - - path: / - pathType: Prefix - -nodeSelector: - beamline: p48 - -tolerations: - - key: beamline - operator: Equal - value: p48 - effect: NoSchedule - - key: nodetype - operator: Equal - value: training-rig - effect: NoSchedule - -extraEnvVars: - - name: BEAMLINE - value: p48 - - name: INSTRUMENT - value: p48 - - name: EPICS_PVA_NAME_SERVERS - value: p48-epics-gateways:9075 - - name: EPICS_CA_NAME_SERVERS - value: p48-epics-gateways:9064 - - name: EPICS_PVA_AUTO_ADDR_LIST - value: "NO" - - name: EPICS_CA_AUTO_ADDR_LIST - value: "NO" - - - name: TILED_WRITER_SECRET - valueFrom: - secretKeyRef: - name: blueapi-secret - key: tiled-writer-secret - -worker: - env: - metadata: - instrument: p48 - sources: - - kind: deviceManager - module: dodal.beamlines.training_rig - - kind: planFunctions - module: dodal.plans - events: - broadcast_status_events: false - stomp: - auth: - username: guest - password: guest - url: tcp://p48-rabbitmq-daq.diamond.ac.uk:61613 - enabled: true - logging: - level: "INFO" - graylog: - enabled: true - - oidc: - well_known_url: "https://identity.diamond.ac.uk/realms/dls/.well-known/openid-configuration" - client_id: "blueapiCli" - client_audience: "account" - logout_redirect_endpoint: "oauth2/sign_out" - - numtracker: - url: https://numtracker-staging.diamond.ac.uk/graphql - - tiled: - enabled: true - url: "https://tiled.diamond.ac.uk/api/v1/" - authentication: - client_id: "p48TiledWriter" - client_secret: ${TILED_WRITER_SECRET} - - scratch: - root: /workspaces - repositories: - - name: dodal - remote_url: https://github.com/DiamondLightSource/dodal.git - - name: htss-rig-bluesky - remote_url: https://github.com/DiamondLightSource/htss-rig-bluesky.git - target_revision: update-copier-template - # - name: ophyd-async - # remote_url: https://github.com/bluesky/ophyd-async.git - -initContainer: - enabled: true - persistentVolume: - enabled: true - size: "2Gi" From 2248f599e7b5a09af519e327cea6808e33ca18d2 Mon Sep 17 00:00:00 2001 From: Daniel Fernandes Date: Tue, 22 Sep 2026 12:47:31 +0000 Subject: [PATCH 4/5] Flesh out message-brokers.md --- docs/explanations/message-brokers.md | 60 ++++++++++++++++++++-------- 1 file changed, 43 insertions(+), 17 deletions(-) diff --git a/docs/explanations/message-brokers.md b/docs/explanations/message-brokers.md index 0b96e6b36..573862037 100644 --- a/docs/explanations/message-brokers.md +++ b/docs/explanations/message-brokers.md @@ -1,11 +1,11 @@ # Message Brokers -Blueapi uses a message broker to communicate with downstream services such as the Nexus Filewriter, and tools such as the Blueapi CLI client. Blueapi is broker agnostic, communicating in STOMP via the [bluesky-stomp](https://github.com/DiamondLightSource/bluesky-stomp) library. However, RabbitMQ is most commonly used in DLS deployments. +Blueapi uses a message broker to communicate with certain downstream services such as the Nexus Filewriter, and tools such as the Blueapi CLI client. Blueapi is relatively broker agnostic, communicating via STOMP using the [bluesky-stomp](https://github.com/DiamondLightSource/bluesky-stomp) library. However, RabbitMQ is most commonly used in DLS deployments. ## Messages -When a plan is run, Blueapi broadcasts all worker events, progress events and data events to the configured message broker. Some example messages can be seen below. These are the messages sent to RabbitMQ during a `sleep` plan (printed with an `[x]` at the start): -``` +When a plan is run, Blueapi broadcasts all worker events, progress events and data events to the configured message broker. The example below are messages broadcast during a `sleep` plan: +``` sh [x] public.worker.event:b'{"state":"RUNNING","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":{"outcome":"success","result":null,"type":"NoneType"},"task_complete":true,"task_failed":false},"errors":[],"warnings":[]}' @@ -15,35 +15,61 @@ When a plan is run, Blueapi broadcasts all worker events, progress events and da ``` ## Topology -The topology of the message broker is entirely deployment-specific. Nonetheless, the following details how deployments at DLS are generally structured. +The topology of the message broker is decoupled from Blueapi, and entirely deployment-specific. Nonetheless, the following describe how deployments at DLS are generally structured. +RabbitMQ receives messages from Blueapi via its default topic exchange, `amq.topic`. All messages have the routing key `public.worker.event`, so will be forwarded to all queues bound to `amq.topic` with a compatible routing key (eg. `public.worker.event`, `*.*.event` etc.). No default queues exist, meaning that if no downstream services have registered a queue, all messages will be dropped. -RabbitMQ receives messages from Blueapi via its default topic exchange, `amq.topic`. All messages have the routing key `public.worker.event`, so will be forwarded to all queues bound to `amq.topic` with a compatible routing key (eg. `public.worker.event`, `*.*.event` etc.). +In RabbitMQ, different queues can have the same routing key, and consequently receive a copy of the same messages without interfering with each other. In Blueapi a common pair of queues to have is `public.worker.event.nexus` and `stomp-subscription-{random GUID}`, both bound to the `amq.topic` exchange with routing key `public.worker.event`. The prior is for the Nexus Filewriter, and the latter the Blueapi CLI client. In this case both have access to all messages published by Blueapi, but due to automatically declaring seperate queues at subscription, each consumer's message consumption will not deprive the other. -No default queues exist, meaning that if no downstream services have registered a queue, all messages will be dropped. +## Subscribing to Blueapi Messages -In RabbitMQ, different queues can have the same routing key, and consequently receive a copy of the same messages without interfering with each other. In Blueapi a common pair of queues to have is `public.worker.event.nexus` and `stomp-subscription-{random GUID}`, both bound to the `amq.topic` exchange with routing key `public.worker.event`. The prior is for the Nexus Filewriter, and the latter the Blueapi CLI client. In this case both have access to all messages published by Blueapi, but due to automatically declaring seperate queues at subscription, each consumer's message consumption will not deprive the other. +When Blueapi is configured to use a message broker, any service can subscribe to the broker to receive events generated during plan execution. Due to the dynamic nature of RabbitMQ, new services can subscribe to the exchange at any point, without reconfiguring the server or interfering with other consumer's subscriptions. In general, having services subscribe directly to the message broker is not recommended, as there should be a more appropriate way to access the information you want. +To to connect a consumer to a RabbitMQ topic, follow RabbitMQ's [Tutorial 5: Topics](https://www.rabbitmq.com/tutorials) in your preferred language. Change the exchange name to `amq.topic` and the routing key to `public.worker.event`. For example: -## Blueapi Messaging Mechanism +``` python +#!/usr/bin/env python +import pika +import sys -When Blueapi starts up, a `StompClient` is instantiated which registers a callback to each event stream (worker, progress and data). When an event occurs, it will be sent to the broker's default topic exchange with the routing key `public.worker.event`. In RabbitMQ, the default topic exchange is `amq.topic`. +connection = pika.BlockingConnection( + pika.ConnectionParameters(host='localhost')) +channel = connection.channel() +channel.exchange_declare(exchange='amq.topic', exchange_type='topic', durable=True) -## Subscribing to Blueapi Messages +result = channel.queue_declare('', exclusive=True) +queue_name = result.method.queue + +binding_keys = sys.argv[1:] +if not binding_keys: + sys.stderr.write("Usage: %s [binding_key]...\n" % sys.argv[0]) + sys.exit(1) + +for binding_key in binding_keys: + channel.queue_bind( + exchange='amq.topic', queue=queue_name, routing_key=binding_key) -When Blueapi is configured to use a message broker, services can subscribe to the broker to receive events generated during plan execution. In general, having services subscribe directly to the message broker is not recommended, as there should be a more appropriate way to access the information you want. +print(' [*] Waiting for logs. To exit press CTRL+C') -Due to the dynamic nature of RabbitMQ, new services can subscribe to the exchange at any point, without reconfiguring the server or interfering with other consumer's subscriptions. -To to connect a consumer to a RabbitMQ topic, follow RabbitMQ's [Tutorial 5: Topics](https://www.rabbitmq.com/tutorials) in your preferred language. Change the exchange name to `amq.topic` and the routing key to `public.worker.event`. +def callback(ch, method, properties, body): + print(f" [x] {method.routing_key}:{body}") + + +channel.basic_consume( + queue=queue_name, on_message_callback=callback, auto_ack=True) + +channel.start_consuming() +``` + +With this script and a local instance of RabbitMQ (see [Run BlueAPI and connect to services locally](../how-to/just-run-blueapi-and-services-locally.md)), the command `python ./client.py public.worker.event` will create a new queue bound to `amq.topic` with routing key `public.worker.event`. When a plan is run, this queue woud be populated with messages, which would in turn be consumed by and printed to console. + The consumer created in this tutorial will capture all messages generated during its runtime. It will not have access to messages that were generated previously. If this is a requirement for your use-case, RabbitMQ can be reconfigured to generate specific queues at startup. This can be achieved through [schema definitions](https://www.rabbitmq.com/docs/definitions). ## Tips -It is important to remember that while all queues are guaranteed to receive the same set of messages, there is no guarantee that each queue has consumed their messages. - -For example, a data analysis service listening for a Stop document may process events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready for analysis. +It is important to remember that while all queues are guaranteed to receive the same set of messages in the same order, there is no guarantee that each queue has consumed their messages. -The only guarantee we make is that all queues will receive all events in the correct order. +For example, a service listening for a Stop document in order to kick off data analysis may consume events faster than a file writing service, which needs to write each event to disk. Receiving a Stop document would then only guarantee that Blueapi has completed the plan, not that the data is written to disk and ready to be used. From 75b05c27b46ac1ba015e48fa6f43a732c6f56fae Mon Sep 17 00:00:00 2001 From: Daniel Fernandes Date: Tue, 22 Sep 2026 12:52:01 +0000 Subject: [PATCH 5/5] Add events docs link --- docs/explanations/message-brokers.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/explanations/message-brokers.md b/docs/explanations/message-brokers.md index 573862037..e8f0b0a5e 100644 --- a/docs/explanations/message-brokers.md +++ b/docs/explanations/message-brokers.md @@ -4,7 +4,7 @@ Blueapi uses a message broker to communicate with certain downstream services su ## Messages -When a plan is run, Blueapi broadcasts all worker events, progress events and data events to the configured message broker. The example below are messages broadcast during a `sleep` plan: +When a plan is run, Blueapi broadcasts all worker events, progress events and data events to the configured message broker (see [Events Emitted by the Worker](events.md)). The example below are messages broadcast during a `sleep` plan: ``` sh [x] public.worker.event:b'{"state":"RUNNING","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}' [x] public.worker.event:b'{"state":"IDLE","task_status":{"task_id":"ebef36e4-47bf-4145-855d-15e343a26424","result":null,"task_complete":false,"task_failed":false},"errors":[],"warnings":[]}'