Browse Source

Added logs for newly added functionality. Added tenantId into edge logs

pull/9175/head
Volodymyr Babak 3 years ago
parent
commit
0b9f7f0518
  1. 14
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  2. 50
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  3. 86
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  4. 43
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmMsgConstructor.java
  5. 11
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/rule/AbstractRuleChainMetadataConstructor.java
  7. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java
  8. 37
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java
  9. 8
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessor.java
  10. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProcessor.java
  11. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProfileProcessor.java
  12. 19
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/BaseDeviceProcessor.java
  13. 9
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/BaseDeviceProfileProcessor.java
  14. 9
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessor.java
  15. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeProcessor.java
  16. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/BaseEntityViewProcessor.java
  17. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java
  18. 25
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java
  19. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/TelemetryEdgeProcessor.java
  20. 14
      application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
  21. 54
      application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java
  22. 10
      application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java
  23. 47
      application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java
  24. 2
      application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java

14
application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java

@ -155,9 +155,9 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Override @Override
public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) { public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) {
log.debug("Pushing notification to edge {}", edgeNotificationMsg); TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB()));
log.debug("[{}] Pushing notification to edge {}", tenantId, edgeNotificationMsg);
try { try {
TenantId tenantId = TenantId.fromUUID(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB()));
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());
ListenableFuture<Void> future; ListenableFuture<Void> future;
switch (type) { switch (type) {
@ -216,7 +216,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
future = tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg); future = tenantProfileEdgeProcessor.processEntityNotification(tenantId, edgeNotificationMsg);
break; break;
default: default:
log.warn("Edge event type [{}] is not designed to be pushed to edge", type); log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type);
future = Futures.immediateFuture(null); future = Futures.immediateFuture(null);
} }
Futures.addCallback(future, new FutureCallback<>() { Futures.addCallback(future, new FutureCallback<>() {
@ -227,16 +227,16 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Override @Override
public void onFailure(Throwable throwable) { public void onFailure(Throwable throwable) {
callBackFailure(edgeNotificationMsg, callback, throwable); callBackFailure(tenantId, edgeNotificationMsg, callback, throwable);
} }
}, dbCallBackExecutor); }, dbCallBackExecutor);
} catch (Exception e) { } catch (Exception e) {
callBackFailure(edgeNotificationMsg, callback, e); callBackFailure(tenantId, edgeNotificationMsg, callback, e);
} }
} }
private void callBackFailure(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback, Throwable throwable) { private void callBackFailure(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback, Throwable throwable) {
log.error("Can't push to edge updates, edgeNotificationMsg [{}]", edgeNotificationMsg, throwable); log.error("[{}] Can't push to edge updates, edgeNotificationMsg [{}]", tenantId, edgeNotificationMsg, throwable);
callback.onFailure(throwable); callback.onFailure(throwable);
} }

