|
|
|
@ -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<Void> 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) { |
|
|
|
|