Browse Source

Merge pull request #8344 from volodymyr-babak/edge-activity-to-rule-chain

[3.5] Push edge connect/disconnect events to rule chain
pull/8351/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
ff71374d8e
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 41
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  2. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java
  3. 43
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNodeTest.java

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

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.service.edge.rpc;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import io.grpc.Server;
@ -25,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants;
@ -35,6 +37,9 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg;
import org.thingsboard.server.common.msg.edge.EdgeSessionMsg;
import org.thingsboard.server.common.msg.edge.FromEdgeSyncResponse;
@ -66,6 +71,10 @@ import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
import static org.thingsboard.server.common.data.DataConstants.CONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.DISCONNECT_EVENT;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@Service
@Slf4j
@ConditionalOnProperty(prefix = "edges", value = "enabled", havingValue = "true")
@ -261,7 +270,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
newEventLock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis());
long lastConnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, lastConnectTs);
pushRuleEngineMessage(edgeGrpcSession.getEdge().getTenantId(), edgeId, lastConnectTs, CONNECT_EVENT);
cancelScheduleEdgeEventsCheck(edgeId);
scheduleEdgeEventsCheck(edgeGrpcSession);
}
@ -365,7 +376,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void onEdgeDisconnect(EdgeId edgeId) {
log.info("[{}] edge disconnected!", edgeId);
sessions.remove(edgeId);
EdgeGrpcSession removed = sessions.remove(edgeId);
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
@ -374,7 +385,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
newEventLock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis());
long lastDisconnectTs = System.currentTimeMillis();
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, lastDisconnectTs);
pushRuleEngineMessage(removed.getEdge().getTenantId(), edgeId, lastDisconnectTs, DISCONNECT_EVENT);
cancelScheduleEdgeEventsCheck(edgeId);
}
@ -423,4 +436,26 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.warn("[{}] Failed to update attribute [{}] with value [{}]", edgeId, key, value, t);
}
}
private void pushRuleEngineMessage(TenantId tenantId, EdgeId edgeId, long ts, String msgType) {
try {
ObjectNode edgeState = JacksonUtil.OBJECT_MAPPER.createObjectNode();
if (msgType.equals(CONNECT_EVENT)) {
edgeState.put(DefaultDeviceStateService.ACTIVITY_STATE, true);
edgeState.put(DefaultDeviceStateService.LAST_CONNECT_TIME, ts);
} else {
edgeState.put(DefaultDeviceStateService.ACTIVITY_STATE, false);
edgeState.put(DefaultDeviceStateService.LAST_DISCONNECT_TIME, ts);
}
String data = JacksonUtil.toString(edgeState);
TbMsgMetaData md = new TbMsgMetaData();
if (!persistToTelemetry) {
md.putValue(DataConstants.SCOPE, SERVER_SCOPE);
}
TbMsg tbMsg = TbMsg.newMsg(msgType, edgeId, md, TbMsgDataType.JSON, data);
clusterService.pushMsgToRuleEngine(tenantId, edgeId, tbMsg, null);
} catch (Exception e) {
log.warn("[{}][{}] Failed to push {}", tenantId, edgeId, msgType, e);
}
}
}

22
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java

@ -77,9 +77,9 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
EdgeEventActionType actionType = getAlarmActionType(msg);
return buildEvent(ctx.getTenantId(), actionType, getUUIDFromMsgData(msg), getAlarmEventType(), null);
} else {
EdgeEventActionType actionType = getEdgeEventActionTypeByMsgType(msgType);
Map<String, Object> entityBody = new HashMap<>();
Map<String, String> metadata = msg.getMetaData().getData();
EdgeEventActionType actionType = getEdgeEventActionTypeByMsgType(msgType, metadata);
Map<String, Object> entityBody = new HashMap<>();
JsonNode dataJson = JacksonUtil.toJsonNode(msg.getData());
switch (actionType) {
case ATTRIBUTES_UPDATED:
@ -148,7 +148,7 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
return scope;
}
protected EdgeEventActionType getEdgeEventActionTypeByMsgType(String msgType) {
protected EdgeEventActionType getEdgeEventActionTypeByMsgType(String msgType, Map<String, String> metadata) {
EdgeEventActionType actionType;
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)
|| DataConstants.TIMESERIES_UPDATED.equals(msgType)) {
@ -159,6 +159,16 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
actionType = EdgeEventActionType.POST_ATTRIBUTES;
} else if (DataConstants.ATTRIBUTES_DELETED.equals(msgType)) {
actionType = EdgeEventActionType.ATTRIBUTES_DELETED;
} else if (DataConstants.CONNECT_EVENT.equals(msgType)
|| DataConstants.DISCONNECT_EVENT.equals(msgType)
|| DataConstants.ACTIVITY_EVENT.equals(msgType)
|| DataConstants.INACTIVITY_EVENT.equals(msgType)) {
String scope = metadata.get(SCOPE);
if ( StringUtils.isEmpty(scope)) {
actionType = EdgeEventActionType.TIMESERIES_UPDATED;
} else {
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
}
} else {
log.warn("Unsupported msg type [{}]", msgType);
throw new IllegalArgumentException("Unsupported msg type: " + msgType);
@ -172,7 +182,11 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
|| DataConstants.ATTRIBUTES_UPDATED.equals(msgType)
|| DataConstants.ATTRIBUTES_DELETED.equals(msgType)
|| DataConstants.TIMESERIES_UPDATED.equals(msgType)
|| DataConstants.ALARM.equals(msgType);
|| DataConstants.ALARM.equals(msgType)
|| DataConstants.CONNECT_EVENT.equals(msgType)
|| DataConstants.DISCONNECT_EVENT.equals(msgType)
|| DataConstants.ACTIVITY_EVENT.equals(msgType)
|| DataConstants.INACTIVITY_EVENT.equals(msgType);
}
protected boolean isSupportedOriginator(EntityType entityType) {

43
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNodeTest.java

@ -19,6 +19,7 @@ import com.google.common.util.concurrent.SettableFuture;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentMatcher;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.junit.MockitoJUnitRunner;
@ -28,6 +29,8 @@ import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
@ -106,4 +109,44 @@ public class TbMsgPushToEdgeNodeTest {
verify(edgeEventService).saveAsync(any());
}
@Test
public void testMiscEventsProcessedAsAttributesUpdated() {
List<String> miscEvents = List.of(DataConstants.CONNECT_EVENT, DataConstants.DISCONNECT_EVENT,
DataConstants.ACTIVITY_EVENT, DataConstants.INACTIVITY_EVENT);
for (String event : miscEvents) {
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue(DataConstants.SCOPE, DataConstants.SERVER_SCOPE);
testEvent(event, metaData, EdgeEventActionType.ATTRIBUTES_UPDATED, "kv");
}
}
@Test
public void testMiscEventsProcessedAsTimeseriesUpdated() {
List<String> miscEvents = List.of(DataConstants.CONNECT_EVENT, DataConstants.DISCONNECT_EVENT,
DataConstants.ACTIVITY_EVENT, DataConstants.INACTIVITY_EVENT);
for (String event : miscEvents) {
testEvent(event, new TbMsgMetaData(), EdgeEventActionType.TIMESERIES_UPDATED, "data");
}
}
private void testEvent(String event, TbMsgMetaData metaData, EdgeEventActionType expectedType, String dataKey) {
Mockito.when(ctx.getTenantId()).thenReturn(tenantId);
Mockito.when(ctx.getEdgeService()).thenReturn(edgeService);
Mockito.when(ctx.getEdgeEventService()).thenReturn(edgeEventService);
Mockito.when(ctx.getDbCallbackExecutor()).thenReturn(dbCallbackExecutor);
Mockito.when(edgeEventService.saveAsync(any())).thenReturn(SettableFuture.create());
TbMsg msg = TbMsg.newMsg(event, new EdgeId(UUID.randomUUID()), metaData,
TbMsgDataType.JSON, "{\"lastConnectTs\":1}", null, null);
node.onMsg(ctx, msg);
ArgumentMatcher<EdgeEvent> eventArgumentMatcher = edgeEvent ->
edgeEvent.getAction().equals(expectedType)
&& edgeEvent.getBody().get(dataKey).get("lastConnectTs").asInt() == 1;
verify(edgeEventService).saveAsync(Mockito.argThat(eventArgumentMatcher));
Mockito.reset(ctx, edgeEventService);
}
}

Loading…
Cancel
Save