50
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -195,17 +195,17 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
switch (msg.getMsgType()) { switch (msg.getMsgType()) {
case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG: case EDGE_EVENT_UPDATE_TO_EDGE_SESSION_MSG:
EdgeEventUpdateMsg edgeEventUpdateMsg = (EdgeEventUpdateMsg) msg; EdgeEventUpdateMsg edgeEventUpdateMsg = (EdgeEventUpdateMsg) msg;
log.trace("[{}] onToEdgeSessionMsg [{}]", edgeEventUpdateMsg.getTenantId(), msg); log.trace("[{}] onToEdgeSessionMsg [{}]", tenantId, msg);
onEdgeEvent(tenantId, edgeEventUpdateMsg.getEdgeId()); onEdgeEvent(tenantId, edgeEventUpdateMsg.getEdgeId());
break; break;
case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG: case EDGE_SYNC_REQUEST_TO_EDGE_SESSION_MSG:
ToEdgeSyncRequest toEdgeSyncRequest = (ToEdgeSyncRequest) msg; ToEdgeSyncRequest toEdgeSyncRequest = (ToEdgeSyncRequest) msg;
log.trace("[{}] toEdgeSyncRequest [{}]", toEdgeSyncRequest.getTenantId(), msg); log.trace("[{}] toEdgeSyncRequest [{}]", tenantId, msg);
startSyncProcess(tenantId, toEdgeSyncRequest.getEdgeId(), toEdgeSyncRequest.getId()); startSyncProcess(tenantId, toEdgeSyncRequest.getEdgeId(), toEdgeSyncRequest.getId());
break; break;
case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG: case EDGE_SYNC_RESPONSE_FROM_EDGE_SESSION_MSG:
FromEdgeSyncResponse fromEdgeSyncResponse = (FromEdgeSyncResponse) msg; FromEdgeSyncResponse fromEdgeSyncResponse = (FromEdgeSyncResponse) msg;
log.trace("[{}] fromEdgeSyncResponse [{}]", fromEdgeSyncResponse.getTenantId(), msg); log.trace("[{}] fromEdgeSyncResponse [{}]", tenantId, msg);
processSyncResponse(fromEdgeSyncResponse); processSyncResponse(fromEdgeSyncResponse);
break; break;
} }
@ -263,7 +263,8 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} }
private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) {
log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId); TenantId tenantId = edgeGrpcSession.getEdge().getTenantId();
log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId);
sessions.put(edgeId, edgeGrpcSession); sessions.put(edgeId, edgeGrpcSession);
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock()); final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock(); newEventLock.lock();
@ -272,10 +273,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} finally { } finally {
newEventLock.unlock(); newEventLock.unlock();
} }
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true); save(tenantId, edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
long lastConnectTs = System.currentTimeMillis(); long lastConnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, lastConnectTs); save(tenantId, edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, lastConnectTs);
pushRuleEngineMessage(edgeGrpcSession.getEdge().getTenantId(), edgeId, lastConnectTs, TbMsgType.CONNECT_EVENT); pushRuleEngineMessage(tenantId, edgeId, lastConnectTs, TbMsgType.CONNECT_EVENT);
cancelScheduleEdgeEventsCheck(edgeId); cancelScheduleEdgeEventsCheck(edgeId);
scheduleEdgeEventsCheck(edgeGrpcSession); scheduleEdgeEventsCheck(edgeGrpcSession);
} }
@ -334,7 +335,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
newEventLock.lock(); newEventLock.lock();
try { try {
if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) {
log.trace("[{}] Set session new events flag to false", edgeId.getId()); log.trace("[{}][{}] Set session new events flag to false", tenantId, edgeId.getId());
sessionNewEvents.put(edgeId, false); sessionNewEvents.put(edgeId, false);
Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() { Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() {
@Override @Override
@ -392,9 +393,10 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} finally { } finally {
newEventLock.unlock(); newEventLock.unlock();
} }
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); TenantId tenantId = toRemove.getEdge().getTenantId();
save(tenantId, edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
long lastDisconnectTs = System.currentTimeMillis(); long lastDisconnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, lastDisconnectTs); save(tenantId, edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, lastDisconnectTs);
pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edgeId, lastDisconnectTs, TbMsgType.DISCONNECT_EVENT); pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edgeId, lastDisconnectTs, TbMsgType.DISCONNECT_EVENT);
cancelScheduleEdgeEventsCheck(edgeId); cancelScheduleEdgeEventsCheck(edgeId);
} else { } else {
@ -402,36 +404,38 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} }
} }
private void save(EdgeId edgeId, String key, long value) { private void save(TenantId tenantId, EdgeId edgeId, String key, long value) {
log.debug("[{}] Updating long edge telemetry [{}] [{}]", edgeId, key, value); log.debug("[{}][{}] Updating long edge telemetry [{}] [{}]", tenantId, edgeId, key, value);
if (persistToTelemetry) { if (persistToTelemetry) {
tsSubService.saveAndNotify( tsSubService.saveAndNotify(
TenantId.SYS_TENANT_ID, edgeId, tenantId, edgeId,
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))),
new AttributeSaveCallback(edgeId, key, value)); new AttributeSaveCallback(tenantId, edgeId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(edgeId, key, value)); tsSubService.saveAttrAndNotify(tenantId, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value));
} }
} }
private void save(EdgeId edgeId, String key, boolean value) { private void save(TenantId tenantId, EdgeId edgeId, String key, boolean value) {
log.debug("[{}] Updating boolean edge telemetry [{}] [{}]", edgeId, key, value); log.debug("[{}][{}] Updating boolean edge telemetry [{}] [{}]", tenantId, edgeId, key, value);
if (persistToTelemetry) { if (persistToTelemetry) {
tsSubService.saveAndNotify( tsSubService.saveAndNotify(
TenantId.SYS_TENANT_ID, edgeId, tenantId, edgeId,
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(key, value))),
new AttributeSaveCallback(edgeId, key, value)); new AttributeSaveCallback(tenantId, edgeId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(edgeId, key, value)); tsSubService.saveAttrAndNotify(tenantId, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value));
} }
} }
private static class AttributeSaveCallback implements FutureCallback<Void> { private static class AttributeSaveCallback implements FutureCallback<Void> {
private final TenantId tenantId;
private final EdgeId edgeId; private final EdgeId edgeId;
private final String key; private final String key;
private final Object value; private final Object value;
AttributeSaveCallback(EdgeId edgeId, String key, Object value) { AttributeSaveCallback(TenantId tenantId, EdgeId edgeId, String key, Object value) {
this.tenantId = tenantId;
this.edgeId = edgeId; this.edgeId = edgeId;
this.key = key; this.key = key;
this.value = value; this.value = value;
@ -439,12 +443,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override @Override
public void onSuccess(@Nullable Void result) { public void onSuccess(@Nullable Void result) {
log.trace("[{}] Successfully updated attribute [{}] with value [{}]", edgeId, key, value); log.trace("[{}][{}] Successfully updated attribute [{}] with value [{}]", tenantId, edgeId, key, value);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.warn("[{}] Failed to update attribute [{}] with value [{}]", edgeId, key, value, t); log.warn("[{}][{}] Failed to update attribute [{}] with value [{}]", tenantId, edgeId, key, value, t);
} }
} }

86
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -105,6 +105,7 @@ public final class EdgeGrpcSession implements Closeable {
private EdgeContextComponent ctx; private EdgeContextComponent ctx;
private Edge edge; private Edge edge;
private TenantId tenantId;
private StreamObserver<RequestMsg> inputStream; private StreamObserver<RequestMsg> inputStream;
private StreamObserver<ResponseMsg> outputStream; private StreamObserver<ResponseMsg> outputStream;
private boolean connected; private boolean connected;
@ -148,7 +149,7 @@ public final class EdgeGrpcSession implements Closeable {
outputStream.onError(new RuntimeException(responseMsg.getErrorMsg())); outputStream.onError(new RuntimeException(responseMsg.getErrorMsg()));
} else { } else {
if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) { if (requestMsg.getConnectRequestMsg().hasMaxInboundMessageSize()) {
log.debug("[{}] Client max inbound message size: {}", sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize()); log.debug("[{}][{}] Client max inbound message size: {}", tenantId, sessionId, requestMsg.getConnectRequestMsg().getMaxInboundMessageSize());
clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize(); clientMaxInboundMessageSize = requestMsg.getConnectRequestMsg().getMaxInboundMessageSize();
} }
connected = true; connected = true;
@ -179,13 +180,13 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onError(Throwable t) { public void onError(Throwable t) {
log.error("[{}] Stream was terminated due to error:", sessionId, t); log.error("[{}][{}] Stream was terminated due to error:", tenantId, sessionId, t);
closeSession(); closeSession();
} }
@Override @Override
public void onCompleted() { public void onCompleted() {
log.info("[{}] Stream was closed and completed successfully!", sessionId); log.info("[{}][{}] Stream was closed and completed successfully!", tenantId, sessionId);
closeSession(); closeSession();
} }
@ -206,7 +207,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
public void startSyncProcess(boolean fullSync) { public void startSyncProcess(boolean fullSync) {
log.trace("[{}][{}][{}] Staring edge sync process", edge.getTenantId(), edge.getId(), this.sessionId); log.trace("[{}][{}][{}] Staring edge sync process", this.tenantId, edge.getId(), this.sessionId);
syncCompleted = false; syncCompleted = false;
interruptGeneralProcessingOnSync(); interruptGeneralProcessingOnSync();
doSync(new EdgeSyncCursor(ctx, edge, fullSync)); doSync(new EdgeSyncCursor(ctx, edge, fullSync));
@ -216,7 +217,7 @@ public final class EdgeGrpcSession implements Closeable {
if (cursor.hasNext()) { if (cursor.hasNext()) {
EdgeEventFetcher next = cursor.getNext(); EdgeEventFetcher next = cursor.getNext();
log.info("[{}][{}] starting sync process, cursor current idx = {}, class = {}", log.info("[{}][{}] starting sync process, cursor current idx = {}, class = {}",
edge.getTenantId(), edge.getId(), cursor.getCurrentIdx(), next.getClass().getSimpleName()); this.tenantId, edge.getId(), cursor.getCurrentIdx(), next.getClass().getSimpleName());
ListenableFuture<Pair<Long, Long>> future = startProcessingEdgeEvents(next); ListenableFuture<Pair<Long, Long>> future = startProcessingEdgeEvents(next);
Futures.addCallback(future, new FutureCallback<>() { Futures.addCallback(future, new FutureCallback<>() {
@Override @Override
@ -226,7 +227,7 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), t); log.error("[{}][{}] Exception during sync process", tenantId, edge.getId(), t);
} }
}, ctx.getGrpcCallbackExecutorService()); }, ctx.getGrpcCallbackExecutorService());
} else { } else {
@ -243,7 +244,7 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Exception during sending sync complete", edge.getTenantId(), edge.getId(), t); log.error("[{}][{}] Exception during sending sync complete", tenantId, edge.getId(), t);
} }
}, ctx.getGrpcCallbackExecutorService()); }, ctx.getGrpcCallbackExecutorService());
} }
@ -279,39 +280,40 @@ public final class EdgeGrpcSession implements Closeable {
try { try {
if (msg.getSuccess()) { if (msg.getSuccess()) {
sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId());
log.debug("[{}] Msg has been processed successfully!Msd Id: [{}], Msg: {}", edge.getRoutingKey(), msg.getDownlinkMsgId(), msg); log.debug("[{}][{}] Msg has been processed successfully!Msd Id: [{}], Msg: {}", this.tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg);
} else { } else {
log.error("[{}] Msg processing failed! Msd Id: [{}], Error msg: {}", edge.getRoutingKey(), msg.getDownlinkMsgId(), msg.getErrorMsg()); log.error("[{}][{}] Msg processing failed! Msd Id: [{}], Error msg: {}", this.tenantId, edge.getRoutingKey(), msg.getDownlinkMsgId(), msg.getErrorMsg());
} }
if (sessionState.getPendingMsgsMap().isEmpty()) { if (sessionState.getPendingMsgsMap().isEmpty()) {
log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey()); log.debug("[{}][{}] Pending msgs map is empty. Stopping current iteration", this.tenantId, edge.getRoutingKey());
stopCurrentSendDownlinkMsgsTask(false); stopCurrentSendDownlinkMsgsTask(false);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e); log.error("[{}][{}] Can't process downlink response message [{}]", this.tenantId, this.sessionId, msg, e);
} }
} }
private void sendDownlinkMsg(ResponseMsg downlinkMsg) { private void sendDownlinkMsg(ResponseMsg downlinkMsg) {
log.trace("[{}] Sending downlink msg [{}]", this.sessionId, downlinkMsg); log.trace("[{}][{}] Sending downlink msg [{}]", this.tenantId, this.sessionId, downlinkMsg);
if (isConnected()) { if (isConnected()) {
downlinkMsgLock.lock(); downlinkMsgLock.lock();
try { try {
outputStream.onNext(downlinkMsg); outputStream.onNext(downlinkMsg);
} catch (Exception e) { } catch (Exception e) {
log.error("[{}] Failed to send downlink message [{}]", this.sessionId, downlinkMsg, e); log.error("[{}][{}] Failed to send downlink message [{}]", this.tenantId, this.sessionId, downlinkMsg, e);
connected = false; connected = false;
sessionCloseListener.accept(edge.getId(), sessionId); sessionCloseListener.accept(edge.getId(), sessionId);
} finally { } finally {
downlinkMsgLock.unlock(); downlinkMsgLock.unlock();
} }
log.trace("[{}] Response msg successfully sent [{}]", this.sessionId, downlinkMsg); log.trace("[{}][{}] Response msg successfully sent [{}]", this.tenantId, this.sessionId, downlinkMsg);
} }
} }
void onConfigurationUpdate(Edge edge) { void onConfigurationUpdate(Edge edge) {
log.debug("[{}] onConfigurationUpdate [{}]", this.sessionId, edge); log.debug("[{}] onConfigurationUpdate [{}]", this.sessionId, edge);
this.edge = edge; this.edge = edge;
this.tenantId = edge.getTenantId();
EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder() EdgeUpdateMsg edgeConfig = EdgeUpdateMsg.newBuilder()
.setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build(); .setConfiguration(ctx.getEdgeMsgConstructor().constructEdgeConfiguration(edge)).build();
ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder() ResponseMsg edgeConfigMsg = ResponseMsg.newBuilder()
@ -322,7 +324,7 @@ public final class EdgeGrpcSession implements Closeable {
ListenableFuture<Boolean> processEdgeEvents() throws Exception { ListenableFuture<Boolean> processEdgeEvents() throws Exception {
SettableFuture<Boolean> result = SettableFuture.create(); SettableFuture<Boolean> result = SettableFuture.create();
log.trace("[{}] starting processing edge events", this.sessionId); log.trace("[{}][{}] starting processing edge events", this.tenantId, this.sessionId);
if (isConnected() && isSyncCompleted()) { if (isConnected() && isSyncCompleted()) {
Pair<Long, Long> startTsAndSeqId = getQueueStartTsAndSeqId().get(); Pair<Long, Long> startTsAndSeqId = getQueueStartTsAndSeqId().get();
this.previousStartTs = startTsAndSeqId.getFirst(); this.previousStartTs = startTsAndSeqId.getFirst();
@ -342,7 +344,7 @@ public final class EdgeGrpcSession implements Closeable {
Futures.addCallback(updateFuture, new FutureCallback<>() { Futures.addCallback(updateFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable List<String> list) { public void onSuccess(@Nullable List<String> list) {
log.debug("[{}] queue offset was updated [{}]", sessionId, newStartTsAndSeqId); log.debug("[{}][{}] queue offset was updated [{}]", tenantId, sessionId, newStartTsAndSeqId);
if (fetcher.isSeqIdNewCycleStarted()) { if (fetcher.isSeqIdNewCycleStarted()) {
seqIdEnd = fetcher.getSeqIdEnd(); seqIdEnd = fetcher.getSeqIdEnd();
boolean newEventsAvailable = isNewEdgeEventsAvailable(); boolean newEventsAvailable = isNewEdgeEventsAvailable();
@ -359,24 +361,24 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}] Failed to update queue offset [{}]", sessionId, newStartTsAndSeqId, t); log.error("[{}][{}] Failed to update queue offset [{}]", tenantId, sessionId, newStartTsAndSeqId, t);
result.setException(t); result.setException(t);
} }
}, ctx.getGrpcCallbackExecutorService()); }, ctx.getGrpcCallbackExecutorService());
} else { } else {
log.trace("[{}] newStartTsAndSeqId is null. Skipping iteration without db update", sessionId); log.trace("[{}][{}] newStartTsAndSeqId is null. Skipping iteration without db update", tenantId, sessionId);
result.set(null); result.set(null);
} }
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}] Failed to process events", sessionId, t); log.error("[{}][{}] Failed to process events", tenantId, sessionId, t);
result.setException(t); result.setException(t);
} }
}, ctx.getGrpcCallbackExecutorService()); }, ctx.getGrpcCallbackExecutorService());
} else { } else {
log.trace("[{}] edge is not connected or sync is not completed. Skipping iteration", sessionId); log.trace("[{}][{}] edge is not connected or sync is not completed. Skipping iteration", tenantId, sessionId);
result.set(null); result.set(null);
} }
return result; return result;
@ -393,13 +395,13 @@ public final class EdgeGrpcSession implements Closeable {
try { try {
PageData<EdgeEvent> pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); PageData<EdgeEvent> pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink);
if (isConnected() && !pageData.getData().isEmpty()) { if (isConnected() && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size()); log.trace("[{}][{}][{}] event(s) are going to be processed.", this.tenantId, this.sessionId, pageData.getData().size());
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData());
Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable Boolean isInterrupted) { public void onSuccess(@Nullable Boolean isInterrupted) {
if (Boolean.TRUE.equals(isInterrupted)) { if (Boolean.TRUE.equals(isInterrupted)) {
log.debug("[{}][{}][{}] Send downlink messages task was interrupted", edge.getTenantId(), edge.getId(), sessionId); log.debug("[{}][{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId(), sessionId);
result.set(null); result.set(null);
} else { } else {
if (isConnected() && pageData.hasNext()) { if (isConnected() && pageData.hasNext()) {
@ -452,14 +454,14 @@ public final class EdgeGrpcSession implements Closeable {
if (isConnected() && sessionState.getPendingMsgsMap().values().size() > 0) { if (isConnected() && sessionState.getPendingMsgsMap().values().size() > 0) {
List<DownlinkMsg> copy = new ArrayList<>(sessionState.getPendingMsgsMap().values()); List<DownlinkMsg> copy = new ArrayList<>(sessionState.getPendingMsgsMap().values());
if (attempt > 1) { if (attempt > 1) {
log.warn("[{}] Failed to deliver the batch: {}, attempt: {}", this.sessionId, copy, attempt); log.warn("[{}][{}] Failed to deliver the batch: {}, attempt: {}", this.tenantId, this.sessionId, copy, attempt);
} }
log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size()); log.trace("[{}][{}][{}] downlink msg(s) are going to be send.", this.tenantId, this.sessionId, copy.size());
for (DownlinkMsg downlinkMsg : copy) { for (DownlinkMsg downlinkMsg : copy) {
if (this.clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > this.clientMaxInboundMessageSize) { if (this.clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > this.clientMaxInboundMessageSize) {
log.error("[{}][{}][{}] Downlink msg size [{}] exceeds client max inbound message size [{}]. Skipping this message. " + log.error("[{}][{}][{}] Downlink msg size [{}] exceeds client max inbound message size [{}]. Skipping this message. " +
"Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE env variable on the edge and restart it." + "Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE env variable on the edge and restart it." +
"Message {}", edge.getTenantId(), edge.getId(), this.sessionId, downlinkMsg.getSerializedSize(), "Message {}", this.tenantId, edge.getId(), this.sessionId, downlinkMsg.getSerializedSize(),
this.clientMaxInboundMessageSize, downlinkMsg); this.clientMaxInboundMessageSize, downlinkMsg);
sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId()); sessionState.getPendingMsgsMap().remove(downlinkMsg.getDownlinkMsgId());
} else { } else {
@ -471,15 +473,15 @@ public final class EdgeGrpcSession implements Closeable {
if (attempt < MAX_DOWNLINK_ATTEMPTS) { if (attempt < MAX_DOWNLINK_ATTEMPTS) {
scheduleDownlinkMsgsPackSend(attempt + 1); scheduleDownlinkMsgsPackSend(attempt + 1);
} else { } else {
log.warn("[{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", log.warn("[{}][{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}",
this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy); this.tenantId, this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy);
stopCurrentSendDownlinkMsgsTask(false); stopCurrentSendDownlinkMsgsTask(false);
} }
} else { } else {
stopCurrentSendDownlinkMsgsTask(false); stopCurrentSendDownlinkMsgsTask(false);
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Failed to send downlink msgs. Error msg {}", this.sessionId, e.getMessage(), e); log.warn("[{}][{}] Failed to send downlink msgs. Error msg {}", this.tenantId, this.sessionId, e.getMessage(), e);
stopCurrentSendDownlinkMsgsTask(true); stopCurrentSendDownlinkMsgsTask(true);
} }
}; };
@ -499,7 +501,7 @@ public final class EdgeGrpcSession implements Closeable {
private List<DownlinkMsg> convertToDownlinkMsgsPack(List<EdgeEvent> edgeEvents) { private List<DownlinkMsg> convertToDownlinkMsgsPack(List<EdgeEvent> edgeEvents) {
List<DownlinkMsg> result = new ArrayList<>(); List<DownlinkMsg> result = new ArrayList<>();
for (EdgeEvent edgeEvent : edgeEvents) { for (EdgeEvent edgeEvent : edgeEvents) {
log.trace("[{}][{}] converting edge event to downlink msg [{}]", edge.getTenantId(), this.sessionId, edgeEvent); log.trace("[{}][{}] converting edge event to downlink msg [{}]", this.tenantId, this.sessionId, edgeEvent);
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
try { try {
switch (edgeEvent.getAction()) { switch (edgeEvent.getAction()) {
@ -518,7 +520,7 @@ public final class EdgeGrpcSession implements Closeable {
case ASSIGNED_TO_CUSTOMER: case ASSIGNED_TO_CUSTOMER:
case UNASSIGNED_FROM_CUSTOMER: case UNASSIGNED_FROM_CUSTOMER:
downlinkMsg = convertEntityEventToDownlink(edgeEvent); downlinkMsg = convertEntityEventToDownlink(edgeEvent);
log.trace("[{}][{}] entity message processed [{}]", edgeEvent.getTenantId(), this.sessionId, downlinkMsg); log.trace("[{}][{}] entity message processed [{}]", this.tenantId, this.sessionId, downlinkMsg);
break; break;
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
case POST_ATTRIBUTES: case POST_ATTRIBUTES:
@ -527,10 +529,10 @@ public final class EdgeGrpcSession implements Closeable {
downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edgeEvent); downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edgeEvent);
break; break;
default: default:
log.warn("[{}][{}] Unsupported action type [{}]", edge.getTenantId(), this.sessionId, edgeEvent.getAction()); log.warn("[{}][{}] Unsupported action type [{}]", this.tenantId, this.sessionId, edgeEvent.getAction());
} }
} catch (Exception e) { } catch (Exception e) {
log.error("[{}][{}] Exception during converting edge event to downlink msg", edge.getTenantId(), this.sessionId, e); log.error("[{}][{}] Exception during converting edge event to downlink msg", this.tenantId, this.sessionId, e);
} }
if (downlinkMsg != null) { if (downlinkMsg != null) {
result.add(downlinkMsg); result.add(downlinkMsg);
@ -566,7 +568,7 @@ public final class EdgeGrpcSession implements Closeable {
PageData<EdgeEvent> edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), 0L, this.previousStartSeqId == 0 ? null : this.previousStartSeqId - 1, pageLink); PageData<EdgeEvent> edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), 0L, this.previousStartSeqId == 0 ? null : this.previousStartSeqId - 1, pageLink);
return !edgeEvents.getData().isEmpty(); return !edgeEvents.getData().isEmpty();
} catch (Exception e) { } catch (Exception e) {
log.error("[{}][{}][{}] Failed to execute isSeqIdStartedNewCycle", edge.getTenantId(), edge.getId(), sessionId, e); log.error("[{}][{}][{}] Failed to execute isSeqIdStartedNewCycle", this.tenantId, edge.getId(), sessionId, e);
} }
return false; return false;
} }
@ -577,7 +579,7 @@ public final class EdgeGrpcSession implements Closeable {
PageData<EdgeEvent> edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), this.newStartSeqId, null, pageLink); PageData<EdgeEvent> edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), this.newStartSeqId, null, pageLink);
return !edgeEvents.getData().isEmpty(); return !edgeEvents.getData().isEmpty();
} catch (Exception e) { } catch (Exception e) {
log.error("[{}][{}][{}] Failed to execute isNewEdgeEventsAvailable", edge.getTenantId(), edge.getId(), sessionId, e); log.error("[{}][{}][{}] Failed to execute isNewEdgeEventsAvailable", this.tenantId, edge.getId(), sessionId, e);
} }
return false; return false;
} }
@ -591,7 +593,7 @@ public final class EdgeGrpcSession implements Closeable {
startSeqId = edgeEvents.getData().get(0).getSeqId() - 1; startSeqId = edgeEvents.getData().get(0).getSeqId() - 1;
} }
} catch (Exception e) { } catch (Exception e) {
log.error("[{}][{}][{}] Failed to execute findStartSeqIdFromOldestEventIfAny", edge.getTenantId(), edge.getId(), sessionId, e); log.error("[{}][{}][{}] Failed to execute findStartSeqIdFromOldestEventIfAny", this.tenantId, edge.getId(), sessionId, e);
} }
return startSeqId; return startSeqId;
} }
@ -607,7 +609,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
private DownlinkMsg convertEntityEventToDownlink(EdgeEvent edgeEvent) { private DownlinkMsg convertEntityEventToDownlink(EdgeEvent edgeEvent) {
log.trace("Executing convertEntityEventToDownlink, edgeEvent [{}], action [{}]", edgeEvent, edgeEvent.getAction()); log.trace("[{}] Executing convertEntityEventToDownlink, edgeEvent [{}], action [{}]", this.tenantId, edgeEvent, edgeEvent.getAction());
switch (edgeEvent.getType()) { switch (edgeEvent.getType()) {
case EDGE: case EDGE:
return ctx.getEdgeProcessor().convertEdgeEventToDownlink(edgeEvent); return ctx.getEdgeProcessor().convertEdgeEventToDownlink(edgeEvent);
@ -650,7 +652,7 @@ public final class EdgeGrpcSession implements Closeable {
case TENANT_PROFILE: case TENANT_PROFILE:
return ctx.getTenantProfileEdgeProcessor().convertTenantProfileEventToDownlink(edgeEvent); return ctx.getTenantProfileEdgeProcessor().convertTenantProfileEventToDownlink(edgeEvent);
default: default:
log.warn("Unsupported edge event type [{}]", edgeEvent); log.warn("[{}] Unsupported edge event type [{}]", this.tenantId, edgeEvent);
return null; return null;
} }
} }
@ -749,7 +751,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
} }
} catch (Exception e) { } catch (Exception e) {
log.error("[{}] Can't process uplink msg [{}]", this.sessionId, uplinkMsg, e); log.error("[{}][{}] Can't process uplink msg [{}]", this.tenantId, this.sessionId, uplinkMsg, e);
return Futures.immediateFailedFuture(e); return Futures.immediateFailedFuture(e);
} }
return Futures.allAsList(result); return Futures.allAsList(result);
@ -791,25 +793,25 @@ public final class EdgeGrpcSession implements Closeable {
@Override @Override
public void close() { public void close() {
log.debug("[{}] Closing session", sessionId); log.debug("[{}][{}] Closing session", this.tenantId, sessionId);
connected = false; connected = false;
try { try {
outputStream.onCompleted(); outputStream.onCompleted();
} catch (Exception e) { } catch (Exception e) {
log.debug("[{}] Failed to close output stream: {}", sessionId, e.getMessage()); log.debug("[{}][{}] Failed to close output stream: {}", this.tenantId, sessionId, e.getMessage());
} }
} }
private void interruptPreviousSendDownlinkMsgsTask() { private void interruptPreviousSendDownlinkMsgsTask() {
if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone() if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()
|| sessionState.getScheduledSendDownlinkTask() != null && !sessionState.getScheduledSendDownlinkTask().isCancelled()) { || sessionState.getScheduledSendDownlinkTask() != null && !sessionState.getScheduledSendDownlinkTask().isCancelled()) {
log.debug("[{}][{}][{}] Previous send downlink future was not properly completed, stopping it now!", edge.getTenantId(), edge.getId(), this.sessionId); log.debug("[{}][{}][{}] Previous send downlink future was not properly completed, stopping it now!", this.tenantId, edge.getId(), this.sessionId);
stopCurrentSendDownlinkMsgsTask(true); stopCurrentSendDownlinkMsgsTask(true);
} }
} }
private void interruptGeneralProcessingOnSync() { private void interruptGeneralProcessingOnSync() {
log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", edge.getTenantId(), edge.getId(), this.sessionId); log.debug("[{}][{}][{}] Sync process started. General processing interrupted!", this.tenantId, edge.getId(), this.sessionId);
stopCurrentSendDownlinkMsgsTask(true); stopCurrentSendDownlinkMsgsTask(true);
} }

