committed by
GitHub
106 changed files with 2021 additions and 773 deletions
@ -1,222 +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.service.edge; |
|
||||
|
|
||||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|
||||
import jakarta.annotation.PostConstruct; |
|
||||
import jakarta.annotation.PreDestroy; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.beans.factory.annotation.Autowired; |
|
||||
import org.springframework.beans.factory.annotation.Value; |
|
||||
import org.springframework.context.ApplicationEventPublisher; |
|
||||
import org.springframework.stereotype.Service; |
|
||||
import org.thingsboard.common.util.JacksonUtil; |
|
||||
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
||||
import org.thingsboard.server.common.data.audit.ActionType; |
|
||||
import org.thingsboard.server.common.data.edge.Edge; |
|
||||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|
||||
import org.thingsboard.server.common.data.id.RuleChainId; |
|
||||
import org.thingsboard.server.common.data.id.TenantId; |
|
||||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|
||||
import org.thingsboard.server.dao.edge.EdgeService; |
|
||||
import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; |
|
||||
import org.thingsboard.server.gen.transport.TransportProtos; |
|
||||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.alarm.AlarmEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.asset.AssetEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.asset.profile.AssetProfileEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.customer.CustomerEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.dashboard.DashboardEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.device.DeviceEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.device.profile.DeviceProfileEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.edge.EdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.entityview.EntityViewEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.notification.NotificationEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.oauth2.OAuth2EdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.ota.OtaPackageEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.queue.QueueEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.relation.RelationEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.rule.RuleChainEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.tenant.TenantEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.tenant.TenantProfileEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.user.UserEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.widget.WidgetBundleEdgeProcessor; |
|
||||
import org.thingsboard.server.service.edge.rpc.processor.widget.WidgetTypeEdgeProcessor; |
|
||||
|
|
||||
import java.util.UUID; |
|
||||
import java.util.concurrent.ExecutorService; |
|
||||
import java.util.concurrent.TimeUnit; |
|
||||
|
|
||||
@Service |
|
||||
@TbCoreComponent |
|
||||
@Slf4j |
|
||||
public class DefaultEdgeNotificationService implements EdgeNotificationService { |
|
||||
|
|
||||
public static final String EDGE_IS_ROOT_BODY_KEY = "isRoot"; |
|
||||
|
|
||||
@Autowired |
|
||||
private EdgeService edgeService; |
|
||||
|
|
||||
@Autowired |
|
||||
private EdgeProcessor edgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private AssetEdgeProcessor assetProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private AssetProfileEdgeProcessor assetProfileEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private DeviceEdgeProcessor deviceProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private DeviceProfileEdgeProcessor deviceProfileEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private EntityViewEdgeProcessor entityViewProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private DashboardEdgeProcessor dashboardProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private RuleChainEdgeProcessor ruleChainProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private UserEdgeProcessor userProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private CustomerEdgeProcessor customerProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private OtaPackageEdgeProcessor otaPackageProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private WidgetBundleEdgeProcessor widgetBundleProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private WidgetTypeEdgeProcessor widgetTypeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private QueueEdgeProcessor queueProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private TenantEdgeProcessor tenantEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private TenantProfileEdgeProcessor tenantProfileEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private AlarmEdgeProcessor alarmProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private RelationEdgeProcessor relationProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private ResourceEdgeProcessor resourceEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private NotificationEdgeProcessor notificationEdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
private OAuth2EdgeProcessor oAuth2EdgeProcessor; |
|
||||
|
|
||||
@Autowired |
|
||||
protected ApplicationEventPublisher eventPublisher; |
|
||||
|
|
||||
@Value("${actors.system.edge_dispatcher_pool_size:4}") |
|
||||
private int edgeDispatcherSize; |
|
||||
|
|
||||
private ExecutorService executor; |
|
||||
|
|
||||
@PostConstruct |
|
||||
public void initExecutor() { |
|
||||
executor = ThingsBoardExecutors.newWorkStealingPool(edgeDispatcherSize, "edge-notifications"); |
|
||||
} |
|
||||
|
|
||||
@PreDestroy |
|
||||
public void shutdownExecutor() { |
|
||||
if (executor != null) { |
|
||||
executor.shutdownNow(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) { |
|
||||
edge.setRootRuleChainId(ruleChainId); |
|
||||
Edge savedEdge = edgeService.saveEdge(edge); |
|
||||
ObjectNode isRootBody = JacksonUtil.newObjectNode(); |
|
||||
isRootBody.put(EDGE_IS_ROOT_BODY_KEY, Boolean.TRUE); |
|
||||
eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).edgeId(edge.getId()).entityId(ruleChainId) |
|
||||
.body(JacksonUtil.toString(isRootBody)).actionType(ActionType.UPDATED).build()); |
|
||||
return savedEdge; |
|
||||
} |
|
||||
|
|
||||
@Override |
|
||||
public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) { |
|
||||
TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB())); |
|
||||
log.debug("[{}] Pushing notification to edge {}", tenantId, edgeNotificationMsg); |
|
||||
final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(60); |
|
||||
try { |
|
||||
executor.submit(() -> { |
|
||||
try { |
|
||||
if (deadline < System.nanoTime()) { |
|
||||
log.warn("[{}] Skipping notification message because deadline reached {}", tenantId, edgeNotificationMsg); |
|
||||
return; |
|
||||
} |
|
||||
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); |
|
||||
switch (type) { |
|
||||
case EDGE -> edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg); |
|
||||
case ASSET -> assetProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case ASSET_PROFILE -> assetProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case DEVICE -> deviceProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case DEVICE_PROFILE -> deviceProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case ENTITY_VIEW -> entityViewProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case DASHBOARD -> dashboardProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case RULE_CHAIN -> ruleChainProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case USER -> userProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case CUSTOMER -> customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg); |
|
||||
case OTA_PACKAGE -> otaPackageProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case WIDGETS_BUNDLE -> widgetBundleProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case WIDGET_TYPE -> widgetTypeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case QUEUE -> queueProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case ALARM -> alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg); |
|
||||
case ALARM_COMMENT -> alarmProcessor.processAlarmCommentNotification(tenantId, edgeNotificationMsg); |
|
||||
case RELATION -> relationProcessor.processRelationNotification(tenantId, edgeNotificationMsg); |
|
||||
case TENANT -> tenantEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case TENANT_PROFILE -> tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case NOTIFICATION_RULE, NOTIFICATION_TARGET, NOTIFICATION_TEMPLATE -> |
|
||||
notificationEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case TB_RESOURCE -> resourceEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
case DOMAIN, OAUTH2_CLIENT -> oAuth2EdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); |
|
||||
default -> log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); |
|
||||
} |
|
||||
} catch (Exception e) { |
|
||||
callBackFailure(tenantId, edgeNotificationMsg, callback, e); |
|
||||
} |
|
||||
}); |
|
||||
callback.onSuccess(); |
|
||||
} catch (Exception e) { |
|
||||
callBackFailure(tenantId, edgeNotificationMsg, callback, e); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private void callBackFailure(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback, Throwable throwable) { |
|
||||
log.error("[{}] Can't push to edge updates, edgeNotificationMsg [{}]", tenantId, edgeNotificationMsg, throwable); |
|
||||
callback.onFailure(throwable); |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -0,0 +1,332 @@ |
|||||
|
/** |
||||
|
* 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.service.queue; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FutureCallback; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import lombok.Data; |
||||
|
import lombok.Getter; |
||||
|
import lombok.Setter; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.checkerframework.checker.nullness.qual.Nullable; |
||||
|
import org.jetbrains.annotations.NotNull; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.scheduling.annotation.Scheduled; |
||||
|
import org.springframework.stereotype.Service; |
||||
|
import org.thingsboard.server.actors.ActorSystemContext; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.edge.Edge; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventType; |
||||
|
import org.thingsboard.server.common.data.id.EdgeId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; |
||||
|
import org.thingsboard.server.common.data.queue.QueueConfig; |
||||
|
import org.thingsboard.server.common.msg.edge.EdgeSessionMsg; |
||||
|
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; |
||||
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
||||
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
||||
|
import org.thingsboard.server.common.stats.StatsFactory; |
||||
|
import org.thingsboard.server.common.util.ProtoUtils; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.EdgeNotificationMsgProto; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg; |
||||
|
import org.thingsboard.server.queue.TbQueueConsumer; |
||||
|
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
||||
|
import org.thingsboard.server.queue.discovery.QueueKey; |
||||
|
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; |
||||
|
import org.thingsboard.server.queue.provider.TbCoreQueueFactory; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.edge.EdgeContextComponent; |
||||
|
import org.thingsboard.server.service.queue.consumer.MainQueueConsumerManager; |
||||
|
import org.thingsboard.server.service.queue.processing.AbstractConsumerService; |
||||
|
import org.thingsboard.server.service.queue.processing.IdMsgPair; |
||||
|
|
||||
|
import javax.annotation.PostConstruct; |
||||
|
import javax.annotation.PreDestroy; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ConcurrentHashMap; |
||||
|
import java.util.concurrent.ConcurrentMap; |
||||
|
import java.util.concurrent.CountDownLatch; |
||||
|
import java.util.concurrent.Future; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.stream.Collectors; |
||||
|
|
||||
|
@Slf4j |
||||
|
@Service |
||||
|
@TbCoreComponent |
||||
|
public class DefaultTbEdgeConsumerService extends AbstractConsumerService<ToEdgeNotificationMsg> implements TbEdgeConsumerService { |
||||
|
|
||||
|
@Value("${queue.edge.pool-interval:25}") |
||||
|
private int pollInterval; |
||||
|
@Value("${queue.edge.pack-processing-timeout:10000}") |
||||
|
private int packProcessingTimeout; |
||||
|
@Value("${queue.edge.consumer-per-partition:false}") |
||||
|
private boolean consumerPerPartition; |
||||
|
@Value("${queue.edge.pack-processing-retries:3}") |
||||
|
private int packProcessingRetries; |
||||
|
@Value("${queue.edge.stats.enabled:false}") |
||||
|
private boolean statsEnabled; |
||||
|
|
||||
|
private final TbCoreQueueFactory queueFactory; |
||||
|
private final EdgeContextComponent edgeCtx; |
||||
|
private final EdgeConsumerStats stats; |
||||
|
|
||||
|
private MainQueueConsumerManager<TbProtoQueueMsg<ToEdgeMsg>, EdgeQueueConfig> mainConsumer; |
||||
|
|
||||
|
public DefaultTbEdgeConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext, |
||||
|
StatsFactory statsFactory, EdgeContextComponent edgeCtx) { |
||||
|
super(actorContext, null, null, null, null, null, |
||||
|
null, null); |
||||
|
this.edgeCtx = edgeCtx; |
||||
|
this.stats = new EdgeConsumerStats(statsFactory); |
||||
|
this.queueFactory = tbCoreQueueFactory; |
||||
|
} |
||||
|
|
||||
|
@PostConstruct |
||||
|
public void init() { |
||||
|
super.init("tb-edge"); |
||||
|
|
||||
|
this.mainConsumer = MainQueueConsumerManager.<TbProtoQueueMsg<ToEdgeMsg>, EdgeQueueConfig>builder() |
||||
|
.queueKey(new QueueKey(ServiceType.TB_CORE).withQueueName(DataConstants.EDGE_QUEUE_NAME)) |
||||
|
.config(EdgeQueueConfig.of(consumerPerPartition, (int) pollInterval)) |
||||
|
.msgPackProcessor(this::processMsgs) |
||||
|
.consumerCreator((config, partitionId) -> queueFactory.createEdgeMsgConsumer()) |
||||
|
.consumerExecutor(consumersExecutor) |
||||
|
.scheduler(scheduler) |
||||
|
.taskExecutor(mgmtExecutor) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@PreDestroy |
||||
|
public void destroy() { |
||||
|
super.destroy(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void startConsumers() { |
||||
|
super.startConsumers(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void onTbApplicationEvent(PartitionChangeEvent event) { |
||||
|
var partitions = event.getEdgePartitions(); |
||||
|
log.debug("Subscribing to partitions: {}", partitions); |
||||
|
mainConsumer.update(partitions); |
||||
|
} |
||||
|
|
||||
|
private void processMsgs(List<TbProtoQueueMsg<ToEdgeMsg>> msgs, TbQueueConsumer<TbProtoQueueMsg<ToEdgeMsg>> consumer, EdgeQueueConfig edgeQueueConfig) throws InterruptedException { |
||||
|
List<IdMsgPair<ToEdgeMsg>> orderedMsgList = msgs.stream().map(msg -> new IdMsgPair<>(UUID.randomUUID(), msg)).toList(); |
||||
|
ConcurrentMap<UUID, TbProtoQueueMsg<ToEdgeMsg>> pendingMap = orderedMsgList.stream().collect( |
||||
|
Collectors.toConcurrentMap(IdMsgPair::getUuid, IdMsgPair::getMsg)); |
||||
|
CountDownLatch processingTimeoutLatch = new CountDownLatch(1); |
||||
|
TbPackProcessingContext<TbProtoQueueMsg<ToEdgeMsg>> ctx = new TbPackProcessingContext<>( |
||||
|
processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>()); |
||||
|
PendingMsgHolder pendingMsgHolder = new PendingMsgHolder(); |
||||
|
Future<?> submitFuture = consumersExecutor.submit(() -> { |
||||
|
orderedMsgList.forEach((element) -> { |
||||
|
UUID id = element.getUuid(); |
||||
|
TbProtoQueueMsg<ToEdgeMsg> msg = element.getMsg(); |
||||
|
TbCallback callback = new TbPackCallback<>(id, ctx); |
||||
|
try { |
||||
|
ToEdgeMsg toEdgeMsg = msg.getValue(); |
||||
|
pendingMsgHolder.setToEdgeMsg(toEdgeMsg); |
||||
|
if (toEdgeMsg.hasEdgeNotificationMsg()) { |
||||
|
pushNotificationToEdge(toEdgeMsg.getEdgeNotificationMsg(), 0, packProcessingRetries, callback); |
||||
|
} |
||||
|
if (statsEnabled) { |
||||
|
stats.log(toEdgeMsg); |
||||
|
} |
||||
|
} catch (Throwable e) { |
||||
|
log.warn("[{}] Failed to process message: {}", id, msg, e); |
||||
|
callback.onFailure(e); |
||||
|
} |
||||
|
}); |
||||
|
}); |
||||
|
if (!processingTimeoutLatch.await(packProcessingTimeout, TimeUnit.MILLISECONDS)) { |
||||
|
if (!submitFuture.isDone()) { |
||||
|
submitFuture.cancel(true); |
||||
|
ToEdgeMsg lastSubmitMsg = pendingMsgHolder.getToEdgeMsg(); |
||||
|
log.info("Timeout to process message: {}", lastSubmitMsg); |
||||
|
} |
||||
|
ctx.getFailedMap().forEach((id, msg) -> log.warn("[{}] Failed to process message: {}", id, msg.getValue())); |
||||
|
} |
||||
|
consumer.commit(); |
||||
|
} |
||||
|
|
||||
|
private static class PendingMsgHolder { |
||||
|
@Getter |
||||
|
@Setter |
||||
|
private volatile ToEdgeMsg toEdgeMsg; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected ServiceType getServiceType() { |
||||
|
return ServiceType.TB_CORE; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected long getNotificationPollDuration() { |
||||
|
return pollInterval; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected long getNotificationPackProcessingTimeout() { |
||||
|
return packProcessingTimeout; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected int getMgmtThreadPoolSize() { |
||||
|
return Math.max(Runtime.getRuntime().availableProcessors(), 4); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected TbQueueConsumer<TbProtoQueueMsg<ToEdgeNotificationMsg>> createNotificationsConsumer() { |
||||
|
return queueFactory.createToEdgeNotificationsMsgConsumer(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void handleNotification(UUID id, TbProtoQueueMsg<ToEdgeNotificationMsg> msg, TbCallback callback) { |
||||
|
ToEdgeNotificationMsg toEdgeNotificationMsg = msg.getValue(); |
||||
|
try { |
||||
|
if (toEdgeNotificationMsg.hasEdgeHighPriority()) { |
||||
|
EdgeSessionMsg edgeSessionMsg = ProtoUtils.fromProto(toEdgeNotificationMsg.getEdgeHighPriority()); |
||||
|
edgeCtx.getEdgeRpcService().onToEdgeSessionMsg(edgeSessionMsg.getTenantId(), edgeSessionMsg); |
||||
|
callback.onSuccess(); |
||||
|
} else if (toEdgeNotificationMsg.hasEdgeEventUpdate()) { |
||||
|
EdgeSessionMsg edgeSessionMsg = ProtoUtils.fromProto(toEdgeNotificationMsg.getEdgeEventUpdate()); |
||||
|
edgeCtx.getEdgeRpcService().onToEdgeSessionMsg(edgeSessionMsg.getTenantId(), edgeSessionMsg); |
||||
|
callback.onSuccess(); |
||||
|
} else if (toEdgeNotificationMsg.hasToEdgeSyncRequest()) { |
||||
|
EdgeSessionMsg edgeSessionMsg = ProtoUtils.fromProto(toEdgeNotificationMsg.getToEdgeSyncRequest()); |
||||
|
edgeCtx.getEdgeRpcService().onToEdgeSessionMsg(edgeSessionMsg.getTenantId(), edgeSessionMsg); |
||||
|
callback.onSuccess(); |
||||
|
} else if (toEdgeNotificationMsg.hasFromEdgeSyncResponse()) { |
||||
|
EdgeSessionMsg edgeSessionMsg = ProtoUtils.fromProto(toEdgeNotificationMsg.getFromEdgeSyncResponse()); |
||||
|
edgeCtx.getEdgeRpcService().onToEdgeSessionMsg(edgeSessionMsg.getTenantId(), edgeSessionMsg); |
||||
|
callback.onSuccess(); |
||||
|
} else if (toEdgeNotificationMsg.hasComponentLifecycle()) { |
||||
|
ComponentLifecycleMsg componentLifecycle = ProtoUtils.fromProto(toEdgeNotificationMsg.getComponentLifecycle()); |
||||
|
TenantId tenantId = componentLifecycle.getTenantId(); |
||||
|
EdgeId edgeId = new EdgeId(componentLifecycle.getEntityId().getId()); |
||||
|
if (ComponentLifecycleEvent.DELETED.equals(componentLifecycle.getEvent())) { |
||||
|
edgeCtx.getEdgeRpcService().deleteEdge(tenantId, edgeId); |
||||
|
} else if (ComponentLifecycleEvent.UPDATED.equals(componentLifecycle.getEvent())) { |
||||
|
Edge edge = edgeCtx.getEdgeService().findEdgeById(tenantId, edgeId); |
||||
|
edgeCtx.getEdgeRpcService().updateEdge(tenantId, edge); |
||||
|
} |
||||
|
callback.onSuccess(); |
||||
|
} |
||||
|
} catch (Exception e) { |
||||
|
log.error("Error processing edge notification message", e); |
||||
|
callback.onFailure(e); |
||||
|
} |
||||
|
|
||||
|
if (statsEnabled) { |
||||
|
stats.log(msg.getValue()); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void pushNotificationToEdge(EdgeNotificationMsgProto edgeNotificationMsg, int retryCount, int retryLimit, TbCallback callback) { |
||||
|
TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB())); |
||||
|
log.debug("[{}] Pushing notification to edge {}", tenantId, edgeNotificationMsg); |
||||
|
try { |
||||
|
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); |
||||
|
ListenableFuture<Void> future; |
||||
|
switch (type) { |
||||
|
case EDGE -> future = edgeCtx.getEdgeProcessor().processEdgeNotification(tenantId, edgeNotificationMsg); |
||||
|
case ASSET -> future = edgeCtx.getAssetProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case ASSET_PROFILE -> future = edgeCtx.getAssetProfileProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case DEVICE -> future = edgeCtx.getDeviceProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case DEVICE_PROFILE -> future = edgeCtx.getDeviceProfileProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case ENTITY_VIEW -> future = edgeCtx.getEntityViewProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case DASHBOARD -> future = edgeCtx.getDashboardProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case RULE_CHAIN -> future = edgeCtx.getRuleChainProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case USER -> future = edgeCtx.getUserProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case CUSTOMER -> future = edgeCtx.getCustomerProcessor().processCustomerNotification(tenantId, edgeNotificationMsg); |
||||
|
case OTA_PACKAGE -> future = edgeCtx.getOtaPackageProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case WIDGETS_BUNDLE -> future = edgeCtx.getWidgetBundleProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case WIDGET_TYPE -> future = edgeCtx.getWidgetTypeProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case QUEUE -> future = edgeCtx.getQueueProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case ALARM -> future = edgeCtx.getAlarmProcessor().processAlarmNotification(tenantId, edgeNotificationMsg); |
||||
|
case ALARM_COMMENT -> future = edgeCtx.getAlarmProcessor().processAlarmCommentNotification(tenantId, edgeNotificationMsg); |
||||
|
case RELATION -> future = edgeCtx.getRelationProcessor().processRelationNotification(tenantId, edgeNotificationMsg); |
||||
|
case TENANT -> future = edgeCtx.getTenantProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case TENANT_PROFILE -> future = edgeCtx.getTenantProfileProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case NOTIFICATION_RULE, NOTIFICATION_TARGET, NOTIFICATION_TEMPLATE -> |
||||
|
future = edgeCtx.getNotificationEdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case TB_RESOURCE -> future = edgeCtx.getResourceProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
case DOMAIN, OAUTH2_CLIENT -> future = edgeCtx.getOAuth2EdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg); |
||||
|
default -> { |
||||
|
future = Futures.immediateFuture(null); |
||||
|
log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); |
||||
|
} |
||||
|
} |
||||
|
Futures.addCallback(future, new FutureCallback<>() { |
||||
|
@Override |
||||
|
public void onSuccess(@Nullable Void unused) { |
||||
|
callback.onSuccess(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void onFailure(@NotNull Throwable throwable) { |
||||
|
if (retryCount < retryLimit) { |
||||
|
log.warn("[{}] Retry {} for message due to failure: {}", tenantId, retryCount + 1, throwable.getMessage()); |
||||
|
pushNotificationToEdge(edgeNotificationMsg, retryCount + 1, retryLimit, callback); |
||||
|
} else { |
||||
|
callBackFailure(tenantId, edgeNotificationMsg, callback, throwable); |
||||
|
} |
||||
|
} |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} catch (Exception e) { |
||||
|
if (retryCount < retryLimit) { |
||||
|
log.warn("[{}] Retry {} for message due to exception: {}", tenantId, retryCount + 1, e.getMessage()); |
||||
|
pushNotificationToEdge(edgeNotificationMsg, retryCount + 1, retryLimit, callback); |
||||
|
} else { |
||||
|
callBackFailure(tenantId, edgeNotificationMsg, callback, e); |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void callBackFailure(TenantId tenantId, EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback, Throwable throwable) { |
||||
|
log.error("[{}] Can't push to edge updates, edgeNotificationMsg [{}]", tenantId, edgeNotificationMsg, throwable); |
||||
|
callback.onFailure(throwable); |
||||
|
} |
||||
|
|
||||
|
@Scheduled(fixedDelayString = "${queue.edge.stats.print-interval-ms}") |
||||
|
public void printStats() { |
||||
|
if (statsEnabled) { |
||||
|
stats.printStats(); |
||||
|
stats.reset(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void stopConsumers() { |
||||
|
super.stopConsumers(); |
||||
|
mainConsumer.stop(); |
||||
|
mainConsumer.awaitStop(); |
||||
|
} |
||||
|
|
||||
|
@Data(staticConstructor = "of") |
||||
|
public static class EdgeQueueConfig implements QueueConfig { |
||||
|
private final boolean consumerPerPartition; |
||||
|
private final int pollInterval; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,99 @@ |
|||||
|
/** |
||||
|
* 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.service.queue; |
||||
|
|
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.server.common.stats.StatsCounter; |
||||
|
import org.thingsboard.server.common.stats.StatsFactory; |
||||
|
import org.thingsboard.server.common.stats.StatsType; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeNotificationMsg; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
|
||||
|
@Slf4j |
||||
|
public class EdgeConsumerStats { |
||||
|
|
||||
|
public static final String TOTAL_MSGS = "totalMsgs"; |
||||
|
public static final String EDGE_NOTIFICATIONS = "edgeNfs"; |
||||
|
public static final String TO_CORE_NF_EDGE_EVENT = "coreNfEdgeHPUpd"; |
||||
|
public static final String TO_CORE_NF_EDGE_EVENT_UPDATE = "coreNfEdgeUpd"; |
||||
|
public static final String TO_CORE_NF_EDGE_SYNC_REQUEST = "coreNfEdgeSyncReq"; |
||||
|
public static final String TO_CORE_NF_EDGE_SYNC_RESPONSE = "coreNfEdgeSyncResp"; |
||||
|
public static final String TO_CORE_NF_EDGE_COMPONENT_LIFECYCLE = "coreNfEdgeCompLfcl"; |
||||
|
|
||||
|
private final StatsCounter totalCounter; |
||||
|
private final StatsCounter edgeNotificationsCounter; |
||||
|
private final StatsCounter edgeHighPriorityCounter; |
||||
|
private final StatsCounter edgeEventUpdateCounter; |
||||
|
private final StatsCounter edgeSyncRequestCounter; |
||||
|
private final StatsCounter edgeSyncResponseCounter; |
||||
|
private final StatsCounter edgeComponentLifecycle; |
||||
|
|
||||
|
private final List<StatsCounter> counters = new ArrayList<>(7); |
||||
|
|
||||
|
public EdgeConsumerStats(StatsFactory statsFactory) { |
||||
|
String statsKey = StatsType.EDGE.getName(); |
||||
|
|
||||
|
this.totalCounter = register(statsFactory.createStatsCounter(statsKey, TOTAL_MSGS)); |
||||
|
this.edgeNotificationsCounter = register(statsFactory.createStatsCounter(statsKey, EDGE_NOTIFICATIONS)); |
||||
|
this.edgeHighPriorityCounter = register(statsFactory.createStatsCounter(statsKey, TO_CORE_NF_EDGE_EVENT)); |
||||
|
this.edgeEventUpdateCounter = register(statsFactory.createStatsCounter(statsKey, TO_CORE_NF_EDGE_EVENT_UPDATE)); |
||||
|
this.edgeSyncRequestCounter = register(statsFactory.createStatsCounter(statsKey, TO_CORE_NF_EDGE_SYNC_REQUEST)); |
||||
|
this.edgeSyncResponseCounter = register(statsFactory.createStatsCounter(statsKey, TO_CORE_NF_EDGE_SYNC_RESPONSE)); |
||||
|
this.edgeComponentLifecycle = register(statsFactory.createStatsCounter(statsKey, TO_CORE_NF_EDGE_COMPONENT_LIFECYCLE)); |
||||
|
} |
||||
|
|
||||
|
private StatsCounter register(StatsCounter counter) { |
||||
|
counters.add(counter); |
||||
|
return counter; |
||||
|
} |
||||
|
|
||||
|
public void log(ToEdgeNotificationMsg msg) { |
||||
|
totalCounter.increment(); |
||||
|
if (msg.hasEdgeHighPriority()) { |
||||
|
edgeHighPriorityCounter.increment(); |
||||
|
} else if (msg.hasEdgeEventUpdate()) { |
||||
|
edgeEventUpdateCounter.increment(); |
||||
|
} else if (msg.hasToEdgeSyncRequest()) { |
||||
|
edgeSyncRequestCounter.increment(); |
||||
|
} else if (msg.hasFromEdgeSyncResponse()) { |
||||
|
edgeSyncResponseCounter.increment(); |
||||
|
} else if (msg.hasComponentLifecycle()) { |
||||
|
edgeComponentLifecycle.increment(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void log(ToEdgeMsg msg) { |
||||
|
totalCounter.increment(); |
||||
|
edgeNotificationsCounter.increment(); |
||||
|
} |
||||
|
|
||||
|
public void printStats() { |
||||
|
int total = totalCounter.get(); |
||||
|
if (total > 0) { |
||||
|
StringBuilder stats = new StringBuilder(); |
||||
|
counters.forEach(counter -> stats.append(counter.getName()).append(" = [").append(counter.get()).append("] ")); |
||||
|
log.info("Edge Stats: {}", stats); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void reset() { |
||||
|
counters.forEach(StatsCounter::clear); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,19 @@ |
|||||
|
/** |
||||
|
* 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.service.queue; |
||||
|
|
||||
|
public interface TbEdgeConsumerService { |
||||
|
} |
||||
@ -0,0 +1,31 @@ |
|||||
|
/** |
||||
|
* 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.settings; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import org.springframework.beans.factory.annotation.Value; |
||||
|
import org.springframework.context.annotation.Lazy; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
|
||||
|
@Lazy |
||||
|
@Data |
||||
|
@Component |
||||
|
public class TbQueueEdgeSettings { |
||||
|
|
||||
|
@Value("${queue.edge.topic}") |
||||
|
private String topic; |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue