This is an automated email from the ASF dual-hosted git repository. riemer pushed a commit to branch remove-broker-dependencies in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 02fa18800f0920c19d781b3b387bfe0a340cc96f Author: Dominik Riemer <[email protected]> AuthorDate: Wed Nov 1 23:36:24 2023 +0100 improvement: Remove broker dependencies from extensions management module --- streampipes-extensions-management/pom.xml | 25 ------------ .../connect/adapter/AdapterPipelineGenerator.java | 39 +------------------ .../elements/SendToBrokerAdapterSink.java | 44 ++++++++++++---------- .../elements/SendToJmsAdapterSink.java | 42 --------------------- .../elements/SendToKafkaAdapterSink.java | 43 --------------------- .../elements/SendToMqttAdapterSink.java | 42 --------------------- .../elements/SendToNatsAdapterSink.java | 43 --------------------- .../elements/SendToPulsarAdapterSink.java | 43 --------------------- .../opcua/config/SpOpcUaConfigExtractor.java | 2 +- .../streampipes-extensions-all-jvm/pom.xml | 5 +++ 10 files changed, 32 insertions(+), 296 deletions(-) diff --git a/streampipes-extensions-management/pom.xml b/streampipes-extensions-management/pom.xml index a63d88a2f..7ce180de6 100644 --- a/streampipes-extensions-management/pom.xml +++ b/streampipes-extensions-management/pom.xml @@ -63,31 +63,6 @@ <artifactId>streampipes-measurement-units</artifactId> <version>0.93.0-SNAPSHOT</version> </dependency> - <dependency> - <groupId>org.apache.streampipes</groupId> - <artifactId>streampipes-messaging-kafka</artifactId> - <version>0.93.0-SNAPSHOT</version> - </dependency> - <dependency> - <groupId>org.apache.streampipes</groupId> - <artifactId>streampipes-messaging-jms</artifactId> - <version>0.93.0-SNAPSHOT</version> - </dependency> - <dependency> - <groupId>org.apache.streampipes</groupId> - <artifactId>streampipes-messaging-mqtt</artifactId> - <version>0.93.0-SNAPSHOT</version> - </dependency> - <dependency> - <groupId>org.apache.streampipes</groupId> - <artifactId>streampipes-messaging-nats</artifactId> - <version>0.93.0-SNAPSHOT</version> - </dependency> - <dependency> - <groupId>org.apache.streampipes</groupId> - <artifactId>streampipes-messaging-pulsar</artifactId> - <version>0.93.0-SNAPSHOT</version> - </dependency> <dependency> <groupId>org.apache.streampipes</groupId> <artifactId>streampipes-sdk</artifactId> diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java index 3fe245255..2e20e276c 100644 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java +++ b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java @@ -19,21 +19,9 @@ package org.apache.streampipes.extensions.management.connect.adapter; import org.apache.streampipes.connect.shared.AdapterPipelineGeneratorBase; -import org.apache.streampipes.extensions.management.client.StreamPipesClientResolver; import org.apache.streampipes.extensions.management.connect.adapter.model.pipeline.AdapterPipeline; import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToBrokerAdapterSink; -import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToJmsAdapterSink; -import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToKafkaAdapterSink; -import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToMqttAdapterSink; -import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToNatsAdapterSink; -import org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToPulsarAdapterSink; -import org.apache.streampipes.model.configuration.MessagingSettings; -import org.apache.streampipes.model.configuration.SpProtocol; import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.JmsTransportProtocol; -import org.apache.streampipes.model.grounding.KafkaTransportProtocol; -import org.apache.streampipes.model.grounding.MqttTransportProtocol; -import org.apache.streampipes.model.grounding.PulsarTransportProtocol; public class AdapterPipelineGenerator extends AdapterPipelineGeneratorBase { @@ -51,31 +39,8 @@ public class AdapterPipelineGenerator extends AdapterPipelineGeneratorBase { } } - private SendToBrokerAdapterSink<?> getAdapterSink(AdapterDescription adapterDescription) { - var prioritizedProtocol = - getMessagingSettings().getPrioritizedProtocols().get(0); - - if (isPrioritized(prioritizedProtocol, JmsTransportProtocol.class)) { - return new SendToJmsAdapterSink(adapterDescription); - } else if (isPrioritized(prioritizedProtocol, KafkaTransportProtocol.class)) { - return new SendToKafkaAdapterSink(adapterDescription); - } else if (isPrioritized(prioritizedProtocol, MqttTransportProtocol.class)) { - return new SendToMqttAdapterSink(adapterDescription); - } else if (isPrioritized(prioritizedProtocol, PulsarTransportProtocol.class)) { - return new SendToPulsarAdapterSink(adapterDescription); - } else { - return new SendToNatsAdapterSink(adapterDescription); - } - } - - private boolean isPrioritized(SpProtocol prioritizedProtocol, - Class<?> protocolClass) { - return prioritizedProtocol.getProtocolClass().equals(protocolClass.getCanonicalName()); - } - - private MessagingSettings getMessagingSettings() { - var client = new StreamPipesClientResolver().makeStreamPipesClientInstance(); - return client.adminApi().getMessagingSettings(); + private SendToBrokerAdapterSink getAdapterSink(AdapterDescription adapterDescription) { + return new SendToBrokerAdapterSink(adapterDescription); } private boolean hasValidGrounding(AdapterDescription adapterDescription) { diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java index b9cebe258..66983df35 100644 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java +++ b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java @@ -19,51 +19,52 @@ package org.apache.streampipes.extensions.management.connect.adapter.preprocessi import org.apache.streampipes.commons.environment.Environment; import org.apache.streampipes.commons.environment.Environments; -import org.apache.streampipes.commons.exceptions.SpRuntimeException; import org.apache.streampipes.dataformat.SpDataFormatDefinition; import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager; import org.apache.streampipes.extensions.management.connect.adapter.util.TransportFormatSelector; import org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger; import org.apache.streampipes.messaging.EventProducer; +import org.apache.streampipes.messaging.SpProtocolManager; import org.apache.streampipes.model.connect.adapter.AdapterDescription; +import org.apache.streampipes.model.grounding.KafkaTransportProtocol; import org.apache.streampipes.model.grounding.TransportFormat; import org.apache.streampipes.model.grounding.TransportProtocol; import java.util.Map; -public abstract class SendToBrokerAdapterSink<T extends TransportProtocol> implements IAdapterPipelineElement { +public class SendToBrokerAdapterSink implements IAdapterPipelineElement { protected AdapterDescription adapterDescription; protected SpDataFormatDefinition dataFormatDefinition; - protected T protocol; + protected TransportProtocol protocol; private final EventProducer producer; - public SendToBrokerAdapterSink(AdapterDescription adapterDescription, - Class<T> protocolClass) { + public SendToBrokerAdapterSink(AdapterDescription adapterDescription) { this.adapterDescription = adapterDescription; - this.protocol = protocolClass.cast(adapterDescription + this.protocol = adapterDescription .getEventGrounding() - .getTransportProtocol()); + .getTransportProtocol(); if (getEnvironment().getSpDebug().getValueOrDefault()) { modifyProtocolForDebugging(this.protocol); } - this.producer = makeProducer(this.protocol); + var producerOpt = SpProtocolManager.INSTANCE.findDefinition(this.protocol); + if (producerOpt.isPresent()) { + this.producer = producerOpt.get().getProducer(this.protocol); - TransportFormat transportFormat = adapterDescription - .getEventGrounding() - .getTransportFormats() - .get(0); + TransportFormat transportFormat = adapterDescription + .getEventGrounding() + .getTransportFormats() + .get(0); - this.dataFormatDefinition = - new TransportFormatSelector(transportFormat).getDataFormatDefinition(); + this.dataFormatDefinition = + new TransportFormatSelector(transportFormat).getDataFormatDefinition(); - try { producer.connect(); - } catch (SpRuntimeException e) { - e.printStackTrace(); + } else { + throw new RuntimeException("Could not find protocol"); } } @@ -86,9 +87,12 @@ public abstract class SendToBrokerAdapterSink<T extends TransportProtocol> imple producer.publish(event); } - protected abstract EventProducer makeProducer(T protocol); - - public abstract void modifyProtocolForDebugging(T protocol); + public void modifyProtocolForDebugging(TransportProtocol protocol) { + protocol.setBrokerHostname("localhost"); + if (protocol instanceof KafkaTransportProtocol) { + ((KafkaTransportProtocol) protocol).setKafkaPort(9094); + } + } private Environment getEnvironment() { return Environments.getEnvironment(); diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java deleted file mode 100644 index 98af0b42c..000000000 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java +++ /dev/null @@ -1,42 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ -package org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements; - -import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; -import org.apache.streampipes.messaging.EventProducer; -import org.apache.streampipes.messaging.jms.ActiveMQPublisher; -import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.JmsTransportProtocol; - -public class SendToJmsAdapterSink extends SendToBrokerAdapterSink<JmsTransportProtocol> - implements IAdapterPipelineElement { - - public SendToJmsAdapterSink(AdapterDescription adapterDescription) { - super(adapterDescription, JmsTransportProtocol.class); - } - - @Override - protected EventProducer makeProducer(JmsTransportProtocol protocol) { - return new ActiveMQPublisher(protocol); - } - - @Override - public void modifyProtocolForDebugging(JmsTransportProtocol protocol) { - protocol.setBrokerHostname("localhost"); - } -} diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java deleted file mode 100644 index 5db6d32b2..000000000 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ -package org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements; - -import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; -import org.apache.streampipes.messaging.EventProducer; -import org.apache.streampipes.messaging.kafka.SpKafkaProducer; -import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.KafkaTransportProtocol; - -public class SendToKafkaAdapterSink extends SendToBrokerAdapterSink<KafkaTransportProtocol> - implements IAdapterPipelineElement { - - public SendToKafkaAdapterSink(AdapterDescription adapterDescription) { - super(adapterDescription, KafkaTransportProtocol.class); - } - - @Override - protected EventProducer makeProducer(KafkaTransportProtocol protocol) { - return new SpKafkaProducer(protocol); - } - - @Override - public void modifyProtocolForDebugging(KafkaTransportProtocol protocol) { - protocol.setBrokerHostname("localhost"); - protocol.setKafkaPort(9094); - } -} diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java deleted file mode 100644 index d2823c1c5..000000000 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java +++ /dev/null @@ -1,42 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ -package org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements; - -import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; -import org.apache.streampipes.messaging.EventProducer; -import org.apache.streampipes.messaging.mqtt.MqttPublisher; -import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.MqttTransportProtocol; - -public class SendToMqttAdapterSink extends SendToBrokerAdapterSink<MqttTransportProtocol> - implements IAdapterPipelineElement { - - public SendToMqttAdapterSink(AdapterDescription adapterDescription) { - super(adapterDescription, MqttTransportProtocol.class); - } - - @Override - protected EventProducer makeProducer(MqttTransportProtocol protocol) { - return new MqttPublisher(protocol); - } - - @Override - public void modifyProtocolForDebugging(MqttTransportProtocol transportProtocol) { - protocol.setBrokerHostname("localhost"); - } -} diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java deleted file mode 100644 index 5216ee0b0..000000000 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements; - -import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; -import org.apache.streampipes.messaging.EventProducer; -import org.apache.streampipes.messaging.nats.NatsPublisher; -import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.NatsTransportProtocol; - -public class SendToNatsAdapterSink extends SendToBrokerAdapterSink<NatsTransportProtocol> - implements IAdapterPipelineElement { - - public SendToNatsAdapterSink(AdapterDescription adapterDescription) { - super(adapterDescription, NatsTransportProtocol.class); - } - - @Override - protected EventProducer makeProducer(NatsTransportProtocol protocol) { - return new NatsPublisher(protocol); - } - - @Override - public void modifyProtocolForDebugging(NatsTransportProtocol protocol) { - protocol.setBrokerHostname("localhost"); - } -} diff --git a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java deleted file mode 100644 index b9f3ba56d..000000000 --- a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java +++ /dev/null @@ -1,43 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements; - -import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement; -import org.apache.streampipes.messaging.EventProducer; -import org.apache.streampipes.messaging.pulsar.PulsarProducer; -import org.apache.streampipes.model.connect.adapter.AdapterDescription; -import org.apache.streampipes.model.grounding.PulsarTransportProtocol; - -public class SendToPulsarAdapterSink extends SendToBrokerAdapterSink<PulsarTransportProtocol> - implements IAdapterPipelineElement { - - public SendToPulsarAdapterSink(AdapterDescription adapterDescription) { - super(adapterDescription, PulsarTransportProtocol.class); - } - - @Override - protected EventProducer makeProducer(PulsarTransportProtocol protocol) { - return new PulsarProducer(protocol); - } - - @Override - public void modifyProtocolForDebugging(PulsarTransportProtocol protocol) { - protocol.setBrokerHostname("pulsar://localhost:6650"); - } -} diff --git a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java index 2c3e3b138..612b3e490 100644 --- a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java +++ b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java @@ -24,7 +24,6 @@ import org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaUtil; import java.util.List; -import static org.apache.kafka.common.config.ConfigDef.Type.PASSWORD; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.ACCESS_MODE; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.ADAPTER_TYPE; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.AVAILABLE_NODES; @@ -33,6 +32,7 @@ import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabe import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_SERVER_PORT; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_SERVER_URL; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_URL; +import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PASSWORD; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULLING_INTERVAL; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULL_MODE; import static org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.UNAUTHENTICATED; diff --git a/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml b/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml index 3b03c5940..bfaa6dc03 100644 --- a/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml +++ b/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml @@ -64,6 +64,11 @@ <classifier>embed</classifier> </dependency> + <dependency> + <groupId>org.apache.streampipes</groupId> + <artifactId>streampipes-messaging-pulsar</artifactId> + <version>0.93.0-SNAPSHOT</version> + </dependency> <dependency> <groupId>org.apache.streampipes</groupId> <artifactId>streampipes-messaging-nats</artifactId>