43
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmMsgConstructor.java

@ -15,20 +15,9 @@
*/ */
package org.thingsboard.server.service.edge.rpc.constructor; package org.thingsboard.server.service.edge.rpc.constructor;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
@ -37,37 +26,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent;
@TbCoreComponent @TbCoreComponent
public class AlarmMsgConstructor { public class AlarmMsgConstructor {
@Autowired public AlarmUpdateMsg constructAlarmUpdatedMsg(UpdateMsgType msgType, Alarm alarm, String entityName) {
private DeviceService deviceService;
@Autowired
private AssetService assetService;
@Autowired
private EntityViewService entityViewService;
public AlarmUpdateMsg constructAlarmUpdatedMsg(TenantId tenantId, UpdateMsgType msgType, Alarm alarm) {
String entityName = null;
switch (alarm.getOriginator().getEntityType()) {
case DEVICE:
Device deviceById = deviceService.findDeviceById(tenantId, new DeviceId(alarm.getOriginator().getId()));
if (deviceById != null) {
entityName = deviceById.getName();
}
break;
case ASSET:
Asset assetById = assetService.findAssetById(tenantId, new AssetId(alarm.getOriginator().getId()));
if (assetById != null) {
entityName = assetById.getName();
}
break;
case ENTITY_VIEW:
EntityView entityViewById = entityViewService.findEntityViewById(tenantId, new EntityViewId(alarm.getOriginator().getId()));
if (entityViewById != null) {
entityName = entityViewById.getName();
}
break;
}
AlarmUpdateMsg.Builder builder = AlarmUpdateMsg.newBuilder() AlarmUpdateMsg.Builder builder = AlarmUpdateMsg.newBuilder()
.setMsgType(msgType) .setMsgType(msgType)
.setIdMSB(alarm.getId().getId().getMostSignificantBits()) .setIdMSB(alarm.getId().getId().getMostSignificantBits())

11
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java

@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.gen.edge.v1.AttributeDeleteMsg; import org.thingsboard.server.gen.edge.v1.AttributeDeleteMsg;
import org.thingsboard.server.gen.edge.v1.EntityDataProto; import org.thingsboard.server.gen.edge.v1.EntityDataProto;
@ -40,7 +41,7 @@ import java.util.List;
@TbCoreComponent @TbCoreComponent
public class EntityDataMsgConstructor { public class EntityDataMsgConstructor {
public EntityDataProto constructEntityDataMsg(EntityId entityId, EdgeEventActionType actionType, JsonElement entityData) { public EntityDataProto constructEntityDataMsg(TenantId tenantId, EntityId entityId, EdgeEventActionType actionType, JsonElement entityData) {
EntityDataProto.Builder builder = EntityDataProto.newBuilder() EntityDataProto.Builder builder = EntityDataProto.newBuilder()
.setEntityIdMSB(entityId.getId().getMostSignificantBits()) .setEntityIdMSB(entityId.getId().getMostSignificantBits())
.setEntityIdLSB(entityId.getId().getLeastSignificantBits()) .setEntityIdLSB(entityId.getId().getLeastSignificantBits())
@ -57,7 +58,7 @@ public class EntityDataMsgConstructor {
} }
builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts)); builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts));
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Can't convert to telemetry proto, entityData [{}]", entityId, entityData, e); log.warn("[{}][{}] Can't convert to telemetry proto, entityData [{}]", tenantId, entityId, entityData, e);
} }
break; break;
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
@ -67,7 +68,7 @@ public class EntityDataMsgConstructor {
builder.setAttributesUpdatedMsg(attributesUpdatedMsg); builder.setAttributesUpdatedMsg(attributesUpdatedMsg);
builder.setPostAttributeScope(getScopeOfDefault(data)); builder.setPostAttributeScope(getScopeOfDefault(data));
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Can't convert to AttributesUpdatedMsg proto, entityData [{}]", entityId, entityData, e); log.warn("[{}][{}] Can't convert to AttributesUpdatedMsg proto, entityData [{}]", tenantId, entityId, entityData, e);
} }
break; break;
case POST_ATTRIBUTES: case POST_ATTRIBUTES:
@ -77,7 +78,7 @@ public class EntityDataMsgConstructor {
builder.setPostAttributesMsg(postAttributesMsg); builder.setPostAttributesMsg(postAttributesMsg);
builder.setPostAttributeScope(getScopeOfDefault(data)); builder.setPostAttributeScope(getScopeOfDefault(data));
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Can't convert to PostAttributesMsg, entityData [{}]", entityId, entityData, e); log.warn("[{}][{}] Can't convert to PostAttributesMsg, entityData [{}]", tenantId, entityId, entityData, e);
} }
break; break;
case ATTRIBUTES_DELETED: case ATTRIBUTES_DELETED:
@ -90,7 +91,7 @@ public class EntityDataMsgConstructor {
attributeDeleteMsg.build(); attributeDeleteMsg.build();
builder.setAttributeDeleteMsg(attributeDeleteMsg); builder.setAttributeDeleteMsg(attributeDeleteMsg);
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Can't convert to AttributeDeleteMsg proto, entityData [{}]", entityId, entityData, e); log.warn("[{}][{}] Can't convert to AttributeDeleteMsg proto, entityData [{}]", tenantId, entityId, entityData, e);
} }
break; break;
} }

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/rule/AbstractRuleChainMetadataConstructor.java

