.. _messaging: Queue support ============= Besides its :ref:`restful`, openFHIR Engine can be driven by a message broker: it subscribes to queues carrying FHIR payloads and/or openEHR compositions, translates every message and publishes the result to another queue. An integration that already moves its data over a broker does not need a REST-calling consumer in front of the engine. .. note:: This is an enterprise-only feature and needs the ``messaging`` license option. Enabling it with a license that does not carry the option stops the engine at startup with a message saying so. Licenses issued before this feature existed do not carry it and have to be re-issued — contact us at info@open-fhir.com. `Apache Kafka `_ is the supported broker. How it works ------------ For every message of a subscribed topic the engine 1. decides whether the message is FHIR or openEHR (see `Payload detection`_), 2. translates it under the subscription's tenant, with the same mappers, templates and terminology the REST operations use — a FHIR message is mapped as ``$toopenehr`` maps it, an openEHR message as ``$tofhir`` does, 3. publishes the result to the subscription's target topic, 4. and only then commits the message as consumed. A message that cannot be translated is published, unchanged, to the subscription's dead-letter topic together with an ``OperationOutcome`` saying why. See `Failures and delivery guarantees`_. Configuration ------------- Queue support is disabled by default. It is configured under ``openfhir.messaging``: .. code-block:: yaml openfhir: messaging: enabled: true kafka: bootstrap-servers: localhost:9092 properties: {} # shared client properties (security.protocol, sasl.*, ssl.*) consumer: {} # consumer-only overrides producer: {} # producer-only overrides (e.g. compression.type, max.request.size) retry: max-attempts: 5 initial-interval-ms: 1000 multiplier: 2 max-interval-ms: 30000 subscriptions: - id: mixed tenant: "123" source: openfhir.in payload: auto # auto | fhir | openehr target: openfhir.out # or target-fhir / target-openehr error-target: openfhir.dlq template-id: # optional default, the header wins openehr-format: canonical # canonical | flat provenance: true concurrency: 1 group-id: openfhir-mixed # default openfhir- .. list-table:: :header-rows: 1 :widths: 40 15 45 * - Environment variable - Default - Description * - ``OPENFHIR_MESSAGING_ENABLED`` - ``false`` - Enables queue support. When ``false``, no broker connection is ever attempted. * - ``OPENFHIR_MESSAGING_KAFKA_BOOTSTRAP_SERVERS`` - ``localhost:9092`` - Kafka bootstrap servers, ``host:port[,host:port]``. * - ``OPENFHIR_MESSAGING_KAFKA_PROPERTIES_*`` - — - Any Kafka client property shared by the consumers and the producer, e.g. ``OPENFHIR_MESSAGING_KAFKA_PROPERTIES_SECURITY_PROTOCOL`` for ``security.protocol``. * - ``OPENFHIR_MESSAGING_KAFKA_CONSUMER_*`` - — - Consumer-only Kafka properties, e.g. ``OPENFHIR_MESSAGING_KAFKA_CONSUMER_MAX_POLL_RECORDS``. * - ``OPENFHIR_MESSAGING_KAFKA_PRODUCER_*`` - — - Producer-only Kafka properties, e.g. ``OPENFHIR_MESSAGING_KAFKA_PRODUCER_COMPRESSION_TYPE``. * - ``OPENFHIR_MESSAGING_RETRY_MAX_ATTEMPTS`` - ``5`` - How many times a message is delivered before it is dead-lettered, the first delivery included. Applies to transient failures only. * - ``OPENFHIR_MESSAGING_RETRY_INITIAL_INTERVAL_MS`` - ``1000`` - Wait before the first retry. * - ``OPENFHIR_MESSAGING_RETRY_MULTIPLIER`` - ``2`` - Factor the wait grows by with every retry. * - ``OPENFHIR_MESSAGING_RETRY_MAX_INTERVAL_MS`` - ``30000`` - Upper bound of the wait. Must stay below the consumer's ``max.poll.interval.ms`` (5 minutes unless overridden). Subscriptions are a list; as environment variables they are indexed, ``OPENFHIR_MESSAGING_SUBSCRIPTIONS_0_ID``, ``OPENFHIR_MESSAGING_SUBSCRIPTIONS_0_SOURCE`` and so on: .. list-table:: :header-rows: 1 :widths: 25 15 60 * - Subscription key - Default - Description * - ``id`` - — - Required and unique. Names the subscription in the log and in ``GET /status``. * - ``tenant`` - — - Required. The tenant every message of this subscription is translated under (see :doc:`multitenancy`). A producer cannot choose or override it. * - ``source`` - — - Required. The topic to consume. * - ``payload`` - ``auto`` - What the source topic carries: ``auto`` (FHIR and openEHR mixed, detected per message), ``fhir`` or ``openehr``. * - ``target`` - — - Topic the results are published to, whatever their kind. * - ``target-fhir`` / ``target-openehr`` - — - Topic for the produced FHIR Bundles / openEHR compositions; each wins over ``target``. They are named after what is **published**: an ``openehr`` subscription produces FHIR and therefore needs ``target-fhir`` (or ``target``). * - ``error-target`` - — - Dead-letter topic. When unset, a message that cannot be translated is logged (without its payload) and dropped; a warning at startup points this out. * - ``template-id`` - — - Template id used for messages that carry no ``openfhir-template-id`` header. Leave it unset on a topic that carries more than one template. * - ``openehr-format`` - ``canonical`` - Format of the produced compositions, ``canonical`` or ``flat``, unless the message's ``openfhir-format`` header says otherwise. * - ``provenance`` - ``true`` - Whether produced Bundles carry the engine-generated ``Provenance`` entry, as ``$tofhir`` responses do. The ``openfhir-who`` / ``openfhir-on-behalf-of`` headers only have an effect through it. * - ``concurrency`` - ``1`` - Consumer threads of this subscription on one node. Effective parallelism is ``min(concurrency, partitions)``. * - ``group-id`` - ``openfhir-`` - Kafka consumer group. All nodes of a cluster share it. One mixed queue or one queue per kind ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ A topic that carries both FHIR and openEHR is one subscription with ``payload: auto``, as in the example above. One topic per kind is two subscriptions: .. code-block:: yaml subscriptions: - id: compositions tenant: "123" source: ehr.compositions payload: openehr target-fhir: fhir.bundles error-target: openfhir.dlq - id: bundles tenant: "123" source: fhir.incoming payload: fhir target-openehr: ehr.incoming error-target: openfhir.dlq The configuration is validated at startup and the engine does not start when it is invalid: subscription ids must be unique, every direction a subscription can produce needs a target, and **no subscription may publish to a topic that any subscription consumes** — the engine would otherwise consume its own output, and with ``payload: auto`` translate it back and forth forever. Securing the connection ^^^^^^^^^^^^^^^^^^^^^^^ Everything the Kafka clients understand can be set under ``kafka.properties``, for example SASL over TLS: .. code-block:: yaml openfhir: messaging: kafka: bootstrap-servers: broker-1:9093,broker-2:9093 properties: security.protocol: SASL_SSL sasl.mechanism: SCRAM-SHA-512 sasl.jaas.config: org.apache.kafka.common.security.scram.ScramLoginModule required username="openfhir" password="..."; ssl.truststore.location: /app/truststore.jks ssl.truststore.password: "..." ``GET /status`` masks these values (see :ref:`restful`). ``spring.kafka.*`` properties are **not** read; ``openfhir.messaging.kafka`` is the only place Kafka is configured. A few client settings guard the delivery guarantee and cannot be overridden: the consumers never auto-commit, and the producer always runs with ``acks=all`` and idempotence enabled. Two consumer defaults differ from Kafka's own and can be overridden under ``kafka.consumer``: ``auto.offset.reset`` is ``earliest`` (a new subscription starts at the beginning of its topic instead of skipping what is already there) and ``max.poll.records`` is ``10``. Message contract ---------------- The body of a message is the payload itself — a FHIR Bundle (or single resource) or an openEHR composition (canonical or flat) as JSON. There is no envelope: everything else travels in headers. Inbound headers ^^^^^^^^^^^^^^^ All of them are optional. .. list-table:: :header-rows: 1 :widths: 30 70 * - Header - Effect * - ``openfhir-payload-type`` - ``fhir`` or ``openehr``. Declares what the body is, for a body that does not say so itself (see `Payload detection`_). * - ``openfhir-template-id`` - Template id to map with. Required for a flat composition, as it is for ``$tofhir``, unless the subscription has a ``template-id``. * - ``openfhir-format`` - ``canonical`` or ``flat``: the format of the composition produced from a FHIR message. * - ``openfhir-request-id`` - Correlation id, echoed on the result and recorded on the mapping insight. Defaults to ``:-@``. * - ``openfhir-ehr-id``, ``openfhir-patient``, ``openfhir-who``, ``openfhir-on-behalf-of`` - The call context of ``$tofhir`` (``ehr_id``, ``patient``, ``who``, ``onBehalfOf``): used to fill empty subject references and the ``Provenance`` of the produced Bundle. Outbound headers ^^^^^^^^^^^^^^^^ .. list-table:: :header-rows: 1 :widths: 30 70 * - Header - Meaning * - ``openfhir-direction`` - ``tofhir`` or ``toopenehr``. Absent on a dead-letter record whose kind could not be determined. * - ``openfhir-status`` - On a target topic ``ok``, or ``partial`` when at least one mapping failed and the body is the result of the remaining ones. On the dead-letter topic ``error``. * - ``openfhir-outcome`` - An ``OperationOutcome`` (JSON) with the issues reported during mapping, or the reason a message was dead-lettered. Capped at 16 KB: errors come first, and a last issue says how many were left out. * - ``openfhir-issue-count`` - The total number of issues, also when ``openfhir-outcome`` was capped. * - ``openfhir-request-id`` - The request id, echoed or generated. * - ``openfhir-template-id`` - The template the message was mapped with. * - ``openfhir-source-topic``, ``openfhir-source-partition``, ``openfhir-source-offset`` - Where the translated message came from. Together they identify it uniquely — use them to detect duplicates. The record key of a message is preserved on its result and on its dead-letter record, so related messages keep landing on the same partition. The body of a result is the bare FHIR Bundle or the bare composition. Unlike the REST operations, issues are not put into the body (no ``OperationOutcome`` Bundle entry, no ``Parameters`` wrapper), so a consumer can use it as it is. How much ``openfhir-outcome`` carries follows ``OPENFHIR_OPERATIONS_OUTCOME_VERBOSITY`` (see :ref:`restful`); with ``none`` neither it nor ``openfhir-issue-count`` is set, while ``openfhir-status`` still tells ``partial`` from ``ok``. A dead-letter record carries the original body, key and headers unchanged — it can be put back on the source topic as it is — with the outbound headers added. Payload detection ^^^^^^^^^^^^^^^^^ A body with a top-level ``resourceType`` is FHIR. A body with a top-level ``_type`` (canonical) or with flat paths as member names is openEHR. For an array, its first element decides. The body is looked at in every mode. A subscription's ``payload`` setting and the ``openfhir-payload-type`` header say what is *expected*; a body that contradicts either of them — a FHIR Bundle on a ``payload: openehr`` subscription, say — is dead-lettered rather than mapped as something it is not. A body that is JSON but says nothing about itself needs the header or a non-``auto`` subscription. A body that is not JSON is dead-lettered. Failures and delivery guarantees -------------------------------- .. list-table:: :header-rows: 1 :widths: 35 65 * - What happened - What the engine does * - A single mapping failed, or elements were skipped - The result is published with ``openfhir-status: partial`` (or ``ok`` for warnings only) and the issues in ``openfhir-outcome`` — exactly what the REST operations answer with ``200`` and an ``OperationOutcome``. * - The message is at fault: not JSON, unknown template, no context mapper, missing template id, … - Dead-lettered straight away, without retrying. These are the failures REST answers with a ``4xx``. * - Something the engine depends on is unavailable: database, CDR, the broker itself - Retried with exponential backoff (``retry.*``). If it still fails after ``max-attempts`` deliveries, the message is dead-lettered. * - The dead-letter topic cannot be written - The message is retried until it can: the partition waits, nothing is lost. Only a record the dead-letter topic rejects for its size is logged and skipped. In every case the subscription keeps consuming afterwards. **Delivery is at-least-once.** A message is committed only after its result (or its dead-letter record) was acknowledged by the broker, so a crash or a rebalance in between means the message is translated and published again. Produced Bundles contain generated ids, so two results of the same message are not byte-identical: consumers that need exactly-once semantics should dedupe on the ``openfhir-source-topic`` / ``-partition`` / ``-offset`` headers. Order is preserved per partition, and therefore per record key. .. note:: A failure of infrastructure *inside* a single mapping — a terminology server that does not answer, for example — surfaces as an ``error`` issue on a ``partial`` result, as it does over REST. It is not retried. Startup, shutdown and clustering -------------------------------- Subscriptions start once the engine is fully up: after the license was validated and the :ref:`bootstrap scan ` ran. A broker that cannot be reached never keeps the engine from starting or from serving REST: the subscription is reported as ``FAILED`` and started again every 30 seconds until it succeeds. A source topic that does not exist yet is waited for. On shutdown the subscriptions are stopped first and the producer is closed after them, so a message in flight can still publish its result. In a :doc:`multi-node deployment ` every node runs the same subscriptions in the same consumer group; Kafka splits the partitions of a topic between the nodes and moves them when a node leaves. No further coordination is needed. A rebalance can deliver a message a second time (see above). Monitoring ---------- ``GET /status`` reports each subscription's state and what happened to the messages this node took since it started: .. code-block:: json { "messaging": { "enabled": true, "subscriptions": [ { "id": "mixed", "tenant": "123", "source": "openfhir.in", "payload": "auto", "state": "RUNNING", "processed": 1284, "partial": 3, "deadLettered": 1 } ] } } ``processed`` counts messages published to a target topic, ``partial`` those among them that carried a mapping error, ``deadLettered`` those that could not be translated (including the ones dropped for want of an ``error-target``). ``state`` is one of ``STARTING``, ``RUNNING``, ``FAILED`` and ``STOPPED``; ``RUNNING`` means the consumer is running, not that the broker currently answers — a consumer that lost its broker keeps reconnecting in the background. Every message is also recorded as a :doc:`mapping insight `, with its request id, direction and — for a dead-lettered message — the error. Sizing and limits ----------------- * Results are bounded by the message size limits of the broker and the producer, 1 MB by default. FHIR Bundles are larger than the compositions they are produced from: consider ``kafka.producer.compression.type: gzip`` and raise ``max.request.size`` (and the topic's ``max.message.bytes``) where needed. A result that is still too large is dead-lettered with an ``OperationOutcome`` saying so. * One mapping insight is written per message. On a high-volume topic consider ``OPENFHIR_INSIGHTS_ENABLED=false``, or at least leave payload storage off. * Retries wait on the consumer thread, so a message that is being retried holds back the messages behind it on the same partition; other partitions are unaffected.