From e33c3b43a597e5a5a899bf552ac8e1a0888cca63 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 16 Oct 2023 13:36:47 +0300 Subject: [PATCH] DefaultEdgeNotificationService - onSuccess() make ASAP --- .../edge/DefaultEdgeNotificationService.java | 156 ++++++++---------- 1 file changed, 72 insertions(+), 84 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index c8e0c09d50..7f52ebb94b 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -16,11 +16,7 @@ package org.thingsboard.server.service.edge; import com.fasterxml.jackson.databind.node.ObjectNode; -import com.google.common.util.concurrent.FutureCallback; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.checkerframework.checker.nullness.qual.Nullable; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; @@ -162,88 +158,80 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { 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(10); - executor.submit(() -> { - try { - if (deadline < System.nanoTime()) { - log.warn("[{}] Skipping notification message because deadline reached {}", tenantId, edgeNotificationMsg); - return; - } - EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); - ListenableFuture future; - switch (type) { - case EDGE: - future = edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg); - break; - case ASSET: - future = assetProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case DEVICE: - future = deviceProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case ENTITY_VIEW: - future = entityViewProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case DASHBOARD: - future = dashboardProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case RULE_CHAIN: - future = ruleChainProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case USER: - future = userProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case CUSTOMER: - future = customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg); - break; - case DEVICE_PROFILE: - future = deviceProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case ASSET_PROFILE: - future = assetProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case OTA_PACKAGE: - future = otaPackageProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case WIDGETS_BUNDLE: - future = widgetBundleProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case WIDGET_TYPE: - future = widgetTypeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case QUEUE: - future = queueProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case ALARM: - future = alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg); - break; - case RELATION: - future = relationProcessor.processRelationNotification(tenantId, edgeNotificationMsg); - break; - case TENANT: - future = tenantEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - case TENANT_PROFILE: - future = tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); - break; - default: - log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); - future = Futures.immediateFuture(null); - } - Futures.addCallback(future, new FutureCallback<>() { - @Override - public void onSuccess(@Nullable Void unused) { - callback.onSuccess(); + try { + executor.submit(() -> { + try { + if (deadline < System.nanoTime()) { + log.warn("[{}] Skipping notification message because deadline reached {}", tenantId, edgeNotificationMsg); + return; } - - @Override - public void onFailure(Throwable throwable) { - callBackFailure(tenantId, edgeNotificationMsg, callback, throwable); + EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); + switch (type) { + case EDGE: + edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg); + break; + case ASSET: + assetProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case DEVICE: + deviceProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case ENTITY_VIEW: + entityViewProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case DASHBOARD: + dashboardProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case RULE_CHAIN: + ruleChainProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case USER: + userProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case CUSTOMER: + customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg); + break; + case DEVICE_PROFILE: + deviceProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case ASSET_PROFILE: + assetProfileProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case OTA_PACKAGE: + otaPackageProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case WIDGETS_BUNDLE: + widgetBundleProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case WIDGET_TYPE: + widgetTypeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case QUEUE: + queueProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case ALARM: + alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg); + break; + case RELATION: + relationProcessor.processRelationNotification(tenantId, edgeNotificationMsg); + break; + case TENANT: + tenantEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + case TENANT_PROFILE: + tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + break; + default: + log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); } - }, executor); - } catch (Exception e) { - callBackFailure(tenantId, edgeNotificationMsg, callback, e); - } - }); + } 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) {