@ -51,7 +51,7 @@ public abstract class AbstractRuleChainMetadataConstructor implements RuleChainM
builder.setMsgType(msgType); builder.setMsgType(msgType);
return builder.build(); return builder.build();
} catch (JsonProcessingException ex) { } catch (JsonProcessingException ex) {
log.error("Can't construct RuleChainMetadataUpdateMsg", ex); log.error("[{}] Can't construct RuleChainMetadataUpdateMsg", tenantId, ex);
} }
return null; return null;
} }

6
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java

@ -85,7 +85,7 @@ public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher {
result.add(EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, result.add(EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS,
EdgeEventActionType.UPDATED, null, JacksonUtil.OBJECT_MAPPER.valueToTree(tenantMailSettings))); EdgeEventActionType.UPDATED, null, JacksonUtil.OBJECT_MAPPER.valueToTree(tenantMailSettings)));
AdminSettings systemMailTemplates = loadMailTemplates(); AdminSettings systemMailTemplates = loadMailTemplates(tenantId);
result.add(EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, result.add(EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS,
EdgeEventActionType.UPDATED, null, JacksonUtil.OBJECT_MAPPER.valueToTree(systemMailTemplates))); EdgeEventActionType.UPDATED, null, JacksonUtil.OBJECT_MAPPER.valueToTree(systemMailTemplates)));
@ -97,7 +97,7 @@ public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher {
return new PageData<>(result, 1, result.size(), false); return new PageData<>(result, 1, result.size(), false);
} }
private AdminSettings loadMailTemplates() throws Exception { private AdminSettings loadMailTemplates(TenantId tenantId) throws Exception {
Map<String, Object> mailTemplates = new HashMap<>(); Map<String, Object> mailTemplates = new HashMap<>();
for (String templatesName : templatesNames) { for (String templatesName : templatesNames) {
Template template = freemarkerConfig.getTemplate(templatesName); Template template = freemarkerConfig.getTemplate(templatesName);
@ -107,7 +107,7 @@ public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher {
if (mailTemplate != null) { if (mailTemplate != null) {
mailTemplates.put(name, mailTemplate); mailTemplates.put(name, mailTemplate);
} else { } else {
log.error("Can't load mail template from file {}", template.getName()); log.error("[{}] Can't load mail template from file {}", tenantId, template.getName());
} }
} }
} }

37
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java

@ -20,15 +20,21 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest; import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest;
import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest; import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
@ -45,7 +51,7 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor {
EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); EntityType.valueOf(alarmUpdateMsg.getOriginatorType()));
AlarmId alarmId = new AlarmId(new UUID(alarmUpdateMsg.getIdMSB(), alarmUpdateMsg.getIdLSB())); AlarmId alarmId = new AlarmId(new UUID(alarmUpdateMsg.getIdMSB(), alarmUpdateMsg.getIdLSB()));
if (originatorId == null) { if (originatorId == null) {
log.warn("Originator not found for the alarm msg {}", alarmUpdateMsg); log.warn("[{}] Originator not found for the alarm msg {}", tenantId, alarmUpdateMsg);
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }
try { try {
@ -129,13 +135,38 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor {
case ALARM_CLEAR: case ALARM_CLEAR:
Alarm alarm = alarmService.findAlarmById(tenantId, alarmId); Alarm alarm = alarmService.findAlarmById(tenantId, alarmId);
if (alarm != null) { if (alarm != null) {
return alarmMsgConstructor.constructAlarmUpdatedMsg(tenantId, msgType, alarm); return alarmMsgConstructor.constructAlarmUpdatedMsg(msgType, alarm, findOriginatorEntityName(tenantId, alarm));
} }
break; break;
case DELETED: case DELETED:
Alarm deletedAlarm = JacksonUtil.OBJECT_MAPPER.convertValue(body, Alarm.class); Alarm deletedAlarm = JacksonUtil.OBJECT_MAPPER.convertValue(body, Alarm.class);
return alarmMsgConstructor.constructAlarmUpdatedMsg(tenantId, msgType, deletedAlarm); return alarmMsgConstructor.constructAlarmUpdatedMsg(msgType, deletedAlarm, findOriginatorEntityName(tenantId, deletedAlarm));
} }
return null; return null;
} }
private String findOriginatorEntityName(TenantId tenantId, Alarm alarm) {
String entityName = null;
switch (alarm.getOriginator().getEntityType()) {
case DEVICE:
Device deviceById = deviceService.findDeviceById(tenantId, new DeviceId(alarm.getOriginator().getId()));
if (deviceById != null) {
entityName = deviceById.getName();
}
break;
case ASSET:
Asset assetById = assetService.findAssetById(tenantId, new AssetId(alarm.getOriginator().getId()));
if (assetById != null) {
entityName = assetById.getName();
}
break;
case ENTITY_VIEW:
EntityView entityViewById = entityViewService.findEntityViewById(tenantId, new EntityViewId(alarm.getOriginator().getId()));
if (entityViewById != null) {
entityName = entityViewById.getName();
}
break;
}
return entityName;
}
} }

