diff --git a/.github/workflows/cdb-base-mqtt.yml b/.github/workflows/cdb-base-mqtt.yml new file mode 100644 index 000000000..f05712045 --- /dev/null +++ b/.github/workflows/cdb-base-mqtt.yml @@ -0,0 +1,45 @@ +name: CDB Base MQTT + +on: + push: + paths: + - "tools/developer_tools/cdb-base-mqtt/**" + - ".github/workflows/cdb-base-mqtt.yml" + pull_request: + paths: + - "tools/developer_tools/cdb-base-mqtt/**" + - ".github/workflows/cdb-base-mqtt.yml" + workflow_dispatch: + +permissions: + contents: read + +jobs: + test: + runs-on: ubuntu-latest + defaults: + run: + working-directory: tools/developer_tools/cdb-base-mqtt + + steps: + - name: Check out repository + uses: actions/checkout@v4 + + - name: Set up Java 8 and Maven cache + uses: actions/setup-java@v4 + with: + distribution: temurin + java-version: "8" + cache: maven + cache-dependency-path: tools/developer_tools/cdb-base-mqtt/pom.xml + + - name: Run tests + run: mvn --batch-mode --no-transfer-progress verify + + - name: Publish test results + if: always() + uses: actions/upload-artifact@v4 + with: + name: cdb-base-mqtt-test-results + path: tools/developer_tools/cdb-base-mqtt/target/surefire-reports/ + if-no-files-found: warn diff --git a/Makefile b/Makefile index f2daa3b76..33d369d64 100644 --- a/Makefile +++ b/Makefile @@ -9,6 +9,7 @@ TOP = . SUBDIRS = src .PHONY: support support-mysql dev-config prepare-release release-python-client +.PHONY: mqtt-configuration configure-mqtt mqtt-configuration-dev configure-mqtt-dev .PHONY: db backup db-dev deploy-web-portal undeploy-web-portal deploy-web-service undeploy-web-service .PHONY: db-dev backup-dev deploy-web-portal-dev undeploy-web-portal-dev deploy-web-service-dev undeploy-web-service-dev @@ -28,6 +29,12 @@ release-python-client: configuration: $(TOP)/sbin/cdb_create_configuration.sh +mqtt-configuration: + $(TOP)/sbin/cdb_create_mqtt_configuration.sh cdb + +configure-mqtt: + $(TOP)/sbin/cdb_configure_mqtt_service.sh cdb + support: $(TOP)/sbin/cdb_install_support.sh @@ -79,6 +86,12 @@ undeploy-web-service: configuration-dev: $(TOP)/sbin/cdb_create_configuration.sh cdb_dev +mqtt-configuration-dev: + $(TOP)/sbin/cdb_create_mqtt_configuration.sh cdb_dev + +configure-mqtt-dev: + $(TOP)/sbin/cdb_configure_mqtt_service.sh cdb_dev + db-dev: $(TOP)/sbin/cdb_create_db.sh cdb_dev diff --git a/sbin/cdb_configure_mqtt_service.sh b/sbin/cdb_configure_mqtt_service.sh new file mode 100755 index 000000000..020298bd0 --- /dev/null +++ b/sbin/cdb_configure_mqtt_service.sh @@ -0,0 +1,170 @@ +#!/bin/bash + +# Copyright (c) UChicago Argonne, LLC. All rights reserved. +# See LICENSE file. + + +# +# Script used for configuring MQTT connector for CDB webapp +# Deploys MQTT Resource Adapter and creates connection pool/resource +# +# Usage: +# +# $0 [CDB_DB_NAME] [mqtt_config_file] +# +# If no config file is specified, defaults to $CDB_INSTALL_DIR/etc/.mqtt.conf +# +# Sample MQTT configuration file contents (mqtt.conf): +# MQTT_HOST=localhost # MQTT broker hostname (default: localhost) +# MQTT_PORT=1883 # MQTT broker port (default: 1883) +# MQTT_USERNAME=admin # MQTT username (optional) +# MQTT_PASSWORD=admin # MQTT password (optional) +# MQTT_CLEAN_SESSION=true # Clean session flag (optional) +# MQTT_QOS=1 # Quality of Service level (optional) +# MQTT_KEEP_ALIVE_INTERVAL=60 # Keep alive interval in seconds (optional) +# MQTT_CONNECTION_TIMEOUT=30 # Connection timeout in seconds (optional) +# MQTT_MAX_INFLIGHT=10 # Maximum number of messages in flight (optional) +# MQTT_AUTOMATIC_RECONNECT=true # Automatic reconnection on disconnect (optional) +# MQTT_FILE_PERSISTANCE=false # Enable file-based message persistence (optional) +# MQTT_PERSISTENCE_DIRECTORY=. # Directory for persistent message storage (optional) +# MQTT_TOPIC_FILTER # MQTT topic filter for subscriptions (optional) + +MY_DIR=`dirname $0` && cd $MY_DIR && MY_DIR=`pwd` +if [ -z "${CDB_ROOT_DIR}" ]; then + CDB_ROOT_DIR=$MY_DIR/.. +fi +CDB_ENV_FILE=${CDB_ROOT_DIR}/setup.sh +if [ ! -f ${CDB_ENV_FILE} ]; then + echo "Environment file ${CDB_ENV_FILE} does not exist." + exit 2 +fi +. ${CDB_ENV_FILE} > /dev/null + +# Constants +MQTT_POOL_NAME="cdb/MQTT/pool" +MQTT_RESOURCE_NAME="cdb/MQTT/resource" +MQTT_RAR_NAME="mqtt-rar-0.8.0" +MQTT_CONNECTOR_RAR_DEPLOYMENT_NAME="cdb-mqtt-rar-deployment" +MQTT_RAR_PATH=$CDB_ROOT_DIR/src/lib/${MQTT_RAR_NAME}.rar + +# Look for the deployment-specific MQTT configuration file +CDB_DB_NAME=${1:-${CDB_DB_NAME:-cdb}} +if [ ! -z "$2" ]; then + mqttConfigFile=$2 +else + mqttConfigFile=$CDB_INSTALL_DIR/etc/${CDB_DB_NAME}.mqtt.conf +fi + +if [ -f $mqttConfigFile ]; then + echo "Using MQTT config file: $mqttConfigFile" + . $mqttConfigFile +else + echo "Error: MQTT config file $mqttConfigFile not found." + echo "You can create one using cdb_create_mqtt_configuration.sh ${CDB_DB_NAME}" + exit 1 +fi + +CDB_HOST_ARCH=$(uname -sm | tr -s '[:upper:][:blank:]' '[:lower:][\-]') +GLASSFISH_DIR=$CDB_SUPPORT_DIR/payara/$CDB_HOST_ARCH + +ASADMIN_CMD=$GLASSFISH_DIR/bin/asadmin + +# MQTT Configuration defaults +MQTT_HOST=${MQTT_HOST:=localhost} +MQTT_PORT=${MQTT_PORT:=1883} + +# Build properties string for connection pool +PROPERTIES="serverURIs=tcp\\://${MQTT_HOST}\\:${MQTT_PORT}" + +if [ ! -z "$MQTT_USERNAME" ]; then + PROPERTIES="${PROPERTIES}:userName=${MQTT_USERNAME}" +fi + +if [ ! -z "$MQTT_PASSWORD" ]; then + PROPERTIES="${PROPERTIES}:password=${MQTT_PASSWORD}" +fi + +if [ ! -z "$MQTT_CLEAN_SESSION" ]; then + PROPERTIES="${PROPERTIES}:cleanSession=${MQTT_CLEAN_SESSION}" +fi + +if [ ! -z "$MQTT_KEEP_ALIVE_INTERVAL" ]; then + PROPERTIES="${PROPERTIES}:keepAliveInterval=${MQTT_KEEP_ALIVE_INTERVAL}" +fi + +if [ ! -z "$MQTT_CONNECTION_TIMEOUT" ]; then + PROPERTIES="${PROPERTIES}:connectionTimeout=${MQTT_CONNECTION_TIMEOUT}" +fi + +if [ ! -z "$MQTT_MAX_INFLIGHT" ]; then + PROPERTIES="${PROPERTIES}:maxInflight=${MQTT_MAX_INFLIGHT}" +fi + +if [ ! -z "$MQTT_AUTOMATIC_RECONNECT" ]; then + PROPERTIES="${PROPERTIES}:automaticReconnect=${MQTT_AUTOMATIC_RECONNECT}" +fi + +if [ ! -z "$MQTT_FILE_PERSISTANCE" ]; then + PROPERTIES="${PROPERTIES}:filePersistance=${MQTT_FILE_PERSISTANCE}" +fi + +if [ ! -z "$MQTT_PERSISTENCE_DIRECTORY" ]; then + PROPERTIES="${PROPERTIES}:persistenceDirectory=${MQTT_PERSISTENCE_DIRECTORY}" +fi + +if [ ! -z "$MQTT_QOS" ]; then + PROPERTIES="${PROPERTIES}:qos=${MQTT_QOS}" +fi + +if [ ! -z "$MQTT_TOPIC_FILTER" ]; then + PROPERTIES="${PROPERTIES}:topicFilter=${MQTT_TOPIC_FILTER}" +fi + +# Deploy MQTT RAR +echo "Deploying MQTT RAR" +if [ -f "$MQTT_RAR_PATH" ]; then + # Check if already deployed and undeploy if needed + $ASADMIN_CMD list-applications | grep -q ${MQTT_CONNECTOR_RAR_DEPLOYMENT_NAME} && { + echo "Undeploying existing MQTT RAR" + # Check if resource exists and delete it + $ASADMIN_CMD list-connector-resources | grep -q ${MQTT_RESOURCE_NAME} && { + echo "Deleting existing MQTT resource" + $ASADMIN_CMD delete-connector-resource ${MQTT_RESOURCE_NAME} || exit 1 + } + + # Check if connection pool exists and delete it + $ASADMIN_CMD list-connector-connection-pools | grep -q ${MQTT_POOL_NAME} && { + echo "Deleting existing MQTT connection pool" + $ASADMIN_CMD delete-connector-connection-pool ${MQTT_POOL_NAME} || exit 1 + } + + echo "Undeploying existing MQTT RAR" + $ASADMIN_CMD undeploy ${MQTT_CONNECTOR_RAR_DEPLOYMENT_NAME} || exit 1 + } + $ASADMIN_CMD deploy --name ${MQTT_CONNECTOR_RAR_DEPLOYMENT_NAME} $MQTT_RAR_PATH || exit 1 +else + echo "Warning: MQTT RAR file not found at $MQTT_RAR_PATH" + exit 1 +fi + +# Create MQTT connection pool +echo "Creating MQTT connection pool ${MQTT_POOL_NAME}" +$ASADMIN_CMD create-connector-connection-pool \ + --raname ${MQTT_CONNECTOR_RAR_DEPLOYMENT_NAME} \ + --connectiondefinition fish.payara.cloud.connectors.mqtt.api.MQTTConnectionFactory \ + --property "${PROPERTIES}" \ + ${MQTT_POOL_NAME} || exit 1 +# Create MQTT resource +echo "Creating MQTT resource ${MQTT_RESOURCE_NAME}" +$ASADMIN_CMD create-connector-resource \ + --poolname ${MQTT_POOL_NAME} \ + ${MQTT_RESOURCE_NAME} || exit 1 + +# The connector logs allocation failures with a full stack trace before returning null. +$ASADMIN_CMD set-log-levels fish.payara.cloud.connectors.mqtt.api.outbound.MQTTConnectionFactoryImpl=OFF || exit 1 + +# Test MQTT connection pool +echo "Testing MQTT connection pool" +$ASADMIN_CMD ping-connection-pool ${MQTT_POOL_NAME} || { echo "Warning: MQTT connection pool ping failed"; exit 1; } + +echo "Restart or redeploy CDB." \ No newline at end of file diff --git a/sbin/cdb_create_mqtt_configuration.sh b/sbin/cdb_create_mqtt_configuration.sh new file mode 100755 index 000000000..000b0f4fb --- /dev/null +++ b/sbin/cdb_create_mqtt_configuration.sh @@ -0,0 +1,135 @@ +#!/bin/bash + +# Script to create MQTT configuration file for CDB service +# Usage: $0 [CDB_DB_NAME] + +MY_DIR=`dirname $0` && cd $MY_DIR && MY_DIR=`pwd` +if [ -z "${CDB_ROOT_DIR}" ]; then + CDB_ROOT_DIR=$MY_DIR/.. +fi +CDB_ENV_FILE=${CDB_ROOT_DIR}/setup.sh +if [ ! -f ${CDB_ENV_FILE} ]; then + echo "Environment file ${CDB_ENV_FILE} does not exist." + exit 2 +fi +. ${CDB_ENV_FILE} > /dev/null + +# Default configuration file location +CDB_DB_NAME=${1:-${CDB_DB_NAME:-cdb}} +MQTT_CONFIG_FILE=${CDB_INSTALL_DIR}/etc/${CDB_DB_NAME}.mqtt.conf + +echo "===================================" +echo "MQTT Configuration Setup for CDB" +echo "===================================" +echo "" +echo "This script will help you create an MQTT configuration file." +echo "Configuration will be saved to: $MQTT_CONFIG_FILE" +echo "" + +# Check if config file already exists +if [ -f "$MQTT_CONFIG_FILE" ]; then + read -p "Configuration file already exists. Overwrite? (y/n): " -n 1 -r + echo + if [[ ! $REPLY =~ ^[Yy]$ ]]; then + echo "Exiting without changes." + exit 0 + fi +fi + +echo "" +echo "Configuration Details:" +echo "- cleanSession: Whether client and server should remember state across reconnects" +echo "- automaticReconnect: Whether client will automatically reconnect if connection is lost" +echo "- filePersistance: Whether to use file persistence for un-acknowledged messages" +echo "- persistenceDirectory: Directory to use for file persistence" +echo "- connectionTimeout: Connection timeout value in seconds" +echo "- maxInflight: Maximum messages that can be sent without acknowledgements" +echo "- keepAliveInterval: Keep alive interval in seconds" +echo "- userName/password: Authentication credentials" +# Disable MDB only variables. +# echo "- topicFilter: Topic Filter (For MDBs only)" +# echo "- qos: Quality of Service for the subscription (For MDBs only)" +echo "" + +# Ensure directory exists +mkdir -p $(dirname "$MQTT_CONFIG_FILE") + +# Prompt for configuration values +read -p "MQTT Host [localhost]: " MQTT_HOST +MQTT_HOST=${MQTT_HOST:-localhost} + +read -p "MQTT Port [1883]: " MQTT_PORT +MQTT_PORT=${MQTT_PORT:-1883} + +read -p "MQTT Username (leave empty for no auth): " MQTT_USERNAME + +if [ ! -z "$MQTT_USERNAME" ]; then + read -s -p "MQTT Password: " MQTT_PASSWORD + echo +fi + +read -p "Clean Session (true/false) [false]: " CLEAN_SESSION +CLEAN_SESSION=${CLEAN_SESSION:-false} + +read -p "Automatic Reconnect (true/false) [true]: " AUTOMATIC_RECONNECT +AUTOMATIC_RECONNECT=${AUTOMATIC_RECONNECT:-true} + +read -p "File Persistance (true/false) [false]: " FILE_PERSISTANCE +FILE_PERSISTANCE=${FILE_PERSISTANCE:-false} + +read -p "Persistence Directory [.]: " PERSISTENCE_DIRECTORY +PERSISTENCE_DIRECTORY=${PERSISTENCE_DIRECTORY:-.} + +read -p "Connection Timeout (seconds) [30]: " CONNECTION_TIMEOUT +CONNECTION_TIMEOUT=${CONNECTION_TIMEOUT:-30} + +read -p "Max Inflight [10]: " MAX_INFLIGHT +MAX_INFLIGHT=${MAX_INFLIGHT:-10} + +read -p "Keep Alive Interval (seconds) [60]: " KEEP_ALIVE_INTERVAL +KEEP_ALIVE_INTERVAL=${KEEP_ALIVE_INTERVAL:-60} + +# MDB only variables +# read -p "Topic Filter (leave empty if not using MDB): " TOPIC_FILTER +# read -p "QoS (0/1/2) [0]: " QOS +# QOS=${QOS:-0} + +# Write configuration file +cat > "$MQTT_CONFIG_FILE" << EOF +# MQTT Configuration for CDB Service +# Generated on $(date) + +MQTT_HOST=$MQTT_HOST +MQTT_PORT=$MQTT_PORT +EOF + +if [ ! -z "$MQTT_USERNAME" ]; then + echo "MQTT_USERNAME=$MQTT_USERNAME" >> "$MQTT_CONFIG_FILE" +fi + +if [ ! -z "$MQTT_PASSWORD" ]; then + echo "MQTT_PASSWORD=$MQTT_PASSWORD" >> "$MQTT_CONFIG_FILE" +fi + +cat >> "$MQTT_CONFIG_FILE" << EOF +MQTT_CLEAN_SESSION=$CLEAN_SESSION +MQTT_AUTOMATIC_RECONNECT=$AUTOMATIC_RECONNECT +MQTT_FILE_PERSISTANCE=$FILE_PERSISTANCE +MQTT_PERSISTENCE_DIRECTORY=$PERSISTENCE_DIRECTORY +MQTT_CONNECTION_TIMEOUT=$CONNECTION_TIMEOUT +MQTT_MAX_INFLIGHT=$MAX_INFLIGHT +MQTT_KEEP_ALIVE_INTERVAL=$KEEP_ALIVE_INTERVAL +EOF + +if [ ! -z "$TOPIC_FILTER" ]; then + echo "MQTT_TOPIC_FILTER=$TOPIC_FILTER" >> "$MQTT_CONFIG_FILE" +fi + +if [ ! -z "$QOS" ]; then + echo "MQTT_QOS=$QOS" >> "$MQTT_CONFIG_FILE" +fi + +echo "" +echo "Configuration file created successfully at: $MQTT_CONFIG_FILE" +echo "" +echo "You can now run cdb_configure_mqtt_service.sh ${CDB_DB_NAME} to apply this configuration." \ No newline at end of file diff --git a/src/java/CdbWebPortal/lib/cdb-base-mqtt-1.0.0.jar b/src/java/CdbWebPortal/lib/cdb-base-mqtt-1.0.0.jar new file mode 100644 index 000000000..04b3e761a Binary files /dev/null and b/src/java/CdbWebPortal/lib/cdb-base-mqtt-1.0.0.jar differ diff --git a/src/java/CdbWebPortal/lib/mqtt-jca-api-1.0.0.jar b/src/java/CdbWebPortal/lib/mqtt-jca-api-1.0.0.jar new file mode 100644 index 000000000..fa423699f Binary files /dev/null and b/src/java/CdbWebPortal/lib/mqtt-jca-api-1.0.0.jar differ diff --git a/src/java/CdbWebPortal/nbproject/build-impl.xml b/src/java/CdbWebPortal/nbproject/build-impl.xml index 2c394668c..5467c0d3e 100644 --- a/src/java/CdbWebPortal/nbproject/build-impl.xml +++ b/src/java/CdbWebPortal/nbproject/build-impl.xml @@ -1041,6 +1041,8 @@ exists or setup the property manually. For example like this: + + @@ -1104,6 +1106,8 @@ exists or setup the property manually. For example like this: + + diff --git a/src/java/CdbWebPortal/nbproject/genfiles.properties b/src/java/CdbWebPortal/nbproject/genfiles.properties index f3a7e4c63..c1e02b10e 100644 --- a/src/java/CdbWebPortal/nbproject/genfiles.properties +++ b/src/java/CdbWebPortal/nbproject/genfiles.properties @@ -1,10 +1,10 @@ -build.xml.data.CRC32=d5de4743 +build.xml.data.CRC32=9c7b4582 build.xml.script.CRC32=67492cbd build.xml.stylesheet.CRC32=1707db4f@1.94.0.1 # This file is used by a NetBeans-based IDE to track changes in generated files such as build-impl.xml. # Do not edit this file. You may delete it but then the IDE will never regenerate such files for you. -nbproject/build-impl.xml.data.CRC32=d5de4743 -nbproject/build-impl.xml.script.CRC32=8f60c160 +nbproject/build-impl.xml.data.CRC32=9c7b4582 +nbproject/build-impl.xml.script.CRC32=e1cf7847 nbproject/build-impl.xml.stylesheet.CRC32=334708a0@1.94.0.1 nbproject/rest-build.xml.data.CRC32=83e13ee9 nbproject/rest-build.xml.script.CRC32=0d1fb1b4 diff --git a/src/java/CdbWebPortal/nbproject/project.properties b/src/java/CdbWebPortal/nbproject/project.properties index 4206c7f64..34eea5b9f 100644 --- a/src/java/CdbWebPortal/nbproject/project.properties +++ b/src/java/CdbWebPortal/nbproject/project.properties @@ -42,6 +42,7 @@ file.reference.bcmail-jdk14-1.38.jar=lib/bcmail-jdk14-1.38.jar file.reference.bcprov-jdk14-1.38.jar=lib/bcprov-jdk14-1.38.jar file.reference.bctsp-jdk14-1.38.jar=lib/bctsp-jdk14-1.38.jar file.reference.classgraph-4.8.31.jar=lib/classgraph-4.8.31.jar +file.reference.cdb-base-mqtt-1.0.0.jar=lib/cdb-base-mqtt-1.0.0.jar file.reference.commons-lang3-3.11.jar=lib/commons-lang3-3.11.jar file.reference.commons-logging-1.2.jar=lib/commons-logging-1.2.jar file.reference.dm-api-3.3.1.jar=lib/dm-api-3.3.1.jar @@ -79,6 +80,7 @@ file.reference.libphonenumber-8.12.3.jar=lib/libphonenumber-8.12.3.jar file.reference.log4j-api-2.26.0.jar=lib/log4j-api-2.26.0.jar file.reference.log4j-core-2.26.0.jar=lib/log4j-core-2.26.0.jar file.reference.metadata-extractor-2.17.0.jar=lib/metadata-extractor-2.17.0.jar +file.reference.mqtt-jca-api-1.0.0.jar=lib/mqtt-jca-api-1.0.0.jar file.reference.omnifaces-3.10.1.jar-1=lib/omnifaces-3.10.1.jar file.reference.pdfbox-2.0.24.jar=lib/pdfbox-2.0.24.jar file.reference.poi-4.0.1.jar=lib/poi-4.0.1.jar @@ -169,7 +171,9 @@ javac.classpath=\ ${file.reference.flexmark-util-html-0.64.8.jar}:\ ${file.reference.flexmark-util-misc-0.64.8.jar}:\ ${file.reference.flexmark-util-sequence-0.64.8.jar}:\ - ${file.reference.flexmark-util-visitor-0.64.8.jar} + ${file.reference.flexmark-util-visitor-0.64.8.jar}:\ + ${file.reference.cdb-base-mqtt-1.0.0.jar}:\ + ${file.reference.mqtt-jca-api-1.0.0.jar} # Space-separated list of extra javac options javac.compilerargs= javac.debug=true diff --git a/src/java/CdbWebPortal/nbproject/project.xml b/src/java/CdbWebPortal/nbproject/project.xml index fe2ceeb99..553e82ac8 100644 --- a/src/java/CdbWebPortal/nbproject/project.xml +++ b/src/java/CdbWebPortal/nbproject/project.xml @@ -242,6 +242,14 @@ ${file.reference.flexmark-util-visitor-0.64.8.jar} WEB-INF/lib + + ${file.reference.cdb-base-mqtt-1.0.0.jar} + WEB-INF/lib + + + ${file.reference.mqtt-jca-api-1.0.0.jar} + WEB-INF/lib + diff --git a/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbEntityEvent.java b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbEntityEvent.java new file mode 100644 index 000000000..86d216187 --- /dev/null +++ b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbEntityEvent.java @@ -0,0 +1,58 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.mqtt; + +import gov.anl.aps.cdb.portal.model.db.entities.CdbEntity; +import gov.anl.aps.cdb.portal.model.db.entities.UserInfo; +import gov.anl.aps.cdb.base.mqtt.MqttEvent; + +/** Scalar snapshot of a completed entity operation. */ +public class CdbEntityEvent extends MqttEvent { + + private final String topic; + private final String operation; + private final String entityId; + private final String entityType; + private final String entityDescription; + private final String triggeredByUsername; + + public CdbEntityEvent(CdbMqttTopic topic, CdbEntity entity, UserInfo triggeredByUser) { + if (topic == null || entity == null) { + throw new IllegalArgumentException("topic and entity are required"); + } + this.topic = topic.getValue(); + this.operation = topic.name().toLowerCase(); + Object id = entity.getId(); + this.entityId = id == null ? null : id.toString(); + this.entityType = entity.getClass().getSimpleName(); + this.entityDescription = entity.getSystemLogString(); + this.triggeredByUsername = triggeredByUser == null ? null : triggeredByUser.getUsername(); + } + + @Override + public String getTopic() { + return topic; + } + + public String getOperation() { + return operation; + } + + public String getEntityId() { + return entityId; + } + + public String getEntityType() { + return entityType; + } + + public String getEntityDescription() { + return entityDescription; + } + + public String getTriggeredByUsername() { + return triggeredByUsername; + } +} diff --git a/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttPublisher.java b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttPublisher.java new file mode 100644 index 000000000..f2bd5f135 --- /dev/null +++ b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttPublisher.java @@ -0,0 +1,50 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.mqtt; + +import com.fasterxml.jackson.databind.ObjectMapper; +import gov.anl.aps.cdb.base.mqtt.MqttEvent; +import gov.anl.aps.cdb.base.mqtt.MqttPublisher; +import gov.anl.aps.cdb.base.mqtt.exception.MqttException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttNotConfiguredException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttPublishException; +import gov.anl.aps.cdb.base.mqtt.jca.JcaMqttPublisher; +import gov.anl.aps.cdb.base.mqtt.jca.JndiMqttConnectionProvider; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +/** Application adapter that makes MQTT publishing best effort. */ +public final class CdbMqttPublisher { + + public static final String JNDI_RESOURCE_NAME = "cdb/MQTT/resource"; + private static final Logger LOGGER = LogManager.getLogger(CdbMqttPublisher.class.getName()); + private static final MqttPublisher PUBLISHER = new JcaMqttPublisher( + new JndiMqttConnectionProvider(JNDI_RESOURCE_NAME), new ObjectMapper()); + + private CdbMqttPublisher() { + } + + public static void publish(MqttEvent event) { + try { + PUBLISHER.publish(event); + } catch (MqttNotConfiguredException ex) { + LOGGER.warn("MQTT is unavailable; skipped event for " + event.getTopic() + + ". " + ex.getMessage()); + } catch (MqttPublishException ex) { + LOGGER.warn("MQTT broker is unavailable or disconnected; skipped event for " + + event.getTopic() + ". " + getRootCauseMessage(ex)); + } catch (MqttException | RuntimeException ex) { + LOGGER.error("Failed to prepare MQTT event for " + event.getTopic(), ex); + } + } + + private static String getRootCauseMessage(Throwable error) { + Throwable cause = error; + while (cause.getCause() != null && cause.getCause() != cause) { + cause = cause.getCause(); + } + return cause.getMessage() == null ? cause.getClass().getSimpleName() : cause.getMessage(); + } +} diff --git a/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttTopic.java b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttTopic.java new file mode 100644 index 000000000..ac3251519 --- /dev/null +++ b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/mqtt/CdbMqttTopic.java @@ -0,0 +1,21 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.mqtt; + +public enum CdbMqttTopic { + ADD("cdb/entity/add"), + UPDATE("cdb/entity/update"), + DELETE("cdb/entity/delete"); + + private final String value; + + CdbMqttTopic(String value) { + this.value = value; + } + + public String getValue() { + return value; + } +} diff --git a/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/portal/controllers/utilities/CdbEntityControllerUtility.java b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/portal/controllers/utilities/CdbEntityControllerUtility.java index cf7dfff81..9355a1298 100644 --- a/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/portal/controllers/utilities/CdbEntityControllerUtility.java +++ b/src/java/CdbWebPortal/src/java/gov/anl/aps/cdb/portal/controllers/utilities/CdbEntityControllerUtility.java @@ -5,6 +5,9 @@ package gov.anl.aps.cdb.portal.controllers.utilities; import gov.anl.aps.cdb.common.exceptions.CdbException; +import gov.anl.aps.cdb.mqtt.CdbEntityEvent; +import gov.anl.aps.cdb.mqtt.CdbMqttPublisher; +import gov.anl.aps.cdb.mqtt.CdbMqttTopic; import gov.anl.aps.cdb.common.utilities.StringUtility; import gov.anl.aps.cdb.portal.constants.SystemLogLevel; import gov.anl.aps.cdb.portal.model.db.beans.CdbEntityFacade; @@ -55,6 +58,7 @@ public EntityType create(EntityType entity, UserInfo createdByUserInfo) throws C addCreatedSystemLog(entity, createdByUserInfo); entity.setPersitanceErrorMessage(null); + publishMqttEvent(CdbMqttTopic.ADD, entity, createdByUserInfo); clearCaches(); return entity; @@ -79,7 +83,10 @@ public void createList(List entities, UserInfo createdByUserInfo) th } getEntityDbFacade().create(entities); - addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Created " + entities.size() + " entities.", createdByUserInfo); + addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Created " + entities.size() + " entities.", createdByUserInfo); + for (EntityType entity : entities) { + publishMqttEvent(CdbMqttTopic.ADD, entity, createdByUserInfo); + } setPersistenceErrorMessageForList(entities, null); clearCaches(); } catch (CdbException ex) { @@ -103,6 +110,7 @@ public EntityType update(EntityType entity, UserInfo updatedByUserInfo) throws C EntityType updatedEntity = getEntityDbFacade().edit(entity); addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Updated: " + entity.getSystemLogString(), updatedByUserInfo); entity.setPersitanceErrorMessage(null); + publishMqttEvent(CdbMqttTopic.UPDATE, updatedEntity, updatedByUserInfo); clearCaches(); return updatedEntity; @@ -127,6 +135,7 @@ public EntityType updateOnRemoval(EntityType entity, UserInfo updatedByUserInfo) logger.debug("Updating " + getDisplayEntityTypeName() + " " + getEntityInstanceName(entity)); prepareEntityUpdateOnRemoval(entity); EntityType updatedEntity = getEntityDbFacade().edit(entity); + publishMqttEvent(CdbMqttTopic.UPDATE, updatedEntity, updatedByUserInfo); clearCaches(); return updatedEntity; @@ -156,6 +165,7 @@ public void updateList(List entities, UserInfo updatedByUserInfo) th for (EntityType entity : entities) { entity.setPersitanceErrorMessage(null); addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Updated: " + entity.getSystemLogString(), updatedByUserInfo); + publishMqttEvent(CdbMqttTopic.UPDATE, entity, updatedByUserInfo); } clearCaches(); } catch (CdbException ex) { @@ -177,7 +187,8 @@ public void destroy(EntityType entity, UserInfo destroyedByUserInfo) throws CdbE prepareEntityDestroy(entity, destroyedByUserInfo); getEntityDbFacade().remove(entity); - addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Deleted: " + entity.getSystemLogString(), destroyedByUserInfo); + addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Deleted: " + entity.getSystemLogString(), destroyedByUserInfo); + publishMqttEvent(CdbMqttTopic.DELETE, entity, destroyedByUserInfo); clearCaches(); } catch (CdbException ex) { entity.setPersitanceErrorMessage(ex.getMessage()); @@ -214,6 +225,14 @@ public void destroyList( getEntityDbFacade().remove(entities, updateEntity); addCdbEntitySystemLog(SystemLogLevel.entityInfo, "Deleted: " + entities.size() + " entities.", destroyedByUserInfo); + for (EntityType entity : entities) { + if (entity != null) { + publishMqttEvent(CdbMqttTopic.DELETE, entity, destroyedByUserInfo); + } + } + if (updateEntity != null) { + publishMqttEvent(CdbMqttTopic.UPDATE, updateEntity, destroyedByUserInfo); + } setPersistenceErrorMessageForList(entities, null); clearCaches(); } catch (CdbException ex) { @@ -229,6 +248,14 @@ public void destroyList( } } + protected void publishMqttEvent(CdbMqttTopic topic, EntityType entity, UserInfo userInfo) { + try { + CdbMqttPublisher.publish(new CdbEntityEvent(topic, entity, userInfo)); + } catch (RuntimeException ex) { + logger.error("Failed to create MQTT event for " + getDisplayEntityTypeName(), ex); + } + } + /** * On database operation clear cache of related cached entity when needed. */ diff --git a/src/lib/mqtt-rar-0.8.0.rar b/src/lib/mqtt-rar-0.8.0.rar new file mode 100644 index 000000000..5f21fdcae Binary files /dev/null and b/src/lib/mqtt-rar-0.8.0.rar differ diff --git a/tools/developer_tools/cdb-base-mqtt/.gitignore b/tools/developer_tools/cdb-base-mqtt/.gitignore new file mode 100644 index 000000000..2f7896d1d --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/.gitignore @@ -0,0 +1 @@ +target/ diff --git a/tools/developer_tools/cdb-base-mqtt/README.md b/tools/developer_tools/cdb-base-mqtt/README.md new file mode 100644 index 000000000..7a00f06e6 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/README.md @@ -0,0 +1,29 @@ +# CDB Base MQTT + +Shared Java 8 MQTT event base class, publishing API, and Payara JCA/JNDI transport used by ComponentDB and BELY. + +`gov.anl.aps.cdb.base.mqtt.MqttEvent` owns the event timestamp and requires applications to supply a string topic. Application-specific payload fields and topic definitions remain in each portal. + +## Build and test + +A JDK and Maven are required. + +```sh +mvn clean verify +sha256sum -c SHA256SUMS +``` + +The versioned thin JAR is written to `target/cdb-base-mqtt-1.0.0.jar`. Jackson is provided by each consuming portal and is intentionally not bundled. + +GitHub Actions runs `mvn verify` on pushes and pull requests that change this module or its workflow. It also supports manual runs and uploads the Surefire reports as a workflow artifact. + +## Release artifact + +The checked-in `SHA256SUMS` records the released JAR checksum and verifies the artifact under `target/`. After changing this library: + +1. Increment the version in `pom.xml`. +2. Run `mvn clean verify`. +3. Update `SHA256SUMS` with the checksum of the generated JAR. +4. Copy that exact JAR into each portal's checked-in `lib/` directory. +5. Register the versioned JAR in each portal's NetBeans Ant classpath. +6. Verify each copied JAR with `sha256sum` before building the portals. diff --git a/tools/developer_tools/cdb-base-mqtt/SHA256SUMS b/tools/developer_tools/cdb-base-mqtt/SHA256SUMS new file mode 100644 index 000000000..c381ed494 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/SHA256SUMS @@ -0,0 +1 @@ +42a4db3b8f35d7a3785849c0f3365af9f5dec7cffa915899c460c5f90e5e8819 target/cdb-base-mqtt-1.0.0.jar diff --git a/tools/developer_tools/cdb-base-mqtt/pom.xml b/tools/developer_tools/cdb-base-mqtt/pom.xml new file mode 100644 index 000000000..3292ba2fe --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/pom.xml @@ -0,0 +1,56 @@ + + 4.0.0 + + gov.anl.aps + cdb-base-mqtt + 1.0.0 + + + 1.8 + 1.8 + 2026-09-30T00:00:00Z + UTF-8 + 5.10.3 + + + + + com.fasterxml.jackson.core + jackson-annotations + 2.11.2 + provided + + + com.fasterxml.jackson.core + jackson-databind + 2.11.2 + provided + + + org.junit.jupiter + junit-jupiter + ${junit.version} + test + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.13.0 + + 8 + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.2.5 + + + + diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttEvent.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttEvent.java new file mode 100644 index 000000000..6a3482bc1 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttEvent.java @@ -0,0 +1,26 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt; + +import com.fasterxml.jackson.annotation.JsonFormat; +import com.fasterxml.jackson.annotation.JsonIgnore; +import java.util.Date; + +/** Base for an application-owned payload and its destination topic. */ +public abstract class MqttEvent { + private final Date eventTimestamp; + + protected MqttEvent() { + eventTimestamp = new Date(); + } + + @JsonIgnore + public abstract String getTopic(); + + @JsonFormat(shape = JsonFormat.Shape.STRING) + public Date getEventTimestamp() { + return new Date(eventTimestamp.getTime()); + } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttPublisher.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttPublisher.java new file mode 100644 index 000000000..873d73a07 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/MqttPublisher.java @@ -0,0 +1,11 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt; + +import gov.anl.aps.cdb.base.mqtt.exception.MqttException; + +public interface MqttPublisher { + void publish(MqttEvent event) throws MqttException; +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttException.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttException.java new file mode 100644 index 000000000..91982bbe0 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttException.java @@ -0,0 +1,10 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.exception; + +public abstract class MqttException extends Exception { + protected MqttException(String message) { super(message); } + protected MqttException(String message, Throwable cause) { super(message, cause); } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttNotConfiguredException.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttNotConfiguredException.java new file mode 100644 index 000000000..dbefec5be --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttNotConfiguredException.java @@ -0,0 +1,10 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.exception; + +public class MqttNotConfiguredException extends MqttException { + public MqttNotConfiguredException(String message) { super(message); } + public MqttNotConfiguredException(String message, Throwable cause) { super(message, cause); } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttPublishException.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttPublishException.java new file mode 100644 index 000000000..529ea5025 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttPublishException.java @@ -0,0 +1,9 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.exception; + +public class MqttPublishException extends MqttException { + public MqttPublishException(String message, Throwable cause) { super(message, cause); } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttSerializationException.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttSerializationException.java new file mode 100644 index 000000000..19ad831ed --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/exception/MqttSerializationException.java @@ -0,0 +1,9 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.exception; + +public class MqttSerializationException extends MqttException { + public MqttSerializationException(String message, Throwable cause) { super(message, cause); } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JcaMqttPublisher.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JcaMqttPublisher.java new file mode 100644 index 000000000..76bac90f1 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JcaMqttPublisher.java @@ -0,0 +1,39 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.jca; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import gov.anl.aps.cdb.base.mqtt.MqttEvent; +import gov.anl.aps.cdb.base.mqtt.MqttPublisher; +import gov.anl.aps.cdb.base.mqtt.exception.MqttException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttPublishException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttSerializationException; +import java.nio.charset.StandardCharsets; + +public class JcaMqttPublisher implements MqttPublisher { + private final MqttConnectionProvider connections; + private final ObjectMapper mapper; + private final int qos; + private final boolean retained; + + public JcaMqttPublisher(MqttConnectionProvider connections, ObjectMapper mapper) { this(connections, mapper, 0, false); } + public JcaMqttPublisher(MqttConnectionProvider connections, ObjectMapper mapper, int qos, boolean retained) { + if (connections == null || mapper == null) throw new IllegalArgumentException("connections and mapper are required"); + if (qos < 0 || qos > 2) throw new IllegalArgumentException("qos must be between 0 and 2"); + this.connections = connections; this.mapper = mapper; this.qos = qos; this.retained = retained; + } + + @Override public void publish(MqttEvent event) throws MqttException { + if (event == null || event.getTopic() == null || event.getTopic().trim().isEmpty()) throw new IllegalArgumentException("event topic is required"); + final byte[] payload; + try { payload = mapper.writeValueAsString(event).getBytes(StandardCharsets.UTF_8); } + catch (JsonProcessingException ex) { throw new MqttSerializationException("Failed to serialize MQTT event", ex); } + try (MqttConnection connection = connections.getConnection()) { + connection.publish(event.getTopic(), payload, qos, retained); + } catch (MqttException ex) { throw ex; } + catch (Exception ex) { throw new MqttPublishException("Failed to publish MQTT event to " + event.getTopic(), ex); } + } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JndiMqttConnectionProvider.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JndiMqttConnectionProvider.java new file mode 100644 index 000000000..cfd83e7e8 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/JndiMqttConnectionProvider.java @@ -0,0 +1,61 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.jca; + +import gov.anl.aps.cdb.base.mqtt.exception.MqttNotConfiguredException; +import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Method; +import javax.naming.InitialContext; +import javax.naming.NamingException; + +/** Looks up a Payara MQTT JCA factory without binding consumers to its API version. */ +public class JndiMqttConnectionProvider implements MqttConnectionProvider { + private final String jndiName; + + public JndiMqttConnectionProvider(String jndiName) { + if (jndiName == null || jndiName.trim().isEmpty()) throw new IllegalArgumentException("jndiName is required"); + this.jndiName = jndiName; + } + + @Override public MqttConnection getConnection() throws MqttNotConfiguredException { + try { + Object factory = new InitialContext().lookup(jndiName); + Object connection = factory.getClass().getMethod("getConnection").invoke(factory); + if (connection == null) { + throw new MqttNotConfiguredException( + "MQTT resource could not connect to its broker: " + jndiName, null); + } + return new ReflectiveConnection(connection); + } catch (NamingException | ReflectiveOperationException | LinkageError ex) { + throw new MqttNotConfiguredException("MQTT resource is unavailable: " + jndiName, unwrap(ex)); + } + } + + private static Throwable unwrap(Throwable error) { + return error instanceof InvocationTargetException && error.getCause() != null ? error.getCause() : error; + } + + private static final class ReflectiveConnection implements MqttConnection { + private final Object delegate; + private final Method publish; + private final Method close; + ReflectiveConnection(Object delegate) throws NoSuchMethodException { + this.delegate = delegate; + publish = delegate.getClass().getMethod("publish", String.class, byte[].class, int.class, boolean.class); + close = delegate.getClass().getMethod("close"); + } + @Override public void publish(String topic, byte[] payload, int qos, boolean retained) throws Exception { + try { publish.invoke(delegate, topic, payload, qos, retained); } + catch (InvocationTargetException ex) { throw asException(unwrap(ex)); } + } + @Override public void close() throws Exception { + try { close.invoke(delegate); } + catch (InvocationTargetException ex) { throw asException(unwrap(ex)); } + } + private static Exception asException(Throwable error) { + return error instanceof Exception ? (Exception) error : new Exception(error); + } + } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnection.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnection.java new file mode 100644 index 000000000..19741d331 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnection.java @@ -0,0 +1,10 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.jca; + +public interface MqttConnection extends AutoCloseable { + void publish(String topic, byte[] payload, int qos, boolean retained) throws Exception; + @Override void close() throws Exception; +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnectionProvider.java b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnectionProvider.java new file mode 100644 index 000000000..e2209bfa9 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/main/java/gov/anl/aps/cdb/base/mqtt/jca/MqttConnectionProvider.java @@ -0,0 +1,11 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt.jca; + +import gov.anl.aps.cdb.base.mqtt.exception.MqttNotConfiguredException; + +public interface MqttConnectionProvider { + MqttConnection getConnection() throws MqttNotConfiguredException; +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JcaMqttPublisherTest.java b/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JcaMqttPublisherTest.java new file mode 100644 index 000000000..9cb4142b7 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JcaMqttPublisherTest.java @@ -0,0 +1,56 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt; + +import com.fasterxml.jackson.databind.ObjectMapper; +import gov.anl.aps.cdb.base.mqtt.exception.MqttNotConfiguredException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttPublishException; +import gov.anl.aps.cdb.base.mqtt.exception.MqttSerializationException; +import gov.anl.aps.cdb.base.mqtt.jca.JcaMqttPublisher; +import gov.anl.aps.cdb.base.mqtt.jca.MqttConnection; +import gov.anl.aps.cdb.base.mqtt.jca.MqttConnectionProvider; +import java.nio.charset.StandardCharsets; +import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.*; + +class JcaMqttPublisherTest { + @Test void publishesConfiguredArgumentsAndCloses() throws Exception { + FakeConnection connection = new FakeConnection(); + JcaMqttPublisher publisher = new JcaMqttPublisher(() -> connection, new ObjectMapper(), 1, true); + publisher.publish(new Event("value")); + assertEquals("cdb/add", connection.topic); + String payload = new String(connection.payload, StandardCharsets.UTF_8); + assertTrue(payload.contains("\"value\":\"value\"")); + assertTrue(payload.contains("\"eventTimestamp\":")); + assertFalse(payload.contains("\"topic\":")); + assertEquals(1, connection.qos); assertTrue(connection.retained); assertTrue(connection.closed); + } + @Test void reportsMissingConfiguration() { + MqttConnectionProvider provider = () -> { throw new MqttNotConfiguredException("missing"); }; + assertThrows(MqttNotConfiguredException.class, () -> new JcaMqttPublisher(provider, new ObjectMapper()).publish(new Event("x"))); + } + @Test void reportsSerializationFailureWithoutOpeningConnection() { + boolean[] opened = {false}; + ObjectMapper mapper = new ObjectMapper() { @Override public String writeValueAsString(Object value) throws com.fasterxml.jackson.core.JsonProcessingException { throw new com.fasterxml.jackson.core.JsonProcessingException("bad") {}; } }; + assertThrows(MqttSerializationException.class, () -> new JcaMqttPublisher(() -> { opened[0]=true; return new FakeConnection(); }, mapper).publish(new Event("x"))); + assertFalse(opened[0]); + } + @Test void closesAfterPublishFailure() { + FakeConnection connection = new FakeConnection(); connection.failPublish = true; + assertThrows(MqttPublishException.class, () -> new JcaMqttPublisher(() -> connection, new ObjectMapper()).publish(new Event("x"))); + assertTrue(connection.closed); + } + @Test void rejectsInvalidQos() { assertThrows(IllegalArgumentException.class, () -> new JcaMqttPublisher(() -> null, new ObjectMapper(), 3, false)); } + + static final class Event extends MqttEvent { + private final String value; Event(String value){this.value=value;} + @Override public String getTopic(){return "cdb/add";} public String getValue(){return value;} + } + static final class FakeConnection implements MqttConnection { + String topic; byte[] payload; int qos; boolean retained; boolean closed; boolean failPublish; + @Override public void publish(String topic, byte[] payload, int qos, boolean retained) throws Exception { this.topic=topic;this.payload=payload;this.qos=qos;this.retained=retained;if(failPublish)throw new Exception("broker"); } + @Override public void close(){closed=true;} + } +} diff --git a/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JndiMqttConnectionProviderTest.java b/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JndiMqttConnectionProviderTest.java new file mode 100644 index 000000000..60d5a0830 --- /dev/null +++ b/tools/developer_tools/cdb-base-mqtt/src/test/java/gov/anl/aps/cdb/base/mqtt/JndiMqttConnectionProviderTest.java @@ -0,0 +1,16 @@ +/* + * Copyright (c) UChicago Argonne, LLC. All rights reserved. + * See LICENSE file. + */ +package gov.anl.aps.cdb.base.mqtt; + +import gov.anl.aps.cdb.base.mqtt.exception.MqttNotConfiguredException; +import gov.anl.aps.cdb.base.mqtt.jca.JndiMqttConnectionProvider; +import org.junit.jupiter.api.Test; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class JndiMqttConnectionProviderTest { + @Test void missingJndiResourceIsReportedAsNotConfigured() { + assertThrows(MqttNotConfiguredException.class, () -> new JndiMqttConnectionProvider("java:comp/env/missing-mqtt").getConnection()); + } +}