getSupportedCertificateTypes() {
+ return Arrays.asList(CertificateType.X_509, CertificateType.RAW_PUBLIC_KEY);
+ }
+
+ /**
+ * Validates the certificate provided by the the peer as part of the
+ * certificate message.
+ *
+ * If a x509 certificate chain is provided in the certificate message,
+ * validate the chain and key usage. If a RawPublicKey certificate is
+ * provided, check, if this public key is trusted.
+ *
+ * @param cid connection ID
+ * @param serverName indicated server names. May be {@code null}, if not
+ * available or SNI is not enabled.
+ * @param remotePeer socket address of remote peer
+ * @param clientUsage indicator to check certificate usage. {@code true},
+ * check key usage for client, {@code false} for server.
+ * @param verifySubject {@code true} to verify the certificate's subjects,
+ * {@code false}, if not.
+ * @param truncateCertificatePath {@code true} truncate certificate path at
+ * a trusted certificate before validation.
+ * @param message certificate message to be validated
+ * @return certificate verification result, or {@code null}, if result is
+ * provided asynchronous.
+ * @since 3.0 (removed DTLSSession session, added remotePeer and
+ * verifySubject)
+ */
+ @Override
+ public CertificateVerificationResult verifyCertificate(ConnectionId cid, ServerNames serverName, InetSocketAddress remotePeer,
+ boolean clientUsage, boolean verifySubject, boolean truncateCertificatePath,
+ CertificateMessage message) {
+ CertPath certChain = message.getCertificateChain();
+ CertificateVerificationResult result;
+
+ if (certChain == null) {
+ PublicKey publicKey = message.getPublicKey();
+ result = new CertificateVerificationResult(cid, publicKey, null);
+ } else {
+ if (message.getCertificateChain().getCertificates().isEmpty()) {
+ result = new CertificateVerificationResult(cid, new HandshakeException("Empty certificate chain",
+ new AlertMessage(AlertLevel.FATAL, AlertDescription.BAD_CERTIFICATE)), null);
+ } else {
+ result = new CertificateVerificationResult(cid, certChain, null);
+ }
+ }
+
+ return result;
+ }
+
+ /**
+ * Return an list of certificate authorities which are trusted
+ * for authenticating peers.
+ *
+ * @return a non-null (possibly empty) list of accepted CA issuers.
+ */
+ @Override
+ public List getAcceptedIssuers() {
+ log.trace("getAcceptedIssuers: return null");
+ return null;
+ }
+
+ /**
+ * Set the handler for asynchronous handshake results.
+ *
+ * Called during initialization of the {link DTLSConnector}. Synchronous
+ * implementations may just ignore this using an empty implementation.
+ *
+ * @param resultHandler handler for asynchronous master secret results. This
+ * handler MUST NOT be called from the thread calling
+ * {@link #verifyCertificate(ConnectionId, ServerNames, InetSocketAddress, boolean, boolean, boolean, CertificateMessage)},
+ * instead just return the result there.
+ */
+ @Override
+ public void setResultHandler(HandshakeResultHandler resultHandler) {
+ if (this.resultHandler != null && resultHandler != null && this.resultHandler != resultHandler) {
+ throw new IllegalStateException("handshake result handler already set!");
+ }
+ this.resultHandler = resultHandler;
+ }
+}
diff --git a/application/src/test/resources/coap/credentials/client/cert.pem b/application/src/test/resources/coap/credentials/client/cert.pem
new file mode 100644
index 0000000000..4d385a588f
--- /dev/null
+++ b/application/src/test/resources/coap/credentials/client/cert.pem
@@ -0,0 +1,13 @@
+-----BEGIN CERTIFICATE-----
+MIIB/TCCAaOgAwIBAgIIVNrVgKT9OE8wCgYIKoZIzj0EAwIwWjEOMAwGA1UEAxMF
+Y2YtY2ExFDASBgNVBAsTC0NhbGlmb3JuaXVtMRQwEgYDVQQKEwtFY2xpcHNlIElv
+VDEPMA0GA1UEBxMGT3R0YXdhMQswCQYDVQQGEwJDQTAeFw0yMzEwMjYwODA4MjJa
+Fw0yNTEwMjUwODA4MjJaMF4xEjAQBgNVBAMTCWNmLWNsaWVudDEUMBIGA1UECxML
+Q2FsaWZvcm5pdW0xFDASBgNVBAoTC0VjbGlwc2UgSW9UMQ8wDQYDVQQHEwZPdHRh
+d2ExCzAJBgNVBAYTAkNBMFkwEwYHKoZIzj0CAQYIKoZIzj0DAQcDQgAEQxYO5/M5
+ie6+3QPOaAy5MD6CkFILZwIb2rOBCX/EWPaocX1H+eynUnaEEbmqxeN6rnI/pH19
+j4PtsegfHLrzzaNPME0wHQYDVR0OBBYEFKwEDLTJ+5cQoZfbjWN1vJ2ssgK+MAsG
+A1UdDwQEAwIHgDAfBgNVHSMEGDAWgBSxVzoI1TL87++hsUb9vQwqODzgUTAKBggq
+hkjOPQQDAgNIADBFAiA2KCOw3n2AK9Vm8u2u1bQREIEs3tKAU7eFjpNFn929NwIh
+AInhBGoEwS2Xlu5bdZSfWnujoRrEQiIiQpStmLxVcIsH
+-----END CERTIFICATE-----
diff --git a/application/src/test/resources/coap/credentials/client/cert_01.pem b/application/src/test/resources/coap/credentials/client/cert_01.pem
new file mode 100644
index 0000000000..3b97ab4aad
--- /dev/null
+++ b/application/src/test/resources/coap/credentials/client/cert_01.pem
@@ -0,0 +1,14 @@
+-----BEGIN CERTIFICATE-----
+MIICIzCCAcmgAwIBAgIUZZCGYm65c9vU0Xfvd/pAnLVDouUwCgYIKoZIzj0EAwIw
+ZzELMAkGA1UEBhMCVUExDTALBgNVBAgMBEtpeXYxDTALBgNVBAcMBEtpeXYxFDAS
+BgNVBAoMC1RoaW5nc2JvYXJkMRIwEAYDVQQLDAlkZXZlbG9wZXIxEDAOBgNVBAMM
+B2NlcnRfMDEwHhcNMjQxMjE4MTU1NjE1WhcNMjUxMjE4MTU1NjE1WjBnMQswCQYD
+VQQGEwJVQTENMAsGA1UECAwES2l5djENMAsGA1UEBwwES2l5djEUMBIGA1UECgwL
+VGhpbmdzYm9hcmQxEjAQBgNVBAsMCWRldmVsb3BlcjEQMA4GA1UEAwwHY2VydF8w
+MTBZMBMGByqGSM49AgEGCCqGSM49AwEHA0IABNU1tE6o/QpqJJqpy+m+UoPuQe5g
+eTgS4M3x0iQS6pzNEJBhzbnOp/BysGMB4wKiAWTRuKdH/gcRXDBTjLd/d7ijUzBR
+MB0GA1UdDgQWBBSiao1iNWYzlsrSbxYqbda116HG1jAfBgNVHSMEGDAWgBSiao1i
+NWYzlsrSbxYqbda116HG1jAPBgNVHRMBAf8EBTADAQH/MAoGCCqGSM49BAMCA0gA
+MEUCIB2aCM/nvDqic9NkoSX/71GwksLiAKiFNkt2BZQykrcHAiEAr2h5IMdkyurN
+Jy/idx2y44CP0tMq/3QV0QLCQFJIi6s=
+-----END CERTIFICATE-----
diff --git a/application/src/test/resources/coap/credentials/client/key.pem b/application/src/test/resources/coap/credentials/client/key.pem
new file mode 100644
index 0000000000..02ca740c93
--- /dev/null
+++ b/application/src/test/resources/coap/credentials/client/key.pem
@@ -0,0 +1,4 @@
+-----BEGIN PRIVATE KEY-----
+MEECAQAwEwYHKoZIzj0CAQYIKoZIzj0DAQcEJzAlAgEBBCDn0+4CuLeX7xwBs0ts
+UUEDB3+HRwRKdIPeJlIbKuvvEQ==
+-----END PRIVATE KEY-----
\ No newline at end of file
diff --git a/application/src/test/resources/coap/credentials/client/key_01.pem b/application/src/test/resources/coap/credentials/client/key_01.pem
new file mode 100644
index 0000000000..d5918e8181
--- /dev/null
+++ b/application/src/test/resources/coap/credentials/client/key_01.pem
@@ -0,0 +1,8 @@
+-----BEGIN EC PARAMETERS-----
+BggqhkjOPQMBBw==
+-----END EC PARAMETERS-----
+-----BEGIN EC PRIVATE KEY-----
+MHcCAQEEIJldU1MBuJUJnNHa9Ob5NGlXc/Os6put9eh1TlIbuScnoAoGCCqGSM49
+AwEHoUQDQgAE1TW0Tqj9CmokmqnL6b5Sg+5B7mB5OBLgzfHSJBLqnM0QkGHNuc6n
+8HKwYwHjAqIBZNG4p0f+BxFcMFOMt393uA==
+-----END EC PRIVATE KEY-----
diff --git a/application/src/test/resources/coap/credentials/coapclientTest.jks b/application/src/test/resources/coap/credentials/coapclientTest.jks
new file mode 100644
index 0000000000..ca8c8ed1d7
Binary files /dev/null and b/application/src/test/resources/coap/credentials/coapclientTest.jks differ
diff --git a/application/src/test/resources/coap/credentials/coapserverTest.jks b/application/src/test/resources/coap/credentials/coapserverTest.jks
new file mode 100644
index 0000000000..4adf1f4f89
Binary files /dev/null and b/application/src/test/resources/coap/credentials/coapserverTest.jks differ
diff --git a/application/src/test/resources/coap/credentials/server/cert.pem b/application/src/test/resources/coap/credentials/server/cert.pem
new file mode 100644
index 0000000000..03eb9e372e
--- /dev/null
+++ b/application/src/test/resources/coap/credentials/server/cert.pem
@@ -0,0 +1,35 @@
+-----BEGIN CERTIFICATE-----
+MIIBaDCCAQ2gAwIBAgIUCY+goBAOhowBs7BHs/qXdAX8XFgwCgYIKoZIzj0EAwIw
+ETEPMA0GA1UEAwwGUm9vdENBMB4XDTI0MTIxOTEzNTY1OFoXDTM0MTIxNzEzNTY1
+OFowETEPMA0GA1UEAwwGU2VydmVyMFkwEwYHKoZIzj0CAQYIKoZIzj0DAQcDQgAE
+/qief3Kjnz0FpkQVaKRqJq3kHmCqqs+y1EGYLEZZAqLFvxmv7xoL6muG4Mj8tzqk
+Ll94JJuz97hG1FiEZsq7O6NDMEEwCwYDVR0PBAQDAgWgMBMGA1UdJQQMMAoGCCsG
+AQUFBwMBMB0GA1UdDgQWBBTK/UPsN0I2ErVPILWKMRV6TSeAmTAKBggqhkjOPQQD
+AgNJADBGAiEA8EhlOwvTbwGlxo55UIOJp9LBbCp0BEIWojlu8PzOVSsCIQDlV24S
+3BUJVCuMRujO5lTfJLxaSKkOEIgRANwIGi88WA==
+-----END CERTIFICATE-----
+-----BEGIN CERTIFICATE-----
+MIIBGzCBwgIUP/PGQOKa5EyvsIXNgvv9PNietyEwCgYIKoZIzj0EAwMwEDEOMAwG
+A1UEAwwFVFJVU1QwHhcNMjQxMjE5MTM1NjU4WhcNMzQxMjE3MTM1NjU4WjARMQ8w
+DQYDVQQDDAZSb290Q0EwWTATBgcqhkjOPQIBBggqhkjOPQMBBwNCAAT+qJ5/cqOf
+PQWmRBVopGomreQeYKqqz7LUQZgsRlkCosW/Ga/vGgvqa4bgyPy3OqQuX3gkm7P3
+uEbUWIRmyrs7MAoGCCqGSM49BAMDA0gAMEUCIQD2DY3UDXbzaIBKrsCtohKlEunH
+ip9LkSeYfSKCnfm23gIgA8AEJdunpRmPkilxgy6wZSLLROqDpGDnhnyv8dsR8cc=
+-----END CERTIFICATE-----
+-----BEGIN CERTIFICATE-----
+MIIBLTCB1AIUcsuauXAqvIS2RQcNPYysETJUAvwwCgYIKoZIzj0EAwMwIzEhMB8G
+A1UEAwwYQUFBIENlcnRpZmljYXRlIFNlcnZpY2VzMB4XDTI0MTIxOTEzNTY1OFoX
+DTM0MTIxNzEzNTY1OFowEDEOMAwGA1UEAwwFVFJVU1QwWTATBgcqhkjOPQIBBggq
+hkjOPQMBBwNCAAT+qJ5/cqOfPQWmRBVopGomreQeYKqqz7LUQZgsRlkCosW/Ga/v
+Ggvqa4bgyPy3OqQuX3gkm7P3uEbUWIRmyrs7MAoGCCqGSM49BAMDA0gAMEUCIQCM
+DV8sfoArfWiXAUF2LNS3kkHD7sgb91jr2+poEHgBBgIgXf9VeJp3K5jHX6lJwtE8
+nd+jW7T9nhTc/5njHg7xons=
+-----END CERTIFICATE-----
+-----BEGIN EC PARAMETERS-----
+BggqhkjOPQMBBw==
+-----END EC PARAMETERS-----
+-----BEGIN EC PRIVATE KEY-----
+MHcCAQEEIB+Z69so6HqCCWo5VOFxGsLXOlTWIYijOtzt+SeNGrgPoAoGCCqGSM49
+AwEHoUQDQgAE/qief3Kjnz0FpkQVaKRqJq3kHmCqqs+y1EGYLEZZAqLFvxmv7xoL
+6muG4Mj8tzqkLl94JJuz97hG1FiEZsq7Ow==
+-----END EC PRIVATE KEY-----
diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/CoapServerService.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/CoapServerService.java
index 0b1ba35709..5f4f1d152b 100644
--- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/CoapServerService.java
+++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/CoapServerService.java
@@ -17,7 +17,6 @@ package org.thingsboard.server.coapserver;
import org.eclipse.californium.core.CoapServer;
-import java.net.InetSocketAddress;
import java.net.UnknownHostException;
import java.util.concurrent.ConcurrentMap;
@@ -25,5 +24,5 @@ public interface CoapServerService {
CoapServer getCoapServer() throws UnknownHostException;
- ConcurrentMap getDtlsSessionsMap();
+ ConcurrentMap getDtlsSessionsMap();
}
diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java
index f3d30bfea4..41b10b31a3 100644
--- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java
+++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/DefaultCoapServerService.java
@@ -78,7 +78,7 @@ public class DefaultCoapServerService implements CoapServerService {
}
@Override
- public ConcurrentMap getDtlsSessionsMap() {
+ public ConcurrentMap getDtlsSessionsMap() {
return tbDtlsCertificateVerifier != null ? tbDtlsCertificateVerifier.getTbCoapDtlsSessionsMap() : null;
}
diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java
index cc7dcad77c..d61d7f279c 100644
--- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java
+++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsCertificateVerifier.java
@@ -109,7 +109,8 @@ public class TbCoapDtlsCertificateVerifier implements NewAdvancedCertificateVeri
if (msg != null && strCert.equals(msg.getCredentials())) {
DeviceProfile deviceProfile = msg.getDeviceProfile();
if (msg.hasDeviceInfo() && deviceProfile != null) {
- tbCoapDtlsSessionInMemoryStorage.put(remotePeer, new TbCoapDtlsSessionInfo(msg, deviceProfile));
+ TbCoapDtlsSessionKey tbCoapDtlsSessionKey = new TbCoapDtlsSessionKey(remotePeer, msg.getCredentials());
+ tbCoapDtlsSessionInMemoryStorage.put(tbCoapDtlsSessionKey, new TbCoapDtlsSessionInfo(msg, deviceProfile));
}
break;
}
@@ -138,7 +139,7 @@ public class TbCoapDtlsCertificateVerifier implements NewAdvancedCertificateVeri
public void setResultHandler(HandshakeResultHandler resultHandler) {
}
- public ConcurrentMap getTbCoapDtlsSessionsMap() {
+ public ConcurrentMap getTbCoapDtlsSessionsMap() {
return tbCoapDtlsSessionInMemoryStorage.getDtlsSessionsMap();
}
diff --git a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInMemoryStorage.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInMemoryStorage.java
index b4101f1763..5ff44561d8 100644
--- a/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInMemoryStorage.java
+++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionInMemoryStorage.java
@@ -18,7 +18,6 @@ package org.thingsboard.server.coapserver;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
-import java.net.InetSocketAddress;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -26,7 +25,7 @@ import java.util.concurrent.ConcurrentMap;
@Data
public class TbCoapDtlsSessionInMemoryStorage {
- private final ConcurrentMap dtlsSessionsMap = new ConcurrentHashMap<>();
+ private final ConcurrentMap dtlsSessionsMap = new ConcurrentHashMap<>();
private long dtlsSessionInactivityTimeout;
private long dtlsSessionReportTimeout;
@@ -36,9 +35,9 @@ public class TbCoapDtlsSessionInMemoryStorage {
this.dtlsSessionReportTimeout = dtlsSessionReportTimeout;
}
- public void put(InetSocketAddress remotePeer, TbCoapDtlsSessionInfo dtlsSessionInfo) {
- log.trace("DTLS session added to in-memory store: [{}] timestamp: [{}]", remotePeer, dtlsSessionInfo.getLastActivityTime());
- dtlsSessionsMap.putIfAbsent(remotePeer, dtlsSessionInfo);
+ public void put(TbCoapDtlsSessionKey tbCoapDtlsSessionKey, TbCoapDtlsSessionInfo dtlsSessionInfo) {
+ log.trace("DTLS session added to in-memory store: [{}] timestamp: [{}]", tbCoapDtlsSessionKey, dtlsSessionInfo.getLastActivityTime());
+ dtlsSessionsMap.putIfAbsent(tbCoapDtlsSessionKey, dtlsSessionInfo);
}
public void evictTimeoutSessions() {
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/AwsSqsTbQueueMsgMetadata.java b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionKey.java
similarity index 53%
rename from common/queue/src/main/java/org/thingsboard/server/queue/sqs/AwsSqsTbQueueMsgMetadata.java
rename to common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionKey.java
index 8d13a543d3..cf3e0b4fec 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/AwsSqsTbQueueMsgMetadata.java
+++ b/common/coap-server/src/main/java/org/thingsboard/server/coapserver/TbCoapDtlsSessionKey.java
@@ -13,16 +13,20 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.queue.sqs;
+package org.thingsboard.server.coapserver;
-import com.amazonaws.http.SdkHttpMetadata;
-import lombok.AllArgsConstructor;
-import lombok.Data;
-import org.thingsboard.server.queue.TbQueueMsgMetadata;
+import java.net.InetSocketAddress;
+import java.util.Objects;
-@Data
-@AllArgsConstructor
-public class AwsSqsTbQueueMsgMetadata implements TbQueueMsgMetadata {
+public record TbCoapDtlsSessionKey(InetSocketAddress peerAddress, String credentials) {
- private final SdkHttpMetadata metadata;
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (o == null || getClass() != o.getClass()) return false;
+ TbCoapDtlsSessionKey that = (TbCoapDtlsSessionKey) o;
+ return Objects.equals(peerAddress, that.peerAddress) &&
+ Objects.equals(credentials, that.credentials);
+ }
}
+
diff --git a/common/queue/pom.xml b/common/queue/pom.xml
index 6fad5efd5a..44e2239e05 100644
--- a/common/queue/pom.xml
+++ b/common/queue/pom.xml
@@ -64,18 +64,10 @@
org.apache.kafka
kafka-clients
-
- com.amazonaws
- aws-java-sdk-sqs
-
com.google.cloud
google-cloud-pubsub
-
- com.microsoft.azure
- azure-servicebus
-
com.rabbitmq
amqp-client
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java
index 87ad934c63..057351a449 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java
@@ -19,21 +19,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
-import org.thingsboard.server.queue.azure.servicebus.TbServiceBusAdmin;
-import org.thingsboard.server.queue.azure.servicebus.TbServiceBusQueueConfigs;
-import org.thingsboard.server.queue.azure.servicebus.TbServiceBusSettings;
import org.thingsboard.server.queue.kafka.TbKafkaAdmin;
import org.thingsboard.server.queue.kafka.TbKafkaSettings;
import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs;
-import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
-import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
-import org.thingsboard.server.queue.pubsub.TbPubSubSubscriptionSettings;
-import org.thingsboard.server.queue.rabbitmq.TbRabbitMqAdmin;
-import org.thingsboard.server.queue.rabbitmq.TbRabbitMqQueueArguments;
-import org.thingsboard.server.queue.rabbitmq.TbRabbitMqSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
@Configuration
public class RuleEngineTbQueueAdminFactory {
@@ -43,56 +31,12 @@ public class RuleEngineTbQueueAdminFactory {
@Autowired(required = false)
private TbKafkaSettings kafkaSettings;
- @Autowired(required = false)
- private TbAwsSqsQueueAttributes awsSqsQueueAttributes;
- @Autowired(required = false)
- private TbAwsSqsSettings awsSqsSettings;
-
- @Autowired(required = false)
- private TbPubSubSubscriptionSettings pubSubSubscriptionSettings;
- @Autowired(required = false)
- private TbPubSubSettings pubSubSettings;
-
- @Autowired(required = false)
- private TbRabbitMqQueueArguments rabbitMqQueueArguments;
- @Autowired(required = false)
- private TbRabbitMqSettings rabbitMqSettings;
-
- @Autowired(required = false)
- private TbServiceBusQueueConfigs serviceBusQueueConfigs;
- @Autowired(required = false)
- private TbServiceBusSettings serviceBusSettings;
-
@ConditionalOnExpression("'${queue.type:null}'=='kafka'")
@Bean
public TbQueueAdmin createKafkaAdmin() {
return new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getRuleEngineConfigs());
}
- @ConditionalOnExpression("'${queue.type:null}'=='aws-sqs'")
- @Bean
- public TbQueueAdmin createAwsSqsAdmin() {
- return new TbAwsSqsAdmin(awsSqsSettings, awsSqsQueueAttributes.getRuleEngineAttributes());
- }
-
- @ConditionalOnExpression("'${queue.type:null}'=='pubsub'")
- @Bean
- public TbQueueAdmin createPubSubAdmin() {
- return new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getRuleEngineSettings());
- }
-
- @ConditionalOnExpression("'${queue.type:null}'=='rabbitmq'")
- @Bean
- public TbQueueAdmin createRabbitMqAdmin() {
- return new TbRabbitMqAdmin(rabbitMqSettings, rabbitMqQueueArguments.getRuleEngineArgs());
- }
-
- @ConditionalOnExpression("'${queue.type:null}'=='service-bus'")
- @Bean
- public TbQueueAdmin createServiceBusAdmin() {
- return new TbServiceBusAdmin(serviceBusSettings, serviceBusQueueConfigs.getRuleEngineConfigs());
- }
-
@ConditionalOnExpression("'${queue.type:null}'=='in-memory'")
@Bean
public TbQueueAdmin createInMemoryAdmin() {
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java
deleted file mode 100644
index 3f920778e6..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java
+++ /dev/null
@@ -1,137 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.azure.servicebus;
-
-import com.microsoft.azure.servicebus.management.ManagementClient;
-import com.microsoft.azure.servicebus.management.QueueDescription;
-import com.microsoft.azure.servicebus.primitives.ConnectionStringBuilder;
-import com.microsoft.azure.servicebus.primitives.MessagingEntityAlreadyExistsException;
-import com.microsoft.azure.servicebus.primitives.ServiceBusException;
-import lombok.extern.slf4j.Slf4j;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.util.PropertyUtils;
-
-import java.io.IOException;
-import java.time.Duration;
-import java.util.Map;
-import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
-
-@Slf4j
-@Deprecated(forRemoval = true, since = "3.9") // for removal in 4.0
-public class TbServiceBusAdmin implements TbQueueAdmin {
- private final String MAX_SIZE = "maxSizeInMb";
- private final String MESSAGE_TIME_TO_LIVE = "messageTimeToLiveInSec";
- private final String LOCK_DURATION = "lockDurationInSec";
-
- private final Map queueConfigs;
- private final Set queues = ConcurrentHashMap.newKeySet();
-
- private final ManagementClient client;
-
- public TbServiceBusAdmin(TbServiceBusSettings serviceBusSettings, Map queueConfigs) {
- this.queueConfigs = queueConfigs;
-
- ConnectionStringBuilder builder = new ConnectionStringBuilder(
- serviceBusSettings.getNamespaceName(),
- "queues",
- serviceBusSettings.getSasKeyName(),
- serviceBusSettings.getSasKey());
-
- client = new ManagementClient(builder);
-
- try {
- client.getQueues().forEach(queueDescription -> queues.add(queueDescription.getPath()));
- } catch (ServiceBusException | InterruptedException e) {
- log.error("Failed to get queues.", e);
- throw new RuntimeException("Failed to get queues.", e);
- }
- }
-
- @Override
- public void createTopicIfNotExists(String topic, String properties) {
- if (queues.contains(topic)) {
- return;
- }
-
- try {
- QueueDescription queueDescription = new QueueDescription(topic);
- queueDescription.setRequiresDuplicateDetection(false);
- setQueueConfigs(queueDescription, PropertyUtils.getProps(queueConfigs, properties));
-
- client.createQueue(queueDescription);
- queues.add(topic);
- } catch (ServiceBusException | InterruptedException e) {
- if (e instanceof MessagingEntityAlreadyExistsException) {
- queues.add(topic);
- log.info("[{}] queue already exists.", topic);
- } else {
- log.error("Failed to create queue: [{}]", topic, e);
- }
- }
- }
-
- @Override
- public void deleteTopic(String topic) {
- if (queues.contains(topic)) {
- doDelete(topic);
- } else {
- try {
- if (client.getQueue(topic) != null) {
- doDelete(topic);
- } else {
- log.warn("Azure Service Bus Queue [{}] is not exist.", topic);
- }
- } catch (ServiceBusException | InterruptedException e) {
- log.error("Failed to delete Azure Service Bus queue [{}]", topic, e);
- }
- }
- }
-
- private void doDelete(String topic) {
- try {
- client.deleteTopic(topic);
- } catch (ServiceBusException | InterruptedException e) {
- log.error("Failed to delete Azure Service Bus queue [{}]", topic, e);
- }
- }
-
- private void setQueueConfigs(QueueDescription queueDescription, Map queueConfigs) {
- queueConfigs.forEach((confKey, confValue) -> {
- switch (confKey) {
- case MAX_SIZE:
- queueDescription.setMaxSizeInMB(Long.parseLong(confValue));
- break;
- case MESSAGE_TIME_TO_LIVE:
- queueDescription.setDefaultMessageTimeToLive(Duration.ofSeconds(Long.parseLong(confValue)));
- break;
- case LOCK_DURATION:
- queueDescription.setLockDuration(Duration.ofSeconds(Long.parseLong(confValue)));
- break;
- default:
- log.error("Unknown config: [{}]", confKey);
- }
- });
- }
-
- public void destroy() {
- try {
- client.close();
- } catch (IOException e) {
- log.error("Failed to close ManagementClient.");
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java
deleted file mode 100644
index 8e2b32aab5..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java
+++ /dev/null
@@ -1,175 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.azure.servicebus;
-
-import com.google.gson.Gson;
-import com.google.protobuf.InvalidProtocolBufferException;
-import com.microsoft.azure.servicebus.TransactionContext;
-import com.microsoft.azure.servicebus.primitives.ConnectionStringBuilder;
-import com.microsoft.azure.servicebus.primitives.CoreMessageReceiver;
-import com.microsoft.azure.servicebus.primitives.MessageWithDeliveryTag;
-import com.microsoft.azure.servicebus.primitives.MessagingEntityType;
-import com.microsoft.azure.servicebus.primitives.MessagingFactory;
-import com.microsoft.azure.servicebus.primitives.SettleModePair;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.qpid.proton.amqp.messaging.Data;
-import org.apache.qpid.proton.amqp.transport.ReceiverSettleMode;
-import org.apache.qpid.proton.amqp.transport.SenderSettleMode;
-import org.springframework.util.CollectionUtils;
-import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueMsg;
-import org.thingsboard.server.queue.TbQueueMsgDecoder;
-import org.thingsboard.server.queue.common.AbstractTbQueueConsumerTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueMsg;
-
-import java.time.Duration;
-import java.util.Collection;
-import java.util.Collections;
-import java.util.HashSet;
-import java.util.List;
-import java.util.Map;
-import java.util.Set;
-import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutionException;
-import java.util.stream.Collectors;
-import java.util.stream.Stream;
-
-@Slf4j
-public class TbServiceBusConsumerTemplate extends AbstractTbQueueConsumerTemplate {
- private final TbQueueAdmin admin;
- private final TbQueueMsgDecoder decoder;
- private final TbServiceBusSettings serviceBusSettings;
-
- private final Gson gson = new Gson();
-
- private Set receivers;
- private final Map> pendingMessages = new ConcurrentHashMap<>();
- private volatile int messagesPerQueue;
-
- public TbServiceBusConsumerTemplate(TbQueueAdmin admin, TbServiceBusSettings serviceBusSettings, String topic, TbQueueMsgDecoder decoder) {
- super(topic);
- this.admin = admin;
- this.decoder = decoder;
- this.serviceBusSettings = serviceBusSettings;
- }
-
- @Override
- protected List doPoll(long durationInMillis) {
- List>> messageFutures =
- receivers.stream()
- .map(receiver -> receiver
- .receiveAsync(messagesPerQueue, Duration.ofMillis(durationInMillis))
- .whenComplete((messages, err) -> {
- if (!CollectionUtils.isEmpty(messages)) {
- pendingMessages.put(receiver, messages);
- } else if (err != null) {
- log.error("Failed to receive messages.", err);
- }
- }))
- .collect(Collectors.toList());
- try {
- return fromList(messageFutures)
- .get()
- .stream()
- .flatMap(messages -> CollectionUtils.isEmpty(messages) ? Stream.empty() : messages.stream())
- .collect(Collectors.toList());
- } catch (InterruptedException | ExecutionException e) {
- if (stopped) {
- log.info("[{}] Service Bus consumer is stopped.", getTopic());
- } else {
- log.error("Failed to receive messages", e);
- }
- return Collections.emptyList();
- }
- }
-
- @Override
- protected void doSubscribe(List topicNames) {
- createReceivers();
- messagesPerQueue = receivers.size() / Math.max(partitions.size(), 1);
- }
-
- @Override
- protected void doCommit() {
- pendingMessages.forEach((receiver, msgs) ->
- msgs.forEach(msg -> receiver.completeMessageAsync(msg.getDeliveryTag(), TransactionContext.NULL_TXN)));
- pendingMessages.clear();
- }
-
- @Override
- protected void doUnsubscribe() {
- receivers.forEach(CoreMessageReceiver::closeAsync);
- }
-
- private void createReceivers() {
- List> receiverFutures = partitions.stream()
- .map(TopicPartitionInfo::getFullTopicName)
- .map(queue -> {
- MessagingFactory factory;
- try {
- factory = MessagingFactory.createFromConnectionStringBuilder(createConnection(queue));
- } catch (InterruptedException | ExecutionException e) {
- log.error("Failed to create factory for the queue [{}]", queue);
- throw new RuntimeException("Failed to create the factory", e);
- }
-
- return CoreMessageReceiver.create(factory, queue, queue, 0,
- new SettleModePair(SenderSettleMode.UNSETTLED, ReceiverSettleMode.SECOND),
- MessagingEntityType.QUEUE);
- }).collect(Collectors.toList());
-
- try {
- receivers = new HashSet<>(fromList(receiverFutures).get());
- } catch (InterruptedException | ExecutionException e) {
- if (stopped) {
- log.info("[{}] Service Bus consumer is stopped.", getTopic());
- } else {
- log.error("Failed to create receivers", e);
- }
- }
- }
-
- private ConnectionStringBuilder createConnection(String queue) {
- admin.createTopicIfNotExists(queue);
- return new ConnectionStringBuilder(
- serviceBusSettings.getNamespaceName(),
- queue,
- serviceBusSettings.getSasKeyName(),
- serviceBusSettings.getSasKey());
- }
-
- private CompletableFuture> fromList(List> futures) {
- @SuppressWarnings("unchecked")
- CompletableFuture>[] arrayFuture = new CompletableFuture[futures.size()];
- futures.toArray(arrayFuture);
-
- return CompletableFuture
- .allOf(arrayFuture)
- .thenApply(v -> futures
- .stream()
- .map(CompletableFuture::join)
- .collect(Collectors.toList()));
- }
-
- @Override
- protected T decode(MessageWithDeliveryTag data) throws InvalidProtocolBufferException {
- DefaultTbQueueMsg msg = gson.fromJson(new String(((Data) data.getMessage().getBody()).getValue().getArray()), DefaultTbQueueMsg.class);
- return decoder.decode(msg);
- }
-
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java
deleted file mode 100644
index 119f28079b..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusProducerTemplate.java
+++ /dev/null
@@ -1,110 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.azure.servicebus;
-
-import com.google.gson.Gson;
-import com.microsoft.azure.servicebus.IMessage;
-import com.microsoft.azure.servicebus.Message;
-import com.microsoft.azure.servicebus.QueueClient;
-import com.microsoft.azure.servicebus.ReceiveMode;
-import com.microsoft.azure.servicebus.primitives.ConnectionStringBuilder;
-import com.microsoft.azure.servicebus.primitives.ServiceBusException;
-import lombok.extern.slf4j.Slf4j;
-import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueCallback;
-import org.thingsboard.server.queue.TbQueueMsg;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.common.DefaultTbQueueMsg;
-
-import java.util.Map;
-import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-
-@Slf4j
-public class TbServiceBusProducerTemplate implements TbQueueProducer {
- private final String defaultTopic;
- private final Gson gson = new Gson();
- private final TbQueueAdmin admin;
- private final TbServiceBusSettings serviceBusSettings;
- private final Map clients = new ConcurrentHashMap<>();
- private final ExecutorService executorService;
-
- public TbServiceBusProducerTemplate(TbQueueAdmin admin, TbServiceBusSettings serviceBusSettings, String defaultTopic) {
- this.admin = admin;
- this.defaultTopic = defaultTopic;
- this.serviceBusSettings = serviceBusSettings;
- executorService = Executors.newCachedThreadPool();
- }
-
- @Override
- public void init() {
-
- }
-
- @Override
- public String getDefaultTopic() {
- return defaultTopic;
- }
-
- @Override
- public void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback) {
- IMessage message = new Message(gson.toJson(new DefaultTbQueueMsg(msg)));
- CompletableFuture future = getClient(tpi.getFullTopicName()).sendAsync(message);
- future.whenCompleteAsync((success, err) -> {
- if (err != null) {
- callback.onFailure(err);
- } else {
- callback.onSuccess(null);
- }
- }, executorService);
- }
-
- @Override
- public void stop() {
- clients.forEach((t, client) -> {
- try {
- client.close();
- } catch (ServiceBusException e) {
- log.error("Failed to close QueueClient.", e);
- }
- });
-
- if (executorService != null) {
- executorService.shutdownNow();
- }
- }
-
- private QueueClient getClient(String topic) {
- return clients.computeIfAbsent(topic, k -> {
- admin.createTopicIfNotExists(topic);
- ConnectionStringBuilder builder =
- new ConnectionStringBuilder(
- serviceBusSettings.getNamespaceName(),
- topic,
- serviceBusSettings.getSasKeyName(),
- serviceBusSettings.getSasKey());
- try {
- return new QueueClient(builder, ReceiveMode.PEEKLOCK);
- } catch (InterruptedException | ServiceBusException e) {
- log.error("Failed to create new client for the Queue: [{}]", topic, e);
- throw new RuntimeException("Failed to create new client for the Queue", e);
- }
- });
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusQueueConfigs.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusQueueConfigs.java
deleted file mode 100644
index 80da7846ea..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusQueueConfigs.java
+++ /dev/null
@@ -1,72 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.azure.servicebus;
-
-import jakarta.annotation.PostConstruct;
-import lombok.Getter;
-import org.springframework.beans.factory.annotation.Value;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.queue.util.PropertyUtils;
-
-import java.util.Map;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='service-bus'")
-public class TbServiceBusQueueConfigs {
-
- @Value("${queue.service-bus.queue-properties.core:}")
- private String coreProperties;
- @Value("${queue.service-bus.queue-properties.rule-engine:}")
- private String ruleEngineProperties;
- @Value("${queue.service-bus.queue-properties.transport-api:}")
- private String transportApiProperties;
- @Value("${queue.service-bus.queue-properties.notifications:}")
- private String notificationsProperties;
- @Value("${queue.service-bus.queue-properties.js-executor:}")
- private String jsExecutorProperties;
- @Value("${queue.service-bus.queue-properties.version-control:}")
- private String vcProperties;
- @Value("${queue.service-bus.queue-properties.edge:}")
- private String edgeProperties;
-
- @Getter
- private Map coreConfigs;
- @Getter
- private Map ruleEngineConfigs;
- @Getter
- private Map transportApiConfigs;
- @Getter
- private Map notificationsConfigs;
- @Getter
- private Map jsExecutorConfigs;
- @Getter
- private Map vcConfigs;
- @Getter
- private Map edgeConfigs;
-
- @PostConstruct
- private void init() {
- coreConfigs = PropertyUtils.getProps(coreProperties);
- ruleEngineConfigs = PropertyUtils.getProps(ruleEngineProperties);
- transportApiConfigs = PropertyUtils.getProps(transportApiProperties);
- notificationsConfigs = PropertyUtils.getProps(notificationsProperties);
- jsExecutorConfigs = PropertyUtils.getProps(jsExecutorProperties);
- vcConfigs = PropertyUtils.getProps(vcProperties);
- edgeConfigs = PropertyUtils.getProps(edgeProperties);
- }
-
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusSettings.java
deleted file mode 100644
index 98bb11a839..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusSettings.java
+++ /dev/null
@@ -1,37 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.azure.servicebus;
-
-import lombok.Data;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.beans.factory.annotation.Value;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.stereotype.Component;
-
-@Slf4j
-@ConditionalOnExpression("'${queue.type:null}'=='service-bus'")
-@Component
-@Data
-public class TbServiceBusSettings {
- @Value("${queue.service_bus.namespace_name}")
- private String namespaceName;
- @Value("${queue.service_bus.sas_key_name}")
- private String sasKeyName;
- @Value("${queue.service_bus.sas_key}")
- private String sasKey;
- @Value("${queue.service_bus.max_messages}")
- private int maxMessages;
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java
deleted file mode 100644
index 310cc1d1c5..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsMonolithQueueFactory.java
+++ /dev/null
@@ -1,317 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.data.queue.Queue;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsRequest;
-import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsResponse;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-import org.thingsboard.server.queue.settings.TbQueueVersionControlSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsConsumerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsProducerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='aws-sqs' && '${service.type:null}'=='monolith'")
-public class AwsSqsMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngineQueueFactory, TbVersionControlQueueFactory {
-
- private final TopicService topicService;
- private final TbQueueCoreSettings coreSettings;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbAwsSqsSettings sqsSettings;
- private final TbQueueVersionControlSettings vcSettings;
- private final TbQueueEdgeSettings edgeSettings;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin transportApiAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin otaAdmin;
- private final TbQueueAdmin vcAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public AwsSqsMonolithQueueFactory(TopicService topicService, TbQueueCoreSettings coreSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TbServiceInfoProvider serviceInfoProvider,
- TbQueueTransportApiSettings transportApiSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TbAwsSqsSettings sqsSettings,
- TbQueueVersionControlSettings vcSettings,
- TbQueueEdgeSettings edgeSettings,
- TbAwsSqsQueueAttributes sqsQueueAttributes,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings) {
- this.topicService = topicService;
- this.coreSettings = coreSettings;
- this.serviceInfoProvider = serviceInfoProvider;
- this.ruleEngineSettings = ruleEngineSettings;
- this.transportApiSettings = transportApiSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.sqsSettings = sqsSettings;
- this.vcSettings = vcSettings;
- this.edgeSettings = edgeSettings;
- this.jsInvokeSettings = jsInvokeSettings;
-
- this.coreAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getCoreAttributes());
- this.ruleEngineAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getRuleEngineAttributes());
- this.jsExecutorAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getJsExecutorAttributes());
- this.transportApiAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getTransportApiAttributes());
- this.notificationAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getNotificationsAttributes());
- this.otaAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getOtaAttributes());
- this.vcAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getVcAttributes());
- this.edgeAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getEdgeAttributes());
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(ruleEngineAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueConsumer> createToVersionControlMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(vcAdmin, sqsSettings, topicService.buildTopicName(vcSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToVersionControlServiceMsg.parseFrom(msg.getData()), msg.getHeaders())
- );
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) {
- return new TbAwsSqsConsumerTemplate<>(ruleEngineAdmin, sqsSettings, topicService.buildTopicName(configuration.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createTransportApiRequestConsumer() {
- return new TbAwsSqsConsumerTemplate<>(transportApiAdmin, sqsSettings, topicService.buildTopicName(transportApiSettings.getRequestsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createTransportApiResponseProducer() {
- return new TbAwsSqsProducerTemplate<>(transportApiAdmin, sqsSettings, topicService.buildTopicName(transportApiSettings.getResponsesTopic()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer> producer = new TbAwsSqsProducerTemplate<>(jsExecutorAdmin, sqsSettings, jsInvokeSettings.getRequestTopic());
- TbQueueConsumer> consumer = new TbAwsSqsConsumerTemplate<>(jsExecutorAdmin, sqsSettings,
- jsInvokeSettings.getResponseTopic() + "_" + serviceInfoProvider.getServiceId(),
- msg -> {
- RemoteJsResponse.Builder builder = RemoteJsResponse.newBuilder();
- JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder);
- return new TbProtoQueueMsg<>(msg.getKey(), builder.build(), msg.getHeaders());
- });
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> builder = DefaultTbQueueRequestTemplate.builder();
- builder.queueAdmin(jsExecutorAdmin);
- builder.requestTemplate(producer);
- builder.responseTemplate(consumer);
- builder.maxPendingRequests(jsInvokeSettings.getMaxPendingRequests());
- builder.maxRequestTimeout(jsInvokeSettings.getMaxRequestsTimeout());
- builder.pollInterval(jsInvokeSettings.getResponsePollInterval());
- return builder.build();
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToUsageStatsServiceMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToUsageStatsServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToOtaPackageStateServiceMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToOtaPackageStateServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createToOtaPackageStateServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()));
- }
-
- @Override
- public TbQueueProducer> createVersionControlMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(vcAdmin, sqsSettings, topicService.buildTopicName(vcSettings.getTopic()));
- }
-
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperReprocessingMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createEdgeMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(edgeAdmin, sqsSettings, topicService.buildTopicName(edgeSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(edgeAdmin, sqsSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToEdgeNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueProducer> createEdgeEventMsgProducer() {
- return null;
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- if (jsExecutorAdmin != null) {
- jsExecutorAdmin.destroy();
- }
- if (transportApiAdmin != null) {
- transportApiAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (otaAdmin != null) {
- otaAdmin.destroy();
- }
- if (vcAdmin != null) {
- vcAdmin.destroy();
- }
- if (edgeAdmin != null) {
- edgeAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java
deleted file mode 100644
index 515c1bba44..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbCoreQueueFactory.java
+++ /dev/null
@@ -1,291 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-import org.thingsboard.server.queue.settings.TbQueueVersionControlSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsConsumerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsProducerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='aws-sqs' && '${service.type:null}'=='tb-core'")
-public class AwsSqsTbCoreQueueFactory implements TbCoreQueueFactory {
-
- private final TbAwsSqsSettings sqsSettings;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TopicService topicService;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbQueueVersionControlSettings vcSettings;
- private final TbQueueEdgeSettings edgeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin transportApiAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin otaAdmin;
- private final TbQueueAdmin vcAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public AwsSqsTbCoreQueueFactory(TbAwsSqsSettings sqsSettings,
- TbQueueCoreSettings coreSettings,
- TbQueueTransportApiSettings transportApiSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TopicService topicService,
- TbQueueVersionControlSettings vcSettings,
- TbQueueEdgeSettings edgeSettings,
- TbServiceInfoProvider serviceInfoProvider,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings,
- TbAwsSqsQueueAttributes sqsQueueAttributes,
- TbQueueTransportNotificationSettings transportNotificationSettings) {
- this.sqsSettings = sqsSettings;
- this.coreSettings = coreSettings;
- this.transportApiSettings = transportApiSettings;
- this.ruleEngineSettings = ruleEngineSettings;
- this.edgeSettings = edgeSettings;
- this.topicService = topicService;
- this.serviceInfoProvider = serviceInfoProvider;
- this.jsInvokeSettings = jsInvokeSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.vcSettings = vcSettings;
-
- this.coreAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getCoreAttributes());
- this.ruleEngineAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getRuleEngineAttributes());
- this.jsExecutorAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getJsExecutorAttributes());
- this.transportApiAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getTransportApiAttributes());
- this.notificationAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getNotificationsAttributes());
- this.otaAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getOtaAttributes());
- this.vcAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getVcAttributes());
- this.edgeAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getEdgeAttributes());
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueConsumer> createToCoreMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createTransportApiRequestConsumer() {
- return new TbAwsSqsConsumerTemplate<>(transportApiAdmin, sqsSettings, topicService.buildTopicName(transportApiSettings.getRequestsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createTransportApiResponseProducer() {
- return new TbAwsSqsProducerTemplate<>(transportApiAdmin, sqsSettings, topicService.buildTopicName(transportApiSettings.getResponsesTopic()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer> producer = new TbAwsSqsProducerTemplate<>(jsExecutorAdmin, sqsSettings, jsInvokeSettings.getRequestTopic());
- TbQueueConsumer> consumer = new TbAwsSqsConsumerTemplate<>(jsExecutorAdmin, sqsSettings,
- jsInvokeSettings.getResponseTopic() + "_" + serviceInfoProvider.getServiceId(),
- msg -> {
- JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder();
- JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder);
- return new TbProtoQueueMsg<>(msg.getKey(), builder.build(), msg.getHeaders());
- });
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> builder = DefaultTbQueueRequestTemplate.builder();
- builder.queueAdmin(jsExecutorAdmin);
- builder.requestTemplate(producer);
- builder.responseTemplate(consumer);
- builder.maxPendingRequests(jsInvokeSettings.getMaxPendingRequests());
- builder.maxRequestTimeout(jsInvokeSettings.getMaxRequestsTimeout());
- builder.pollInterval(jsInvokeSettings.getResponsePollInterval());
- return builder.build();
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToUsageStatsServiceMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToUsageStatsServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToOtaPackageStateServiceMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToOtaPackageStateServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createToOtaPackageStateServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()));
- }
-
- @Override
- public TbQueueProducer> createVersionControlMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(vcAdmin, sqsSettings, topicService.buildTopicName(vcSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperReprocessingMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createEdgeMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(edgeAdmin, sqsSettings, topicService.buildTopicName(edgeSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(edgeAdmin, sqsSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToEdgeNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- if (jsExecutorAdmin != null) {
- jsExecutorAdmin.destroy();
- }
- if (transportApiAdmin != null) {
- transportApiAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (otaAdmin != null) {
- otaAdmin.destroy();
- }
- if (vcAdmin != null) {
- vcAdmin.destroy();
- }
- if (edgeAdmin != null) {
- edgeAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java
deleted file mode 100644
index a93aba2764..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbRuleEngineQueueFactory.java
+++ /dev/null
@@ -1,208 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.data.queue.Queue;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsConsumerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsProducerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='aws-sqs' && '${service.type:null}'=='tb-rule-engine'")
-public class AwsSqsTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
-
- private final TopicService topicService;
- private final TbQueueCoreSettings coreSettings;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbAwsSqsSettings sqsSettings;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbQueueEdgeSettings edgeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin otaAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public AwsSqsTbRuleEngineQueueFactory(TopicService topicService, TbQueueCoreSettings coreSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TbServiceInfoProvider serviceInfoProvider,
- TbAwsSqsSettings sqsSettings,
- TbAwsSqsQueueAttributes sqsQueueAttributes,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TbQueueEdgeSettings edgeSettings) {
- this.topicService = topicService;
- this.coreSettings = coreSettings;
- this.serviceInfoProvider = serviceInfoProvider;
- this.ruleEngineSettings = ruleEngineSettings;
- this.sqsSettings = sqsSettings;
- this.jsInvokeSettings = jsInvokeSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.edgeSettings = edgeSettings;
-
- this.coreAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getCoreAttributes());
- this.ruleEngineAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getRuleEngineAttributes());
- this.jsExecutorAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getJsExecutorAttributes());
- this.notificationAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getNotificationsAttributes());
- this.otaAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getOtaAttributes());
- this.edgeAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getEdgeAttributes());
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(ruleEngineAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(edgeAdmin, sqsSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) {
- return new TbAwsSqsConsumerTemplate<>(ruleEngineAdmin, sqsSettings, topicService.buildTopicName(configuration.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings,
- topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer> producer = new TbAwsSqsProducerTemplate<>(jsExecutorAdmin, sqsSettings, jsInvokeSettings.getRequestTopic());
- TbQueueConsumer> consumer = new TbAwsSqsConsumerTemplate<>(jsExecutorAdmin, sqsSettings,
- jsInvokeSettings.getResponseTopic() + "_" + serviceInfoProvider.getServiceId(),
- msg -> {
- JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder();
- JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder);
- return new TbProtoQueueMsg<>(msg.getKey(), builder.build(), msg.getHeaders());
- });
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> builder = DefaultTbQueueRequestTemplate.builder();
- builder.queueAdmin(jsExecutorAdmin);
- builder.requestTemplate(producer);
- builder.responseTemplate(consumer);
- builder.maxPendingRequests(jsInvokeSettings.getMaxPendingRequests());
- builder.maxRequestTimeout(jsInvokeSettings.getMaxRequestsTimeout());
- builder.pollInterval(jsInvokeSettings.getResponsePollInterval());
- return builder.build();
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueProducer> createToOtaPackageStateServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(otaAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- if (jsExecutorAdmin != null) {
- jsExecutorAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (otaAdmin != null) {
- otaAdmin.destroy();
- }
- }
-
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java
deleted file mode 100644
index 67ff0393bd..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTbVersionControlQueueFactory.java
+++ /dev/null
@@ -1,98 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueVersionControlSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsConsumerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsProducerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='aws-sqs' && '${service.type:null}'=='tb-vc-executor'")
-public class AwsSqsTbVersionControlQueueFactory implements TbVersionControlQueueFactory {
-
- private final TbAwsSqsSettings sqsSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbQueueVersionControlSettings vcSettings;
- private final TopicService topicService;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin vcAdmin;
-
- public AwsSqsTbVersionControlQueueFactory(TbAwsSqsSettings sqsSettings,
- TbQueueCoreSettings coreSettings,
- TbQueueVersionControlSettings vcSettings,
- TbAwsSqsQueueAttributes sqsQueueAttributes,
- TopicService topicService
- ) {
- this.sqsSettings = sqsSettings;
- this.coreSettings = coreSettings;
- this.vcSettings = vcSettings;
- this.topicService = topicService;
-
- this.coreAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getCoreAttributes());
- this.notificationAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getNotificationsAttributes());
- this.vcAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getVcAttributes());
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToVersionControlMsgConsumer() {
- return new TbAwsSqsConsumerTemplate<>(vcAdmin, sqsSettings, topicService.buildTopicName(vcSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportProtos.ToVersionControlServiceMsg.parseFrom(msg.getData()), msg.getHeaders())
- );
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (vcAdmin != null) {
- vcAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java
deleted file mode 100644
index 407ef86a99..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/AwsSqsTransportQueueFactory.java
+++ /dev/null
@@ -1,153 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import jakarta.annotation.PreDestroy;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.gen.transport.TransportProtos;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-import org.thingsboard.server.queue.sqs.TbAwsSqsAdmin;
-import org.thingsboard.server.queue.sqs.TbAwsSqsConsumerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsProducerTemplate;
-import org.thingsboard.server.queue.sqs.TbAwsSqsQueueAttributes;
-import org.thingsboard.server.queue.sqs.TbAwsSqsSettings;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='aws-sqs' && (('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport')")
-@Slf4j
-public class AwsSqsTransportQueueFactory implements TbTransportQueueFactory {
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbAwsSqsSettings sqsSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TopicService topicService;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin transportApiAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
-
- public AwsSqsTransportQueueFactory(TbQueueTransportApiSettings transportApiSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TbAwsSqsSettings sqsSettings,
- TbServiceInfoProvider serviceInfoProvider,
- TbQueueCoreSettings coreSettings,
- TbAwsSqsQueueAttributes sqsQueueAttributes,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TopicService topicService) {
- this.transportApiSettings = transportApiSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.sqsSettings = sqsSettings;
- this.serviceInfoProvider = serviceInfoProvider;
- this.coreSettings = coreSettings;
- this.ruleEngineSettings = ruleEngineSettings;
- this.topicService = topicService;
-
- this.coreAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getCoreAttributes());
- this.transportApiAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getTransportApiAttributes());
- this.notificationAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getNotificationsAttributes());
- this.ruleEngineAdmin = new TbAwsSqsAdmin(sqsSettings, sqsQueueAttributes.getRuleEngineAttributes());
- }
-
- @Override
- public TbQueueRequestTemplate, TbProtoQueueMsg> createTransportApiRequestTemplate() {
- TbQueueProducer> producerTemplate =
- new TbAwsSqsProducerTemplate<>(transportApiAdmin, sqsSettings, topicService.buildTopicName(transportApiSettings.getRequestsTopic()));
-
- TbQueueConsumer> consumerTemplate =
- new TbAwsSqsConsumerTemplate<>(transportApiAdmin, sqsSettings,
- topicService.buildTopicName(transportApiSettings.getResponsesTopic() + "_" + serviceInfoProvider.getServiceId()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiResponseMsg.parseFrom(msg.getData()), msg.getHeaders()));
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> templateBuilder = DefaultTbQueueRequestTemplate.builder();
- templateBuilder.queueAdmin(transportApiAdmin);
- templateBuilder.requestTemplate(producerTemplate);
- templateBuilder.responseTemplate(consumerTemplate);
- templateBuilder.maxPendingRequests(transportApiSettings.getMaxPendingRequests());
- templateBuilder.maxRequestTimeout(transportApiSettings.getMaxRequestsTimeout());
- templateBuilder.pollInterval(transportApiSettings.getResponsePollInterval());
- return templateBuilder.build();
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(ruleEngineAdmin, sqsSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createTransportNotificationsConsumer() {
- return new TbAwsSqsConsumerTemplate<>(notificationAdmin, sqsSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic() + "_" + serviceInfoProvider.getServiceId()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToTransportMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbAwsSqsProducerTemplate<>(coreAdmin, sqsSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (transportApiAdmin != null) {
- transportApiAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java
deleted file mode 100644
index 05bef819b7..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueFactory.java
+++ /dev/null
@@ -1,315 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.data.queue.Queue;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsRequest;
-import org.thingsboard.server.gen.js.JsInvokeProtos.RemoteJsResponse;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
-import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
-import org.thingsboard.server.queue.pubsub.TbPubSubSubscriptionSettings;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-import org.thingsboard.server.queue.settings.TbQueueVersionControlSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='monolith'")
-public class PubSubMonolithQueueFactory implements TbCoreQueueFactory, TbRuleEngineQueueFactory, TbVersionControlQueueFactory {
-
- private final TbPubSubSettings pubSubSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TopicService topicService;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
- private final TbQueueVersionControlSettings vcSettings;
- private final TbQueueEdgeSettings edgeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin transportApiAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin vcAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public PubSubMonolithQueueFactory(TbPubSubSettings pubSubSettings,
- TbQueueCoreSettings coreSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TbQueueTransportApiSettings transportApiSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TopicService topicService,
- TbServiceInfoProvider serviceInfoProvider,
- TbPubSubSubscriptionSettings pubSubSubscriptionSettings,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings,
- TbQueueVersionControlSettings vcSettings,
- TbQueueEdgeSettings edgeSettings) {
- this.pubSubSettings = pubSubSettings;
- this.coreSettings = coreSettings;
- this.ruleEngineSettings = ruleEngineSettings;
- this.transportApiSettings = transportApiSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.topicService = topicService;
- this.serviceInfoProvider = serviceInfoProvider;
- this.vcSettings = vcSettings;
- this.edgeSettings = edgeSettings;
-
- this.coreAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getCoreSettings());
- this.ruleEngineAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getRuleEngineSettings());
- this.jsExecutorAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getJsExecutorSettings());
- this.transportApiAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getTransportApiSettings());
- this.notificationAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getNotificationsSettings());
- this.vcAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getVcSettings());
- this.edgeAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getEdgeSettings());
-
- this.jsInvokeSettings = jsInvokeSettings;
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbPubSubProducerTemplate<>(ruleEngineAdmin, pubSubSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
-
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToVersionControlMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(vcAdmin, pubSubSettings, topicService.buildTopicName(vcSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToVersionControlServiceMsg.parseFrom(msg.getData()), msg.getHeaders())
- );
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) {
- return new TbPubSubConsumerTemplate<>(ruleEngineAdmin, pubSubSettings, topicService.buildTopicName(configuration.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createTransportApiRequestConsumer() {
- return new TbPubSubConsumerTemplate<>(transportApiAdmin, pubSubSettings, topicService.buildTopicName(transportApiSettings.getRequestsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createTransportApiResponseProducer() {
- return new TbPubSubProducerTemplate<>(transportApiAdmin, pubSubSettings, topicService.buildTopicName(transportApiSettings.getResponsesTopic()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer> producer = new TbPubSubProducerTemplate<>(jsExecutorAdmin, pubSubSettings, jsInvokeSettings.getRequestTopic());
- TbQueueConsumer> consumer = new TbPubSubConsumerTemplate<>(jsExecutorAdmin, pubSubSettings,
- jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId(),
- msg -> {
- RemoteJsResponse.Builder builder = RemoteJsResponse.newBuilder();
- JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder);
- return new TbProtoQueueMsg<>(msg.getKey(), builder.build(), msg.getHeaders());
- });
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> builder = DefaultTbQueueRequestTemplate.builder();
- builder.queueAdmin(jsExecutorAdmin);
- builder.requestTemplate(producer);
- builder.responseTemplate(consumer);
- builder.maxPendingRequests(jsInvokeSettings.getMaxPendingRequests());
- builder.maxRequestTimeout(jsInvokeSettings.getMaxRequestsTimeout());
- builder.pollInterval(jsInvokeSettings.getResponsePollInterval());
- return builder.build();
- }
-
- @Override
- public TbQueueConsumer> createToUsageStatsServiceMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToUsageStatsServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToOtaPackageStateServiceMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToOtaPackageStateServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createToOtaPackageStateServiceMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()));
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueProducer> createVersionControlMsgProducer() {
- return new TbPubSubProducerTemplate<>(vcAdmin, pubSubSettings, topicService.buildTopicName(vcSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperReprocessingMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbPubSubProducerTemplate<>(edgeAdmin, pubSubSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createEdgeMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(edgeAdmin, pubSubSettings, topicService.buildTopicName(edgeSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToEdgeNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueProducer> createEdgeEventMsgProducer() {
- return null;
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- if (jsExecutorAdmin != null) {
- jsExecutorAdmin.destroy();
- }
- if (transportApiAdmin != null) {
- transportApiAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (vcAdmin != null) {
- vcAdmin.destroy();
- }
- if (edgeAdmin != null) {
- edgeAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java
deleted file mode 100644
index 919d895a97..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueFactory.java
+++ /dev/null
@@ -1,278 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToVersionControlServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
-import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
-import org.thingsboard.server.queue.pubsub.TbPubSubSubscriptionSettings;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='tb-core'")
-public class PubSubTbCoreQueueFactory implements TbCoreQueueFactory {
-
- private final TbPubSubSettings pubSubSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TopicService topicService;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbQueueEdgeSettings edgeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin transportApiAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public PubSubTbCoreQueueFactory(TbPubSubSettings pubSubSettings,
- TbQueueCoreSettings coreSettings,
- TbQueueTransportApiSettings transportApiSettings,
- TopicService topicService,
- TbServiceInfoProvider serviceInfoProvider,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TbQueueEdgeSettings edgeSettings,
- TbPubSubSubscriptionSettings pubSubSubscriptionSettings) {
- this.pubSubSettings = pubSubSettings;
- this.coreSettings = coreSettings;
- this.transportApiSettings = transportApiSettings;
- this.topicService = topicService;
- this.serviceInfoProvider = serviceInfoProvider;
- this.jsInvokeSettings = jsInvokeSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.ruleEngineSettings = ruleEngineSettings;
- this.edgeSettings = edgeSettings;
-
- this.coreAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getCoreSettings());
- this.jsExecutorAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getJsExecutorSettings());
- this.transportApiAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getTransportApiSettings());
- this.notificationAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getNotificationsSettings());
- this.ruleEngineAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getRuleEngineSettings());
- this.edgeAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getEdgeSettings());
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueConsumer> createToCoreMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToCoreNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createTransportApiRequestConsumer() {
- return new TbPubSubConsumerTemplate<>(transportApiAdmin, pubSubSettings, topicService.buildTopicName(transportApiSettings.getRequestsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createTransportApiResponseProducer() {
- return new TbPubSubProducerTemplate<>(transportApiAdmin, pubSubSettings, topicService.buildTopicName(transportApiSettings.getResponsesTopic()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer> producer = new TbPubSubProducerTemplate<>(jsExecutorAdmin, pubSubSettings, jsInvokeSettings.getRequestTopic());
- TbQueueConsumer> consumer = new TbPubSubConsumerTemplate<>(jsExecutorAdmin, pubSubSettings,
- jsInvokeSettings.getResponseTopic() + "." + serviceInfoProvider.getServiceId(),
- msg -> {
- JsInvokeProtos.RemoteJsResponse.Builder builder = JsInvokeProtos.RemoteJsResponse.newBuilder();
- JsonFormat.parser().ignoringUnknownFields().merge(new String(msg.getData(), StandardCharsets.UTF_8), builder);
- return new TbProtoQueueMsg<>(msg.getKey(), builder.build(), msg.getHeaders());
- });
-
- DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
- , TbProtoQueueMsg> builder = DefaultTbQueueRequestTemplate.builder();
- builder.queueAdmin(jsExecutorAdmin);
- builder.requestTemplate(producer);
- builder.responseTemplate(consumer);
- builder.maxPendingRequests(jsInvokeSettings.getMaxPendingRequests());
- builder.maxRequestTimeout(jsInvokeSettings.getMaxRequestsTimeout());
- builder.pollInterval(jsInvokeSettings.getResponsePollInterval());
- return builder.build();
- }
-
- @Override
- public TbQueueConsumer> createToUsageStatsServiceMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToUsageStatsServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToOtaPackageStateServiceMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToOtaPackageStateServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createToOtaPackageStateServiceMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getOtaPackageTopic()));
- }
-
- @Override
- public TbQueueProducer> createToUsageStatsServiceMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getUsageStatsTopic()));
- }
-
- @Override
- public TbQueueProducer> createVersionControlMsgProducer() {
- //TODO: version-control
- return null;
- }
-
- @Override
- public TbQueueProducer> createHousekeeperMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createHousekeeperReprocessingMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()));
- }
-
- @Override
- public TbQueueConsumer> createHousekeeperReprocessingMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getHousekeeperReprocessingTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToHousekeeperServiceMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createEdgeMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(edgeAdmin, pubSubSettings, topicService.buildTopicName(edgeSettings.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbPubSubProducerTemplate<>(edgeAdmin, pubSubSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueConsumer> createToEdgeNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToEdgeNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @PreDestroy
- private void destroy() {
- if (coreAdmin != null) {
- coreAdmin.destroy();
- }
- if (jsExecutorAdmin != null) {
- jsExecutorAdmin.destroy();
- }
- if (transportApiAdmin != null) {
- transportApiAdmin.destroy();
- }
- if (notificationAdmin != null) {
- notificationAdmin.destroy();
- }
- if (ruleEngineAdmin != null) {
- ruleEngineAdmin.destroy();
- }
- if (edgeAdmin != null) {
- edgeAdmin.destroy();
- }
- }
-}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java
deleted file mode 100644
index 671da15f20..0000000000
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueFactory.java
+++ /dev/null
@@ -1,209 +0,0 @@
-/**
- * Copyright © 2016-2024 The Thingsboard Authors
- *
- * Licensed 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.thingsboard.server.queue.provider;
-
-import com.google.protobuf.util.JsonFormat;
-import jakarta.annotation.PreDestroy;
-import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
-import org.springframework.context.annotation.Bean;
-import org.springframework.stereotype.Component;
-import org.thingsboard.server.common.data.queue.Queue;
-import org.thingsboard.server.common.msg.queue.ServiceType;
-import org.thingsboard.server.gen.js.JsInvokeProtos;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToOtaPackageStateServiceMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
-import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
-import org.thingsboard.server.queue.TbQueueAdmin;
-import org.thingsboard.server.queue.TbQueueConsumer;
-import org.thingsboard.server.queue.TbQueueProducer;
-import org.thingsboard.server.queue.TbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
-import org.thingsboard.server.queue.common.TbProtoJsQueueMsg;
-import org.thingsboard.server.queue.common.TbProtoQueueMsg;
-import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
-import org.thingsboard.server.queue.discovery.TopicService;
-import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
-import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
-import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
-import org.thingsboard.server.queue.pubsub.TbPubSubSubscriptionSettings;
-import org.thingsboard.server.queue.settings.TbQueueCoreSettings;
-import org.thingsboard.server.queue.settings.TbQueueEdgeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRemoteJsInvokeSettings;
-import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbQueueTransportNotificationSettings;
-
-import java.nio.charset.StandardCharsets;
-
-@Component
-@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='tb-rule-engine'")
-public class PubSubTbRuleEngineQueueFactory implements TbRuleEngineQueueFactory {
-
- private final TbPubSubSettings pubSubSettings;
- private final TbQueueCoreSettings coreSettings;
- private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TopicService topicService;
- private final TbServiceInfoProvider serviceInfoProvider;
- private final TbQueueRemoteJsInvokeSettings jsInvokeSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
- private final TbQueueEdgeSettings edgeSettings;
-
- private final TbQueueAdmin coreAdmin;
- private final TbQueueAdmin ruleEngineAdmin;
- private final TbQueueAdmin jsExecutorAdmin;
- private final TbQueueAdmin notificationAdmin;
- private final TbQueueAdmin edgeAdmin;
-
- public PubSubTbRuleEngineQueueFactory(TbPubSubSettings pubSubSettings,
- TbQueueCoreSettings coreSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TopicService topicService,
- TbServiceInfoProvider serviceInfoProvider,
- TbQueueRemoteJsInvokeSettings jsInvokeSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings,
- TbPubSubSubscriptionSettings pubSubSubscriptionSettings,
- TbQueueEdgeSettings edgeSettings) {
- this.pubSubSettings = pubSubSettings;
- this.coreSettings = coreSettings;
- this.ruleEngineSettings = ruleEngineSettings;
- this.topicService = topicService;
- this.serviceInfoProvider = serviceInfoProvider;
- this.jsInvokeSettings = jsInvokeSettings;
- this.transportNotificationSettings = transportNotificationSettings;
- this.edgeSettings = edgeSettings;
-
- this.coreAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getCoreSettings());
- this.ruleEngineAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getRuleEngineSettings());
- this.jsExecutorAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getJsExecutorSettings());
- this.notificationAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getNotificationsSettings());
- this.edgeAdmin = new TbPubSubAdmin(pubSubSettings, pubSubSubscriptionSettings.getEdgeSettings());
- }
-
- @Override
- public TbQueueProducer> createTransportNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(transportNotificationSettings.getNotificationsTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineMsgProducer() {
- return new TbPubSubProducerTemplate<>(ruleEngineAdmin, pubSubSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createRuleEngineNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(ruleEngineSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreMsgProducer() {
- return new TbPubSubProducerTemplate<>(coreAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createTbCoreNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.buildTopicName(coreSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createEdgeMsgProducer() {
- return new TbPubSubProducerTemplate<>(edgeAdmin, pubSubSettings, topicService.buildTopicName(edgeSettings.getTopic()));
- }
-
- @Override
- public TbQueueProducer> createEdgeNotificationsMsgProducer() {
- return new TbPubSubProducerTemplate<>(notificationAdmin, pubSubSettings, topicService.getEdgeNotificationsTopic(serviceInfoProvider.getServiceId()).getFullTopicName());
- }
-
- @Override
- public TbQueueProducer> createEdgeEventMsgProducer() {
- return null;
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineMsgConsumer(Queue configuration) {
- return new TbPubSubConsumerTemplate<>(ruleEngineAdmin, pubSubSettings, topicService.buildTopicName(configuration.getTopic()),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- public TbQueueConsumer> createToRuleEngineNotificationsMsgConsumer() {
- return new TbPubSubConsumerTemplate<>(notificationAdmin, pubSubSettings,
- topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
- msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
- }
-
- @Override
- @Bean
- public TbQueueRequestTemplate, TbProtoQueueMsg> createRemoteJsRequestTemplate() {
- TbQueueProducer