8
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetProfileEdgeProcessor.java

@ -69,7 +69,7 @@ public class AssetProfileEdgeProcessor extends BaseAssetProfileProcessor {
return handleUnsupportedMsgType(assetProfileUpdateMsg.getMsgType()); return handleUnsupportedMsgType(assetProfileUpdateMsg.getMsgType());
} }
} catch (DataValidationException e) { } catch (DataValidationException e) {
log.warn("Failed to process AssetProfileUpdateMsg from Edge [{}]", assetProfileUpdateMsg, e); log.warn("[{}] Failed to process AssetProfileUpdateMsg from Edge [{}]", tenantId, assetProfileUpdateMsg, e);
return Futures.immediateFailedFuture(e); return Futures.immediateFailedFuture(e);
} finally { } finally {
edgeSynchronizationManager.getSync().remove(); edgeSynchronizationManager.getSync().remove();
@ -98,16 +98,16 @@ public class AssetProfileEdgeProcessor extends BaseAssetProfileProcessor {
tbClusterService.pushMsgToRuleEngine(tenantId, assetProfileId, tbMsg, new TbQueueCallback() { tbClusterService.pushMsgToRuleEngine(tenantId, assetProfileId, tbMsg, new TbQueueCallback() {
@Override @Override
public void onSuccess(TbQueueMsgMetadata metadata) { public void onSuccess(TbQueueMsgMetadata metadata) {
log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", assetProfile); log.debug("[{}] Successfully send ENTITY_CREATED EVENT to rule engine [{}]", tenantId, assetProfile);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.warn("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", assetProfile, t); log.warn("[{}] Failed to send ENTITY_CREATED EVENT to rule engine [{}]", tenantId, assetProfile, t);
} }
}); });
} catch (JsonProcessingException | IllegalArgumentException e) { } catch (JsonProcessingException | IllegalArgumentException e) {
log.warn("[{}] Failed to push asset profile action to rule engine: {}", assetProfileId, DataConstants.ENTITY_CREATED, e); log.warn("[{}][{}] Failed to push asset profile action to rule engine: {}", tenantId, assetProfileId, DataConstants.ENTITY_CREATED, e);
} }
} }

7
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProcessor.java

@ -49,8 +49,8 @@ public abstract class BaseAssetProcessor extends BaseEdgeProcessor {
Asset assetByName = assetService.findAssetByTenantIdAndName(tenantId, assetName); Asset assetByName = assetService.findAssetByTenantIdAndName(tenantId, assetName);
if (assetByName != null && !assetByName.getId().equals(assetId)) { if (assetByName != null && !assetByName.getId().equals(assetId)) {
assetName = assetName + "_" + StringUtils.randomAlphanumeric(15); assetName = assetName + "_" + StringUtils.randomAlphanumeric(15);
log.warn("Asset with name {} already exists. Renaming asset name to {}", log.warn("[{}] Asset with name {} already exists. Renaming asset name to {}",
assetUpdateMsg.getName(), assetName); tenantId, assetUpdateMsg.getName(), assetName);
assetNameUpdated = true; assetNameUpdated = true;
} }
asset.setName(assetName); asset.setName(assetName);
@ -69,6 +69,9 @@ public abstract class BaseAssetProcessor extends BaseEdgeProcessor {
asset.setId(assetId); asset.setId(assetId);
} }
assetService.saveAsset(asset, false); assetService.saveAsset(asset, false);
} catch (Exception e) {
log.error("[{}] Failed to process asset update msg [{}]", tenantId, assetUpdateMsg, e);
throw e;
} finally { } finally {
assetCreationLock.unlock(); assetCreationLock.unlock();
} }

7
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/BaseAssetProfileProcessor.java

@ -46,8 +46,8 @@ public abstract class BaseAssetProfileProcessor extends BaseEdgeProcessor {
AssetProfile assetProfileByName = assetProfileService.findAssetProfileByName(tenantId, assetProfileName); AssetProfile assetProfileByName = assetProfileService.findAssetProfileByName(tenantId, assetProfileName);
if (assetProfileByName != null && !assetProfileByName.getId().equals(assetProfileId)) { if (assetProfileByName != null && !assetProfileByName.getId().equals(assetProfileId)) {
assetProfileName = assetProfileName + "_" + StringUtils.randomAlphabetic(15); assetProfileName = assetProfileName + "_" + StringUtils.randomAlphabetic(15);
log.warn("Asset profile with name {} already exists. Renaming asset profile name to {}", log.warn("[{}] Asset profile with name {} already exists. Renaming asset profile name to {}",
assetProfileUpdateMsg.getName(), assetProfileName); tenantId, assetProfileUpdateMsg.getName(), assetProfileName);
assetProfileNameUpdated = true; assetProfileNameUpdated = true;
} }
assetProfile.setName(assetProfileName); assetProfile.setName(assetProfileName);
@ -66,6 +66,9 @@ public abstract class BaseAssetProfileProcessor extends BaseEdgeProcessor {
assetProfile.setId(assetProfileId); assetProfile.setId(assetProfileId);
} }
assetProfileService.saveAssetProfile(assetProfile, false); assetProfileService.saveAssetProfile(assetProfile, false);
} catch (Exception e) {
log.error("[{}] Failed to process asset profile update msg [{}]", tenantId, assetProfileUpdateMsg, e);
throw e;
} finally { } finally {
assetCreationLock.unlock(); assetCreationLock.unlock();
} }

19
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/BaseDeviceProcessor.java

@ -61,8 +61,8 @@ public abstract class BaseDeviceProcessor extends BaseEdgeProcessor {
Device deviceByName = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); Device deviceByName = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName);
if (deviceByName != null && !deviceByName.getId().equals(deviceId)) { if (deviceByName != null && !deviceByName.getId().equals(deviceId)) {
deviceName = deviceName + "_" + StringUtils.randomAlphabetic(15); deviceName = deviceName + "_" + StringUtils.randomAlphabetic(15);
log.warn("Device with name {} already exists. Renaming device name to {}", log.warn("[{}] Device with name {} already exists. Renaming device name to {}",
deviceUpdateMsg.getName(), deviceName); tenantId, deviceUpdateMsg.getName(), deviceName);
deviceNameUpdated = true; deviceNameUpdated = true;
} }
device.setName(deviceName); device.setName(deviceName);
@ -98,7 +98,10 @@ public abstract class BaseDeviceProcessor extends BaseEdgeProcessor {
deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials);
} }
tbClusterService.onDeviceUpdated(savedDevice, created ? null : device); tbClusterService.onDeviceUpdated(savedDevice, created ? null : device);
} finally { } catch (Exception e) {
log.error("[{}] Failed to process device update msg [{}]", tenantId, deviceUpdateMsg, e);
throw e;
} finally {
deviceCreationLock.unlock(); deviceCreationLock.unlock();
} }
return Pair.of(created, deviceNameUpdated); return Pair.of(created, deviceNameUpdated);
@ -110,8 +113,8 @@ public abstract class BaseDeviceProcessor extends BaseEdgeProcessor {
return dbCallbackExecutorService.submit(() -> { return dbCallbackExecutorService.submit(() -> {
Device device = deviceService.findDeviceById(tenantId, deviceId); Device device = deviceService.findDeviceById(tenantId, deviceId);
if (device != null) { if (device != null) {
log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", log.debug("[{}] Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]",
device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); tenantId, device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue());
try { try {
edgeSynchronizationManager.getSync().set(true); edgeSynchronizationManager.getSync().set(true);
@ -122,14 +125,14 @@ public abstract class BaseDeviceProcessor extends BaseEdgeProcessor {
? deviceCredentialsUpdateMsg.getCredentialsValue() : null); ? deviceCredentialsUpdateMsg.getCredentialsValue() : null);
deviceCredentialsService.updateDeviceCredentials(tenantId, deviceCredentials); deviceCredentialsService.updateDeviceCredentials(tenantId, deviceCredentials);
} catch (Exception e) { } catch (Exception e) {
log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", log.error("[{}] Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]",
device.getName(), deviceCredentialsUpdateMsg, e); tenantId, device.getName(), deviceCredentialsUpdateMsg, e);
throw new RuntimeException(e); throw new RuntimeException(e);
} finally { } finally {
edgeSynchronizationManager.getSync().remove(); edgeSynchronizationManager.getSync().remove();
} }
} else { } else {
log.warn("Can't find device by id [{}], deviceCredentialsUpdateMsg [{}]", deviceId, deviceCredentialsUpdateMsg); log.warn("[{}] Can't find device by id [{}], deviceCredentialsUpdateMsg [{}]", tenantId, deviceId, deviceCredentialsUpdateMsg);
} }
return null; return null;
}); });

9
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/BaseDeviceProfileProcessor.java

