|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 21 KiB After Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
@ -0,0 +1,37 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.config.mqtt; |
||||
|
|
||||
|
import jakarta.validation.constraints.PositiveOrZero; |
||||
|
import lombok.Data; |
||||
|
import org.springframework.boot.context.properties.ConfigurationProperties; |
||||
|
import org.springframework.context.annotation.Configuration; |
||||
|
import org.springframework.validation.annotation.Validated; |
||||
|
|
||||
|
@Data |
||||
|
@Validated |
||||
|
@Configuration |
||||
|
@ConfigurationProperties(prefix = "mqtt.client.retransmission") |
||||
|
public class MqttClientRetransmissionSettingsComponent { |
||||
|
|
||||
|
@PositiveOrZero |
||||
|
private int maxAttempts; |
||||
|
@PositiveOrZero |
||||
|
private long initialDelayMillis; |
||||
|
@PositiveOrZero |
||||
|
private double jitterFactor; |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,47 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.config.mqtt; |
||||
|
|
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.ToString; |
||||
|
import org.springframework.context.annotation.Configuration; |
||||
|
import org.thingsboard.rule.engine.api.MqttClientSettings; |
||||
|
|
||||
|
@ToString |
||||
|
@EqualsAndHashCode |
||||
|
@Configuration |
||||
|
@RequiredArgsConstructor |
||||
|
public class MqttClientSettingsComponent implements MqttClientSettings { |
||||
|
|
||||
|
private final MqttClientRetransmissionSettingsComponent retransmissionSettingsComponent; |
||||
|
|
||||
|
@Override |
||||
|
public int getRetransmissionMaxAttempts() { |
||||
|
return retransmissionSettingsComponent.getMaxAttempts(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public long getRetransmissionInitialDelayMillis() { |
||||
|
return retransmissionSettingsComponent.getInitialDelayMillis(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public double getRetransmissionJitterFactor() { |
||||
|
return retransmissionSettingsComponent.getJitterFactor(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,46 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.service.security.auth; |
||||
|
|
||||
|
import jakarta.servlet.FilterChain; |
||||
|
import jakarta.servlet.http.HttpServletRequest; |
||||
|
import jakarta.servlet.http.HttpServletResponse; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.security.core.AuthenticationException; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.springframework.web.filter.OncePerRequestFilter; |
||||
|
import org.thingsboard.server.exception.ThingsboardErrorResponseHandler; |
||||
|
|
||||
|
@Component |
||||
|
@RequiredArgsConstructor |
||||
|
@Slf4j |
||||
|
public class AuthExceptionHandler extends OncePerRequestFilter { |
||||
|
|
||||
|
private final ThingsboardErrorResponseHandler errorResponseHandler; |
||||
|
|
||||
|
@Override |
||||
|
protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain filterChain) { |
||||
|
try { |
||||
|
filterChain.doFilter(request, response); |
||||
|
} catch (AuthenticationException e) { |
||||
|
throw e; |
||||
|
} catch (Exception e) { |
||||
|
errorResponseHandler.handle(e, response); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,130 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.service.subscription; |
||||
|
|
||||
|
import org.apache.commons.lang3.RandomStringUtils; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.testcontainers.shaded.org.apache.commons.lang3.RandomUtils; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.LongDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
||||
|
import org.thingsboard.server.common.data.page.PageData; |
||||
|
import org.thingsboard.server.common.data.query.EntityData; |
||||
|
import org.thingsboard.server.common.data.query.EntityKeyType; |
||||
|
import org.thingsboard.server.common.data.query.TsValue; |
||||
|
import org.thingsboard.server.service.ws.WebSocketService; |
||||
|
import org.thingsboard.server.service.ws.WebSocketSessionRef; |
||||
|
import org.thingsboard.server.service.ws.telemetry.cmd.v2.CmdUpdate; |
||||
|
import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate; |
||||
|
import org.thingsboard.server.service.ws.telemetry.sub.TelemetrySubscriptionUpdate; |
||||
|
|
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.BDDMockito.then; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class TbEntityDataSubCtxTest { |
||||
|
|
||||
|
private final DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
||||
|
|
||||
|
private final Integer cmdId = RandomUtils.nextInt(); |
||||
|
private final Integer subscriptionId = RandomUtils.nextInt(); |
||||
|
private final String serviceId = RandomStringUtils.randomAlphanumeric(10); |
||||
|
private final String sessionId = RandomStringUtils.randomAlphanumeric(10); |
||||
|
|
||||
|
private final int maxEntitiesPerDataSubscription = 100; |
||||
|
|
||||
|
private TbEntityDataSubCtx subCtx; |
||||
|
@Mock |
||||
|
private WebSocketService webSocketService; |
||||
|
@Mock |
||||
|
private WebSocketSessionRef webSocketSessionRef; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() { |
||||
|
when(webSocketSessionRef.getSessionId()).thenReturn(sessionId); |
||||
|
subCtx = new TbEntityDataSubCtx(serviceId, webSocketService, mock(), mock(), mock(), mock(), webSocketSessionRef, cmdId, maxEntitiesPerDataSubscription); |
||||
|
|
||||
|
Map<Integer, EntityId> subToEntityIdMap = new HashMap<>(); |
||||
|
subToEntityIdMap.put(subscriptionId, deviceId); |
||||
|
ReflectionTestUtils.setField(subCtx, "subToEntityIdMap", subToEntityIdMap); |
||||
|
|
||||
|
long now = System.currentTimeMillis(); |
||||
|
long oldTs = now - 1_000_000; |
||||
|
|
||||
|
Map<String, TsValue> latestCtxValues = new HashMap<>(); |
||||
|
latestCtxValues.put("key", new TsValue(oldTs, "15")); |
||||
|
Map<EntityKeyType, Map<String, TsValue>> latest = new HashMap<>(); |
||||
|
latest.put(EntityKeyType.TIME_SERIES, latestCtxValues); |
||||
|
|
||||
|
EntityData entityData = new EntityData(); |
||||
|
entityData.setEntityId(deviceId); |
||||
|
entityData.setLatest(latest); |
||||
|
|
||||
|
PageData<EntityData> data = new PageData<>(List.of(entityData), 1, 1, true); |
||||
|
ReflectionTestUtils.setField(subCtx, "data", data); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testSendLatestWsMsg() { |
||||
|
long ts = System.currentTimeMillis(); |
||||
|
List<TsKvEntry> telemetry = List.of( |
||||
|
new BasicTsKvEntry(ts - 50000, new LongDataEntry("key", 42L), 34L), |
||||
|
new BasicTsKvEntry(ts - 20000, new LongDataEntry("key", 17L), 78L) |
||||
|
); |
||||
|
|
||||
|
TelemetrySubscriptionUpdate subUpdate = new TelemetrySubscriptionUpdate(subscriptionId, telemetry); |
||||
|
|
||||
|
subCtx.sendWsMsg(sessionId, subUpdate, EntityKeyType.TIME_SERIES, true); |
||||
|
|
||||
|
Map<EntityKeyType, Map<String, TsValue>> expectedLatest = new HashMap<>(); |
||||
|
Map<String, TsValue> expectedLatestCtxValues = new HashMap<>(); |
||||
|
expectedLatestCtxValues.put("key", new TsValue(ts - 20000, "17")); // use latest telemetry
|
||||
|
expectedLatest.put(EntityKeyType.TIME_SERIES, expectedLatestCtxValues); |
||||
|
|
||||
|
EntityData expectedEntityData = new EntityData(); |
||||
|
expectedEntityData.setEntityId(deviceId); |
||||
|
expectedEntityData.setLatest(expectedLatest); |
||||
|
|
||||
|
List<EntityData> expected = List.of(expectedEntityData); |
||||
|
|
||||
|
ArgumentCaptor<CmdUpdate> cmdUpdateCaptor = ArgumentCaptor.forClass(CmdUpdate.class); |
||||
|
then(webSocketService).should().sendUpdate(eq(sessionId), cmdUpdateCaptor.capture()); |
||||
|
CmdUpdate cmdUpdate = cmdUpdateCaptor.getValue(); |
||||
|
assertThat(cmdUpdate).isInstanceOf(EntityDataUpdate.class); |
||||
|
EntityDataUpdate entityDataUpdate = (EntityDataUpdate) cmdUpdate; |
||||
|
assertThat(entityDataUpdate.getCmdId()).isEqualTo(cmdId); |
||||
|
assertThat(entityDataUpdate.getData()).isNull(); |
||||
|
assertThat(entityDataUpdate.getUpdate()).isEqualTo(expected); |
||||
|
assertThat(entityDataUpdate.getAllowedEntities()).isEqualTo(maxEntitiesPerDataSubscription); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,18 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.common.data.edqs; |
||||
|
|
||||
|
public interface EdqsObjectKey {} |
||||
@ -0,0 +1,72 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.common.data.edqs; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; |
||||
|
import lombok.Getter; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
import lombok.Setter; |
||||
|
import org.apache.commons.lang3.BooleanUtils; |
||||
|
|
||||
|
@Getter |
||||
|
@NoArgsConstructor |
||||
|
@JsonIgnoreProperties(ignoreUnknown = true) |
||||
|
public class EdqsState { |
||||
|
|
||||
|
private Boolean edqsReady; |
||||
|
@Setter |
||||
|
private EdqsSyncStatus syncStatus; |
||||
|
@Setter |
||||
|
private EdqsApiMode apiMode; |
||||
|
|
||||
|
public boolean setEdqsReady(boolean ready) { |
||||
|
boolean changed = BooleanUtils.toBooleanDefaultIfNull(this.edqsReady, false) != ready; |
||||
|
this.edqsReady = ready; |
||||
|
return changed; |
||||
|
} |
||||
|
|
||||
|
public boolean isApiReady() { |
||||
|
return edqsReady && syncStatus == EdqsSyncStatus.FINISHED; |
||||
|
} |
||||
|
|
||||
|
public boolean isApiEnabled() { |
||||
|
return apiMode != null && (apiMode == EdqsApiMode.ENABLED || apiMode == EdqsApiMode.AUTO_ENABLED); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public String toString() { |
||||
|
return '[' + |
||||
|
"EDQS ready: " + edqsReady + |
||||
|
", sync status: " + syncStatus + |
||||
|
", API mode: " + apiMode + |
||||
|
']'; |
||||
|
} |
||||
|
|
||||
|
public enum EdqsSyncStatus { |
||||
|
REQUESTED, |
||||
|
STARTED, |
||||
|
FINISHED, |
||||
|
FAILED |
||||
|
} |
||||
|
|
||||
|
public enum EdqsApiMode { |
||||
|
ENABLED, |
||||
|
AUTO_ENABLED, |
||||
|
DISABLED, |
||||
|
AUTO_DISABLED |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,30 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.edqs.util; |
||||
|
|
||||
|
import org.thingsboard.server.common.data.ObjectType; |
||||
|
import org.thingsboard.server.common.data.edqs.EdqsObject; |
||||
|
import org.thingsboard.server.common.data.edqs.EdqsObjectKey; |
||||
|
|
||||
|
public interface EdqsMapper { |
||||
|
|
||||
|
<T extends EdqsObject> byte[] serialize(T value); |
||||
|
|
||||
|
EdqsObject deserialize(ObjectType type, byte[] bytes, boolean onlyKey); |
||||
|
|
||||
|
<T extends EdqsObject> EdqsObjectKey getKey(T object); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,164 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.common; |
||||
|
|
||||
|
import lombok.Builder; |
||||
|
import lombok.Getter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
||||
|
import org.thingsboard.server.common.stats.MessagesStats; |
||||
|
import org.thingsboard.server.queue.TbQueueConsumer; |
||||
|
import org.thingsboard.server.queue.TbQueueHandler; |
||||
|
import org.thingsboard.server.queue.TbQueueMsg; |
||||
|
import org.thingsboard.server.queue.TbQueueProducer; |
||||
|
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
import java.util.concurrent.TimeoutException; |
||||
|
import java.util.concurrent.atomic.AtomicInteger; |
||||
|
import java.util.function.Function; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class PartitionedQueueResponseTemplate<Request extends TbQueueMsg, Response extends TbQueueMsg> extends AbstractTbQueueTemplate { |
||||
|
|
||||
|
@Getter |
||||
|
private final PartitionedQueueConsumerManager<Request> requestConsumer; |
||||
|
private final TbQueueProducer<Response> responseProducer; |
||||
|
|
||||
|
private final TbQueueHandler<Request, Response> handler; |
||||
|
private final long pollInterval; |
||||
|
private final int maxPendingRequests; |
||||
|
private final long requestTimeout; |
||||
|
private final MessagesStats stats; |
||||
|
|
||||
|
private final ScheduledExecutorService scheduler; |
||||
|
private final ExecutorService callbackExecutor; |
||||
|
|
||||
|
private final AtomicInteger pendingRequestCount = new AtomicInteger(); |
||||
|
|
||||
|
@Builder |
||||
|
public PartitionedQueueResponseTemplate(String key, |
||||
|
TbQueueHandler<Request, Response> handler, |
||||
|
String requestsTopic, |
||||
|
Function<TopicPartitionInfo, TbQueueConsumer<Request>> consumerCreator, |
||||
|
TbQueueProducer<Response> responseProducer, |
||||
|
long pollInterval, |
||||
|
long requestTimeout, |
||||
|
int maxPendingRequests, |
||||
|
ExecutorService consumerExecutor, |
||||
|
ExecutorService callbackExecutor, |
||||
|
ExecutorService consumerTaskExecutor, |
||||
|
MessagesStats stats) { |
||||
|
this.scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor(key + "-queue-response-template-scheduler"); |
||||
|
this.callbackExecutor = callbackExecutor; |
||||
|
this.handler = handler; |
||||
|
this.requestConsumer = PartitionedQueueConsumerManager.<Request>create() |
||||
|
.queueKey(key + "-requests") |
||||
|
.topic(requestsTopic) |
||||
|
.pollInterval(pollInterval) |
||||
|
.msgPackProcessor((requests, consumer, config) -> processRequests(requests, consumer)) |
||||
|
.consumerCreator((config, tpi) -> consumerCreator.apply(tpi)) |
||||
|
.consumerExecutor(consumerExecutor) |
||||
|
.scheduler(scheduler) |
||||
|
.taskExecutor(consumerTaskExecutor) |
||||
|
.build(); |
||||
|
this.responseProducer = responseProducer; |
||||
|
this.pollInterval = pollInterval; |
||||
|
this.maxPendingRequests = maxPendingRequests; |
||||
|
this.requestTimeout = requestTimeout; |
||||
|
this.stats = stats; |
||||
|
} |
||||
|
|
||||
|
private void processRequests(List<Request> requests, TbQueueConsumer<Request> consumer) { |
||||
|
while (pendingRequestCount.get() >= maxPendingRequests) { |
||||
|
try { |
||||
|
Thread.sleep(pollInterval); |
||||
|
} catch (InterruptedException e) { |
||||
|
log.trace("Failed to wait until the server has capacity to handle new requests", e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
requests.forEach(request -> { |
||||
|
long currentTime = System.currentTimeMillis(); |
||||
|
long expireTs = bytesToLong(request.getHeaders().get(EXPIRE_TS_HEADER)); |
||||
|
if (expireTs >= currentTime) { |
||||
|
byte[] requestIdHeader = request.getHeaders().get(REQUEST_ID_HEADER); |
||||
|
if (requestIdHeader == null) { |
||||
|
log.error("[{}] Missing requestId in header", request); |
||||
|
return; |
||||
|
} |
||||
|
byte[] responseTopicHeader = request.getHeaders().get(RESPONSE_TOPIC_HEADER); |
||||
|
if (responseTopicHeader == null) { |
||||
|
log.error("[{}] Missing response topic in header", request); |
||||
|
return; |
||||
|
} |
||||
|
UUID requestId = bytesToUuid(requestIdHeader); |
||||
|
String responseTopic = bytesToString(responseTopicHeader); |
||||
|
try { |
||||
|
pendingRequestCount.getAndIncrement(); |
||||
|
stats.incrementTotal(); |
||||
|
AsyncCallbackTemplate.withCallbackAndTimeout(handler.handle(request), |
||||
|
response -> { |
||||
|
pendingRequestCount.decrementAndGet(); |
||||
|
response.getHeaders().put(REQUEST_ID_HEADER, uuidToBytes(requestId)); |
||||
|
responseProducer.send(TopicPartitionInfo.builder().topic(responseTopic).build(), response, null); |
||||
|
stats.incrementSuccessful(); |
||||
|
}, |
||||
|
e -> { |
||||
|
pendingRequestCount.decrementAndGet(); |
||||
|
if (e.getCause() != null && e.getCause() instanceof TimeoutException) { |
||||
|
log.warn("[{}] Timeout to process the request: {}", requestId, request, e); |
||||
|
} else { |
||||
|
log.trace("[{}] Failed to process the request: {}", requestId, request, e); |
||||
|
} |
||||
|
stats.incrementFailed(); |
||||
|
}, |
||||
|
requestTimeout, |
||||
|
scheduler, |
||||
|
callbackExecutor); |
||||
|
} catch (Throwable e) { |
||||
|
pendingRequestCount.decrementAndGet(); |
||||
|
log.warn("[{}] Failed to process the request: {}", requestId, request, e); |
||||
|
stats.incrementFailed(); |
||||
|
} |
||||
|
} |
||||
|
}); |
||||
|
consumer.commit(); |
||||
|
} |
||||
|
|
||||
|
public void subscribe(Set<TopicPartitionInfo> partitions) { |
||||
|
requestConsumer.update(partitions); |
||||
|
} |
||||
|
|
||||
|
public void stop() { |
||||
|
if (requestConsumer != null) { |
||||
|
requestConsumer.stop(); |
||||
|
requestConsumer.awaitStop(); |
||||
|
} |
||||
|
if (responseProducer != null) { |
||||
|
responseProducer.stop(); |
||||
|
} |
||||
|
if (scheduler != null) { |
||||
|
scheduler.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,70 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 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.edqs; |
||||
|
|
||||
|
import com.google.common.util.concurrent.ListeningExecutorService; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import jakarta.annotation.PostConstruct; |
||||
|
import jakarta.annotation.PreDestroy; |
||||
|
import lombok.Getter; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import org.springframework.context.annotation.Lazy; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
||||
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
||||
|
|
||||
|
import java.util.concurrent.ExecutorService; |
||||
|
import java.util.concurrent.Executors; |
||||
|
import java.util.concurrent.ScheduledExecutorService; |
||||
|
|
||||
|
@Lazy |
||||
|
@Component |
||||
|
@Getter |
||||
|
@RequiredArgsConstructor |
||||
|
public class EdqsExecutors { |
||||
|
|
||||
|
private final EdqsConfig edqsConfig; |
||||
|
|
||||
|
private ExecutorService consumersExecutor; |
||||
|
private ExecutorService consumerTaskExecutor; |
||||
|
private ScheduledExecutorService scheduler; |
||||
|
private ListeningExecutorService requestExecutor; |
||||
|
|
||||
|
@PostConstruct |
||||
|
private void init() { |
||||
|
consumersExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("edqs-consumer")); |
||||
|
consumerTaskExecutor = ThingsBoardExecutors.newWorkStealingPool(4, "edqs-consumer-task-executor"); |
||||
|
scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("edqs-scheduler"); |
||||
|
requestExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(edqsConfig.getRequestExecutorSize(), "edqs-requests")); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
private void destroy() { |
||||
|
if (consumersExecutor != null) { |
||||
|
consumersExecutor.shutdownNow(); |
||||
|
} |
||||
|
if (consumerTaskExecutor != null) { |
||||
|
consumerTaskExecutor.shutdownNow(); |
||||
|
} |
||||
|
if (scheduler != null) { |
||||
|
scheduler.shutdownNow(); |
||||
|
} |
||||
|
if (requestExecutor != null) { |
||||
|
requestExecutor.shutdownNow(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||