@ -58,8 +58,8 @@ public abstract class BaseDeviceProfileProcessor extends BaseEdgeProcessor {
DeviceProfile deviceProfileByName = deviceProfileService.findDeviceProfileByName(tenantId, deviceProfileName); DeviceProfile deviceProfileByName = deviceProfileService.findDeviceProfileByName(tenantId, deviceProfileName);
if (deviceProfileByName != null && !deviceProfileByName.getId().equals(deviceProfileId)) { if (deviceProfileByName != null && !deviceProfileByName.getId().equals(deviceProfileId)) {
deviceProfileName = deviceProfileName + "_" + StringUtils.randomAlphabetic(15); deviceProfileName = deviceProfileName + "_" + StringUtils.randomAlphabetic(15);
log.warn("Device profile with name {} already exists. Renaming device profile name to {}", log.warn("[{}] Device profile with name {} already exists. Renaming device profile name to {}",
deviceProfileUpdateMsg.getName(), deviceProfileName); tenantId, deviceProfileUpdateMsg.getName(), deviceProfileName);
deviceProfileNameUpdated = true; deviceProfileNameUpdated = true;
} }
deviceProfile.setName(deviceProfileName); deviceProfile.setName(deviceProfileName);
@ -98,7 +98,10 @@ public abstract class BaseDeviceProfileProcessor extends BaseEdgeProcessor {
deviceProfile.setId(deviceProfileId); deviceProfile.setId(deviceProfileId);
} }
deviceProfileService.saveDeviceProfile(deviceProfile, false); deviceProfileService.saveDeviceProfile(deviceProfile, false);
} finally { } catch (Exception e) {
log.error("[{}] Failed to process device profile update msg [{}]", tenantId, deviceProfileUpdateMsg, e);
throw e;
} finally {
deviceCreationLock.unlock(); deviceCreationLock.unlock();
} }
return Pair.of(created, deviceProfileNameUpdated); return Pair.of(created, deviceProfileNameUpdated);

9
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceProfileEdgeProcessor.java

@ -52,7 +52,6 @@ import java.util.UUID;
@TbCoreComponent @TbCoreComponent
public class DeviceProfileEdgeProcessor extends BaseDeviceProfileProcessor { public class DeviceProfileEdgeProcessor extends BaseDeviceProfileProcessor {
public ListenableFuture<Void> processDeviceProfileMsgFromEdge(TenantId tenantId, Edge edge, DeviceProfileUpdateMsg deviceProfileUpdateMsg) { public ListenableFuture<Void> processDeviceProfileMsgFromEdge(TenantId tenantId, Edge edge, DeviceProfileUpdateMsg deviceProfileUpdateMsg) {
log.trace("[{}] executing processDeviceProfileMsgFromEdge [{}] from edge [{}]", tenantId, deviceProfileUpdateMsg, edge.getName()); log.trace("[{}] executing processDeviceProfileMsgFromEdge [{}] from edge [{}]", tenantId, deviceProfileUpdateMsg, edge.getName());
DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(deviceProfileUpdateMsg.getIdMSB(), deviceProfileUpdateMsg.getIdLSB())); DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(deviceProfileUpdateMsg.getIdMSB(), deviceProfileUpdateMsg.getIdLSB()));
@ -70,7 +69,7 @@ public class DeviceProfileEdgeProcessor extends BaseDeviceProfileProcessor {
return handleUnsupportedMsgType(deviceProfileUpdateMsg.getMsgType()); return handleUnsupportedMsgType(deviceProfileUpdateMsg.getMsgType());
} }
} catch (DataValidationException e) { } catch (DataValidationException e) {
log.warn("Failed to process DeviceProfileUpdateMsg from Edge [{}]", deviceProfileUpdateMsg, e); log.warn("[{}] Failed to process DeviceProfileUpdateMsg from Edge [{}]", tenantId, deviceProfileUpdateMsg, e);
return Futures.immediateFailedFuture(e); return Futures.immediateFailedFuture(e);
} finally { } finally {
edgeSynchronizationManager.getSync().remove(); edgeSynchronizationManager.getSync().remove();
@ -99,16 +98,16 @@ public class DeviceProfileEdgeProcessor extends BaseDeviceProfileProcessor {
tbClusterService.pushMsgToRuleEngine(tenantId, deviceProfileId, tbMsg, new TbQueueCallback() { tbClusterService.pushMsgToRuleEngine(tenantId, deviceProfileId, tbMsg, new TbQueueCallback() {
@Override @Override
public void onSuccess(TbQueueMsgMetadata metadata) { public void onSuccess(TbQueueMsgMetadata metadata) {
log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", deviceProfile); log.debug("[{}] Successfully send ENTITY_CREATED EVENT to rule engine [{}]", tenantId, deviceProfile);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.warn("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", deviceProfile, t); log.warn("[{}] Failed to send ENTITY_CREATED EVENT to rule engine [{}]", tenantId, deviceProfile, t);
} }
}); });
} catch (JsonProcessingException | IllegalArgumentException e) { } catch (JsonProcessingException | IllegalArgumentException e) {
log.warn("[{}] Failed to push device profile action to rule engine: {}", deviceProfileId, DataConstants.ENTITY_CREATED, e); log.warn("[{}][{}] Failed to push device profile action to rule engine: {}", tenantId, deviceProfileId, DataConstants.ENTITY_CREATED, e);
} }
} }

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeProcessor.java

@ -85,7 +85,7 @@ public class EdgeProcessor extends BaseEdgeProcessor {
do { do {
pageData = userService.findCustomerUsers(tenantId, customerId, pageLink); pageData = userService.findCustomerUsers(tenantId, customerId, pageLink);
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] user(s) are going to be added to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}][{}][{}] user(s) are going to be added to edge.", tenantId, edge.getId(), pageData.getData().size());
for (User user : pageData.getData()) { for (User user : pageData.getData()) {
futures.add(saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null)); futures.add(saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null));
} }
@ -108,7 +108,7 @@ public class EdgeProcessor extends BaseEdgeProcessor {
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during processing edge event", e); log.error("[{}] Exception during processing edge event", tenantId, e);
return Futures.immediateFailedFuture(e); return Futures.immediateFailedFuture(e);
} }
} }

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/entityview/BaseEntityViewProcessor.java

@ -49,8 +49,8 @@ public abstract class BaseEntityViewProcessor extends BaseEdgeProcessor {
EntityView entityViewByName = entityViewService.findEntityViewByTenantIdAndName(tenantId, entityViewName); EntityView entityViewByName = entityViewService.findEntityViewByTenantIdAndName(tenantId, entityViewName);
if (entityViewByName != null && !entityViewByName.getId().equals(entityViewId)) { if (entityViewByName != null && !entityViewByName.getId().equals(entityViewId)) {
entityViewName = entityViewName + "_" + StringUtils.randomAlphanumeric(15); entityViewName = entityViewName + "_" + StringUtils.randomAlphanumeric(15);
log.warn("Entity view with name {} already exists. Renaming entity view name to {}", log.warn("[{}] Entity view with name {} already exists. Renaming entity view name to {}",
entityViewUpdateMsg.getName(), entityViewName); tenantId, entityViewUpdateMsg.getName(), entityViewName);
entityViewNameUpdated = true; entityViewNameUpdated = true;
} }
entityView.setName(entityViewName); entityView.setName(entityViewName);

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java

@ -59,7 +59,7 @@ public abstract class BaseRelationProcessor extends BaseEdgeProcessor {
relationService.saveRelation(tenantId, entityRelation); relationService.saveRelation(tenantId, entityRelation);
break; break;
} else { } else {
log.warn("Skipping relating update msg because from/to entity doesn't exists on edge, {}", relationUpdateMsg); log.warn("[{}] Skipping relating update msg because from/to entity doesn't exists on edge, {}", tenantId, relationUpdateMsg);
break; break;
} }
case ENTITY_DELETED_RPC_MESSAGE: case ENTITY_DELETED_RPC_MESSAGE:

25
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java

@ -130,7 +130,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
result.add(processAttributeDeleteMsg(tenantId, entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); result.add(processAttributeDeleteMsg(tenantId, entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType()));
} }
} else { } else {
log.warn("Skipping telemetry update msg because entity doesn't exists on edge, {}", entityData); log.warn("[{}] Skipping telemetry update msg because entity doesn't exists on edge, {}", tenantId, entityData);
} }
return result; return result;
} }
@ -172,7 +172,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
} }
break; break;
default: default:
log.debug("Using empty metadata for entityId [{}]", entityId); log.debug("[{}] Using empty metadata for entityId [{}]", tenantId, entityId);
break; break;
} }
return new ImmutablePair<>(metaData, customerId != null ? customerId : new CustomerId(ModelConstants.NULL_UUID)); return new ImmutablePair<>(metaData, customerId != null ? customerId : new CustomerId(ModelConstants.NULL_UUID));
@ -193,7 +193,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Can't process post telemetry [{}]", msg, t); log.error("[{}] Can't process post telemetry [{}]", tenantId, msg, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}); });
@ -207,7 +207,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
if (EntityType.DEVICE.equals(entityId.getEntityType())) { if (EntityType.DEVICE.equals(entityId.getEntityType())) {
DeviceProfile deviceProfile = deviceProfileCache.get(tenantId, new DeviceId(entityId.getId())); DeviceProfile deviceProfile = deviceProfileCache.get(tenantId, new DeviceId(entityId.getId()));
if (deviceProfile == null) { if (deviceProfile == null) {
log.warn("[{}] Device profile is null!", entityId); log.warn("[{}][{}] Device profile is null!", tenantId, entityId);
} else { } else {
ruleChainId = deviceProfile.getDefaultRuleChainId(); ruleChainId = deviceProfile.getDefaultRuleChainId();
queueName = deviceProfile.getDefaultQueueName(); queueName = deviceProfile.getDefaultQueueName();
@ -215,7 +215,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
} else if (EntityType.ASSET.equals(entityId.getEntityType())) { } else if (EntityType.ASSET.equals(entityId.getEntityType())) {
AssetProfile assetProfile = assetProfileCache.get(tenantId, new AssetId(entityId.getId())); AssetProfile assetProfile = assetProfileCache.get(tenantId, new AssetId(entityId.getId()));
if (assetProfile == null) { if (assetProfile == null) {
log.warn("[{}] Asset profile is null!", entityId); log.warn("[{}][{}] Asset profile is null!", tenantId, entityId);
} else { } else {
ruleChainId = assetProfile.getDefaultRuleChainId(); ruleChainId = assetProfile.getDefaultRuleChainId();
queueName = assetProfile.getDefaultQueueName(); queueName = assetProfile.getDefaultQueueName();
@ -237,7 +237,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Can't process post attributes [{}]", msg, t); log.error("[{}] Can't process post attributes [{}]", tenantId, msg, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}); });
@ -267,7 +267,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Can't process attributes update [{}]", msg, t); log.error("[{}] Can't process attributes update [{}]", tenantId, msg, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}); });
@ -275,7 +275,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Can't process attributes update [{}]", msg, t); log.error("[{}] Can't process attributes update [{}]", tenantId, msg, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}); });
@ -300,7 +300,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); log.error("[{}] Can't process attribute delete msg [{}]", tenantId, attributeDeleteMsg, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}); });
@ -311,7 +311,8 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
} }
public EntityDataProto convertTelemetryEventToEntityDataProto(EntityType entityType, public EntityDataProto convertTelemetryEventToEntityDataProto(TenantId tenantId,
EntityType entityType,
UUID entityUUID, UUID entityUUID,
EdgeEventActionType actionType, EdgeEventActionType actionType,
JsonNode body) throws JsonProcessingException { JsonNode body) throws JsonProcessingException {
@ -342,11 +343,11 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
entityId = new EdgeId(entityUUID); entityId = new EdgeId(entityUUID);
break; break;
default: default:
log.warn("Unsupported edge event type [{}]", entityType); log.warn("[{}] Unsupported edge event type [{}]", tenantId, entityType);
return null; return null;
} }
JsonElement entityData = JsonParser.parseString(JacksonUtil.OBJECT_MAPPER.writeValueAsString(body)); JsonElement entityData = JsonParser.parseString(JacksonUtil.OBJECT_MAPPER.writeValueAsString(body));
return entityDataMsgConstructor.constructEntityDataMsg(entityId, actionType, entityData); return entityDataMsgConstructor.constructEntityDataMsg(tenantId, entityId, actionType, entityData);
} }
} }

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/TelemetryEdgeProcessor.java

@ -38,7 +38,8 @@ public class TelemetryEdgeProcessor extends BaseTelemetryProcessor {
public DownlinkMsg convertTelemetryEventToDownlink(EdgeEvent edgeEvent) throws JsonProcessingException { public DownlinkMsg convertTelemetryEventToDownlink(EdgeEvent edgeEvent) throws JsonProcessingException {
EntityType entityType = EntityType.valueOf(edgeEvent.getType().name()); EntityType entityType = EntityType.valueOf(edgeEvent.getType().name());
EntityDataProto entityDataProto = convertTelemetryEventToEntityDataProto(entityType, edgeEvent.getEntityId(), EntityDataProto entityDataProto = convertTelemetryEventToEntityDataProto(
edgeEvent.getTenantId(), entityType, edgeEvent.getEntityId(),
edgeEvent.getAction(), edgeEvent.getBody()); edgeEvent.getAction(), edgeEvent.getBody());
return DownlinkMsg.newBuilder() return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .setDownlinkMsgId(EdgeUtils.nextPositiveInt())

14
application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java

@ -177,7 +177,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
entityData.put("kv", attributes); entityData.put("kv", attributes);
entityData.put("scope", scope); entityData.put("scope", scope);
JsonNode body = JacksonUtil.OBJECT_MAPPER.valueToTree(entityData); JsonNode body = JacksonUtil.OBJECT_MAPPER.valueToTree(entityData);
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body); log.debug("[{}] Sending attributes data msg, entityId [{}], attributes [{}]", tenantId, entityId, body);
future = saveEdgeEvent(tenantId, edge.getId(), entityType, EdgeEventActionType.ATTRIBUTES_UPDATED, entityId, body); future = saveEdgeEvent(tenantId, edge.getId(), entityType, EdgeEventActionType.ATTRIBUTES_UPDATED, entityId, body);
} else { } else {
future = Futures.immediateFuture(null); future = Futures.immediateFuture(null);
@ -185,7 +185,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
} }
return Futures.transformAsync(future, v -> processLatestTimeseriesAndAddToEdgeQueue(tenantId, entityId, edge, entityType), dbCallbackExecutorService); return Futures.transformAsync(future, v -> processLatestTimeseriesAndAddToEdgeQueue(tenantId, entityId, edge, entityType), dbCallbackExecutorService);
} catch (Exception e) { } catch (Exception e) {
String errMsg = String.format("[%s] Failed to save attribute updates to the edge [%s]", edge.getId(), attributesRequestMsg); String errMsg = String.format("[%s][%s] Failed to save attribute updates to the edge [%s]", tenantId, edge.getId(), attributesRequestMsg);
log.error(errMsg, e); log.error(errMsg, e);
return Futures.immediateFailedFuture(new RuntimeException(errMsg, e)); return Futures.immediateFailedFuture(new RuntimeException(errMsg, e));
} }
@ -239,7 +239,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
if (relationsList != null && !relationsList.isEmpty()) { if (relationsList != null && !relationsList.isEmpty()) {
List<ListenableFuture<Void>> futures = new ArrayList<>(); List<ListenableFuture<Void>> futures = new ArrayList<>();
for (List<EntityRelation> entityRelations : relationsList) { for (List<EntityRelation> entityRelations : relationsList) {
log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size()); log.trace("[{}][{}][{}][{}] relation(s) are going to be pushed to edge.", tenantId, edge.getId(), entityId, entityRelations.size());
for (EntityRelation relation : entityRelations) { for (EntityRelation relation : entityRelations) {
try { try {
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) &&
@ -252,7 +252,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
JacksonUtil.OBJECT_MAPPER.valueToTree(relation))); JacksonUtil.OBJECT_MAPPER.valueToTree(relation)));
} }
} catch (Exception e) { } catch (Exception e) {
String errMsg = String.format("[%s] Exception during loading relation [%s] to edge on sync!", edge.getId(), relation); String errMsg = String.format("[%s][%s] Exception during loading relation [%s] to edge on sync!", tenantId, edge.getId(), relation);
log.error(errMsg, e); log.error(errMsg, e);
futureToSet.setException(new RuntimeException(errMsg, e)); futureToSet.setException(new RuntimeException(errMsg, e));
return; return;
@ -267,7 +267,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
@Override @Override
public void onFailure(Throwable throwable) { public void onFailure(Throwable throwable) {
String errMsg = String.format("[%s] Exception during saving edge events [%s]!", edge.getId(), relationRequestMsg); String errMsg = String.format("[%s][%s] Exception during saving edge events [%s]!", tenantId, edge.getId(), relationRequestMsg);
log.error(errMsg, throwable); log.error(errMsg, throwable);
futureToSet.setException(new RuntimeException(errMsg, throwable)); futureToSet.setException(new RuntimeException(errMsg, throwable));
} }
@ -276,7 +276,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
futureToSet.set(null); futureToSet.set(null);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading relation(s) to edge on sync!", e); log.error("[{}] Exception during loading relation(s) to edge on sync!", tenantId, e);
futureToSet.setException(e); futureToSet.setException(e);
} }
} }
@ -374,7 +374,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("Exception during loading relation to edge on sync!", t); log.error("[{}] Exception during loading relation to edge on sync!", tenantId, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}, dbCallbackExecutorService); }, dbCallbackExecutorService);

54
application/src/test/java/org/thingsboard/server/edge/AssetProfileEdgeTest.java

@ -20,6 +20,7 @@ import com.google.protobuf.AbstractMessage;
import com.google.protobuf.ByteString; import com.google.protobuf.ByteString;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.asset.AssetProfile;
import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.DashboardId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
@ -30,6 +31,7 @@ import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -85,7 +87,7 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest {
@Test @Test
public void testSendAssetProfileToCloud() throws Exception { public void testSendAssetProfileToCloud() throws Exception {
RuleChainId ruleChainId = createEdgeRuleChainAndAssignToEdge("Asset Profile Rule Chain"); RuleChainId edgeRuleChainId = createEdgeRuleChainAndAssignToEdge("Asset Profile Rule Chain");
DashboardId dashboardId = createDashboardAndAssignToEdge("Asset Profile Dashboard"); DashboardId dashboardId = createDashboardAndAssignToEdge("Asset Profile Dashboard");
UUID uuid = Uuids.timeBased(); UUID uuid = Uuids.timeBased();
@ -96,8 +98,8 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest {
assetProfileUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); assetProfileUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
assetProfileUpdateMsgBuilder.setName("Asset Profile On Edge"); assetProfileUpdateMsgBuilder.setName("Asset Profile On Edge");
assetProfileUpdateMsgBuilder.setDefault(false); assetProfileUpdateMsgBuilder.setDefault(false);
assetProfileUpdateMsgBuilder.setDefaultRuleChainIdMSB(ruleChainId.getId().getMostSignificantBits()); assetProfileUpdateMsgBuilder.setDefaultRuleChainIdMSB(edgeRuleChainId.getId().getMostSignificantBits());
assetProfileUpdateMsgBuilder.setDefaultRuleChainIdLSB(ruleChainId.getId().getLeastSignificantBits()); assetProfileUpdateMsgBuilder.setDefaultRuleChainIdLSB(edgeRuleChainId.getId().getLeastSignificantBits());
assetProfileUpdateMsgBuilder.setDefaultDashboardIdMSB(dashboardId.getId().getMostSignificantBits()); assetProfileUpdateMsgBuilder.setDefaultDashboardIdMSB(dashboardId.getId().getMostSignificantBits());
assetProfileUpdateMsgBuilder.setDefaultDashboardIdLSB(dashboardId.getId().getLeastSignificantBits()); assetProfileUpdateMsgBuilder.setDefaultDashboardIdLSB(dashboardId.getId().getLeastSignificantBits());
assetProfileUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); assetProfileUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
@ -117,6 +119,9 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest {
AssetProfile assetProfile = doGet("/api/assetProfile/" + uuid, AssetProfile.class); AssetProfile assetProfile = doGet("/api/assetProfile/" + uuid, AssetProfile.class);
Assert.assertNotNull(assetProfile); Assert.assertNotNull(assetProfile);
Assert.assertEquals("Asset Profile On Edge", assetProfile.getName()); Assert.assertEquals("Asset Profile On Edge", assetProfile.getName());
Assert.assertEquals(dashboardId, assetProfile.getDefaultDashboardId());
Assert.assertNull(assetProfile.getDefaultRuleChainId());
Assert.assertEquals(edgeRuleChainId, assetProfile.getDefaultEdgeRuleChainId());
// delete profile // delete profile
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
@ -132,6 +137,47 @@ public class AssetProfileEdgeTest extends AbstractEdgeTest {
// cleanup // cleanup
unAssignFromEdgeAndDeleteDashboard(dashboardId); unAssignFromEdgeAndDeleteDashboard(dashboardId);
unAssignFromEdgeAndDeleteRuleChain(ruleChainId); unAssignFromEdgeAndDeleteRuleChain(edgeRuleChainId);
}
@Test
public void testSendAssetProfileToCloudWithNameThatAlreadyExistsOnCloud() throws Exception {
String assetProfileOnCloudName = StringUtils.randomAlphanumeric(15);
edgeImitator.expectMessageAmount(1);
AssetProfile assetProfileOnCloud = this.createAssetProfile(assetProfileOnCloudName);
assetProfileOnCloud = doPost("/api/assetProfile", assetProfileOnCloud, AssetProfile.class);
Assert.assertTrue(edgeImitator.waitForMessages());
UUID uuid = Uuids.timeBased();
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AssetProfileUpdateMsg.Builder assetProfileUpdateMsgBuilder = AssetProfileUpdateMsg.newBuilder();
assetProfileUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
assetProfileUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
assetProfileUpdateMsgBuilder.setName(assetProfileOnCloudName);
assetProfileUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
uplinkMsgBuilder.addAssetProfileUpdateMsg(assetProfileUpdateMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
Assert.assertTrue(edgeImitator.waitForResponses());
Assert.assertTrue(edgeImitator.waitForMessages());
Optional<AssetProfileUpdateMsg> assetProfileUpdateMsgOpt = edgeImitator.findMessageByType(AssetProfileUpdateMsg.class);
Assert.assertTrue(assetProfileUpdateMsgOpt.isPresent());
AssetProfileUpdateMsg latestAssetProfileUpdateMsg = assetProfileUpdateMsgOpt.get();
Assert.assertNotEquals(assetProfileOnCloudName, latestAssetProfileUpdateMsg.getName());
Assert.assertNotEquals(assetProfileOnCloud.getUuidId(), uuid);
AssetProfile assetProfile = doGet("/api/assetProfile/" + uuid, AssetProfile.class);
Assert.assertNotNull(assetProfile);
Assert.assertNotEquals(assetProfileOnCloudName, assetProfile.getName());
} }
} }

10
application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java

@ -33,7 +33,6 @@ import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.gen.edge.v1.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.v1.DashboardUpdateMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkMsg;
import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
@ -45,12 +44,18 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
@DaoSqlTest @DaoSqlTest
public class DashboardEdgeTest extends AbstractEdgeTest { public class DashboardEdgeTest extends AbstractEdgeTest {
private static final int MOBILE_ORDER = 5;
private static final String IMAGE = "data:image/png;base64,iVBORw0KGgoA";
@Test @Test
public void testDashboards() throws Exception { public void testDashboards() throws Exception {
// create dashboard and assign to edge // create dashboard and assign to edge
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
Dashboard dashboard = new Dashboard(); Dashboard dashboard = new Dashboard();
dashboard.setTitle("Edge Test Dashboard"); dashboard.setTitle("Edge Test Dashboard");
dashboard.setMobileHide(true);
dashboard.setImage(IMAGE);
dashboard.setMobileOrder(MOBILE_ORDER);
Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class); Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class);
doPost("/api/edge/" + edge.getUuidId() doPost("/api/edge/" + edge.getUuidId()
+ "/dashboard/" + savedDashboard.getUuidId(), Dashboard.class); + "/dashboard/" + savedDashboard.getUuidId(), Dashboard.class);
@ -62,6 +67,9 @@ public class DashboardEdgeTest extends AbstractEdgeTest {
Assert.assertEquals(savedDashboard.getUuidId().getMostSignificantBits(), dashboardUpdateMsg.getIdMSB()); Assert.assertEquals(savedDashboard.getUuidId().getMostSignificantBits(), dashboardUpdateMsg.getIdMSB());
Assert.assertEquals(savedDashboard.getUuidId().getLeastSignificantBits(), dashboardUpdateMsg.getIdLSB()); Assert.assertEquals(savedDashboard.getUuidId().getLeastSignificantBits(), dashboardUpdateMsg.getIdLSB());
Assert.assertEquals(savedDashboard.getTitle(), dashboardUpdateMsg.getTitle()); Assert.assertEquals(savedDashboard.getTitle(), dashboardUpdateMsg.getTitle());
Assert.assertTrue(dashboardUpdateMsg.getMobileHide());
Assert.assertEquals(IMAGE, dashboardUpdateMsg.getImage());
Assert.assertEquals(MOBILE_ORDER, dashboardUpdateMsg.getMobileOrder());
testAutoGeneratedCodeByProtobuf(dashboardUpdateMsg); testAutoGeneratedCodeByProtobuf(dashboardUpdateMsg);
// update dashboard // update dashboard

47
application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java

@ -25,7 +25,7 @@ import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceProfileType;
import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.device.data.PowerMode; import org.thingsboard.server.common.data.device.data.PowerMode;
import org.thingsboard.server.common.data.device.data.PowerSavingConfiguration; import org.thingsboard.server.common.data.device.data.PowerSavingConfiguration;
import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration;
@ -303,7 +303,7 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
UplinkResponseMsg latestResponseMsg = edgeImitator.getLatestResponseMsg(); UplinkResponseMsg latestResponseMsg = edgeImitator.getLatestResponseMsg();
Assert.assertTrue(latestResponseMsg.getSuccess()); Assert.assertTrue(latestResponseMsg.getSuccess());
AssetProfile deviceProfile = doGet("/api/deviceProfile/" + uuid, AssetProfile.class); DeviceProfile deviceProfile = doGet("/api/deviceProfile/" + uuid, DeviceProfile.class);
Assert.assertNotNull(deviceProfile); Assert.assertNotNull(deviceProfile);
Assert.assertEquals("Device Profile On Edge", deviceProfile.getName()); Assert.assertEquals("Device Profile On Edge", deviceProfile.getName());
@ -412,4 +412,47 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
transportConfiguration.setCoapDeviceTypeConfiguration(coapDeviceTypeConfiguration); transportConfiguration.setCoapDeviceTypeConfiguration(coapDeviceTypeConfiguration);
return transportConfiguration; return transportConfiguration;
} }
@Test
public void testSendDeviceProfileToCloudWithNameThatAlreadyExistsOnCloud() throws Exception {
String deviceProfileOnCloudName = StringUtils.randomAlphanumeric(15);
edgeImitator.expectMessageAmount(1);
DeviceProfile deviceProfileOnCloud = this.createDeviceProfile(deviceProfileOnCloudName);
deviceProfileOnCloud = doPost("/api/deviceProfile", deviceProfileOnCloud, DeviceProfile.class);
Assert.assertTrue(edgeImitator.waitForMessages());
UUID uuid = Uuids.timeBased();
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
DeviceProfileUpdateMsg.Builder deviceProfileUpdateMsgBuilder = DeviceProfileUpdateMsg.newBuilder();
deviceProfileUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
deviceProfileUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
deviceProfileUpdateMsgBuilder.setName(deviceProfileOnCloudName);
deviceProfileUpdateMsgBuilder.setType(DeviceProfileType.DEFAULT.name());
deviceProfileUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
deviceProfileUpdateMsgBuilder.setProfileDataBytes(ByteString.copyFrom(dataDecodingEncodingService.encode(deviceProfileOnCloud.getProfileData())));
uplinkMsgBuilder.addDeviceProfileUpdateMsg(deviceProfileUpdateMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
Assert.assertTrue(edgeImitator.waitForResponses());
Assert.assertTrue(edgeImitator.waitForMessages());
Optional<DeviceProfileUpdateMsg> deviceProfileUpdateMsgOpt = edgeImitator.findMessageByType(DeviceProfileUpdateMsg.class);
Assert.assertTrue(deviceProfileUpdateMsgOpt.isPresent());
DeviceProfileUpdateMsg latestDeviceProfileUpdateMsg = deviceProfileUpdateMsgOpt.get();
Assert.assertNotEquals(deviceProfileOnCloudName, latestDeviceProfileUpdateMsg.getName());
Assert.assertNotEquals(deviceProfileOnCloud.getUuidId(), uuid);
DeviceProfile deviceProfile = doGet("/api/deviceProfile/" + uuid, DeviceProfile.class);
Assert.assertNotNull(deviceProfile);
Assert.assertNotEquals(deviceProfileOnCloudName, deviceProfile.getName());
}
} }

2
application/src/test/java/org/thingsboard/server/edge/WidgetEdgeTest.java

@ -58,6 +58,7 @@ public class WidgetEdgeTest extends AbstractEdgeTest {
ObjectNode descriptor = JacksonUtil.newObjectNode(); ObjectNode descriptor = JacksonUtil.newObjectNode();
descriptor.put("key", "value"); descriptor.put("key", "value");
widgetType.setDescriptor(descriptor); widgetType.setDescriptor(descriptor);
widgetType.setDeprecated(true);
WidgetType savedWidgetType = doPost("/api/widgetType", widgetType, WidgetType.class); WidgetType savedWidgetType = doPost("/api/widgetType", widgetType, WidgetType.class);
Assert.assertTrue(edgeImitator.waitForMessages()); Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage(); latestMessage = edgeImitator.getLatestMessage();
@ -68,6 +69,7 @@ public class WidgetEdgeTest extends AbstractEdgeTest {
Assert.assertEquals(savedWidgetType.getUuidId().getLeastSignificantBits(), widgetTypeUpdateMsg.getIdLSB()); Assert.assertEquals(savedWidgetType.getUuidId().getLeastSignificantBits(), widgetTypeUpdateMsg.getIdLSB());
Assert.assertEquals(savedWidgetType.getFqn(), widgetTypeUpdateMsg.getFqn()); Assert.assertEquals(savedWidgetType.getFqn(), widgetTypeUpdateMsg.getFqn());
Assert.assertEquals(savedWidgetType.getName(), widgetTypeUpdateMsg.getName()); Assert.assertEquals(savedWidgetType.getName(), widgetTypeUpdateMsg.getName());
Assert.assertTrue(widgetTypeUpdateMsg.getDeprecated());
Assert.assertEquals(JacksonUtil.toJsonNode(widgetTypeUpdateMsg.getDescriptorJson()), savedWidgetType.getDescriptor()); Assert.assertEquals(JacksonUtil.toJsonNode(widgetTypeUpdateMsg.getDescriptorJson()), savedWidgetType.getDescriptor());
// update widget bundle // update widget bundle

Loading…
Cancel
Save