Browse Source

Merge RPC request/response actions

pull/7592/head
Volodymyr Babak 4 years ago
parent
commit
0a56f065f7
  1. 2
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 35
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java
  4. 21
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
  5. 2
      application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
  6. 3
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java
  7. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java
  8. 10
      ui-ngx/src/app/shared/models/edge.models.ts
  9. 2
      ui-ngx/src/assets/locale/locale.constant-en_US.json

2
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -827,7 +827,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
body.put("retries", msg.getRetries());
body.put("additionalInfo", msg.getAdditionalInfo());
EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL_REQUEST, deviceId, body);
EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body);
return Futures.transform(systemContext.getEdgeEventService().saveAsync(edgeEvent), unused -> {
systemContext.getClusterService().onEdgeEventUpdate(tenantId, edgeId);

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

@ -463,8 +463,7 @@ public final class EdgeGrpcSession implements Closeable {
case UNASSIGNED_FROM_CUSTOMER:
case CREDENTIALS_REQUEST:
case ENTITY_MERGE_REQUEST:
case RPC_CALL_REQUEST:
case RPC_CALL_RESPONSE:
case RPC_CALL:
downlinkMsg = convertEntityEventToDownlink(edgeEvent);
log.trace("[{}][{}] entity message processed [{}]", edgeEvent.getTenantId(), this.sessionId, downlinkMsg);
break;

35
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java

@ -21,7 +21,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg;
@ -97,30 +96,22 @@ public class DeviceMsgConstructor {
.setIdLSB(deviceId.getId().getLeastSignificantBits()).build();
}
public DeviceRpcCallMsg constructDeviceRpcRequestMsg(UUID deviceId, JsonNode body) {
public DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode body) {
DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body);
String method = body.get("method").asText();
String params = body.get("params").asText();
RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder();
requestBuilder.setMethod(method);
requestBuilder.setParams(params);
builder.setRequestMsg(requestBuilder.build());
return builder.build();
}
public DeviceRpcCallMsg constructDeviceRpcResponseMsg(UUID deviceId, JsonNode body) {
DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body);
RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder();
if (body.has("error")) {
responseBuilder.setError(body.get("error").asText());
if (body.has("error") || body.has("response")) {
RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder();
if (body.has("error")) {
responseBuilder.setError(body.get("error").asText());
} else {
responseBuilder.setResponse(body.get("response").asText());
}
builder.setResponseMsg(responseBuilder.build());
} else {
responseBuilder.setResponse(body.get("response").asText());
RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder();
requestBuilder.setMethod(body.get("method").asText());
requestBuilder.setParams(body.get("params").asText());
builder.setRequestMsg(requestBuilder.build());
}
builder.setResponseMsg(responseBuilder.build());
return builder.build();
}

21
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java

@ -451,10 +451,8 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
.build();
}
break;
case RPC_CALL_REQUEST:
return convertRpcCallRequestEventToDownlink(edgeEvent);
case RPC_CALL_RESPONSE:
return convertRpcCallResponseEventToDownlink(edgeEvent);
case RPC_CALL:
return convertRpcCallEventToDownlink(edgeEvent);
case CREDENTIALS_REQUEST:
return convertCredentialsRequestEventToDownlink(edgeEvent);
case ENTITY_MERGE_REQUEST:
@ -463,20 +461,11 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
return downlinkMsg;
}
private DownlinkMsg convertRpcCallRequestEventToDownlink(EdgeEvent edgeEvent) {
log.trace("Executing convertRpcCallRequestEventToDownlink, edgeEvent [{}]", edgeEvent);
private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) {
log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent);
return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcRequestMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
.build();
}
private DownlinkMsg convertRpcCallResponseEventToDownlink(EdgeEvent edgeEvent) {
log.trace("Executing convertRpcCallResponseEventToDownlink, edgeEvent [{}]", edgeEvent);
DeviceRpcCallMsg deviceRpcCallMsg = deviceMsgConstructor.constructDeviceRpcResponseMsg(edgeEvent.getEntityId(), edgeEvent.getBody());
return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceRpcCallMsg(deviceRpcCallMsg)
.addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
.build();
}

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

@ -512,7 +512,7 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
body.put("method", "test_method");
body.put("params", "{\"param1\":\"value1\"}");
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL_REQUEST,
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL,
device.getId().getId(), EdgeEventType.DEVICE, body);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent).get();

3
common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java

@ -28,8 +28,7 @@ public enum EdgeEventActionType {
UNASSIGNED_FROM_CUSTOMER,
RELATION_ADD_OR_UPDATE,
RELATION_DELETED,
RPC_CALL_REQUEST,
RPC_CALL_RESPONSE,
RPC_CALL,
ALARM_ACK,
ALARM_CLEAR,
ASSIGNED_TO_EDGE,

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java

@ -96,7 +96,7 @@ public class TbSendRPCReplyNode implements TbNode {
EdgeId edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID)));
DeviceId deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID)));
EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE,
EdgeEventActionType.RPC_CALL_RESPONSE, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body));
EdgeEventActionType.RPC_CALL, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body));
ListenableFuture<Void> future = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<Void>() {
@Override

10
ui-ngx/src/app/shared/models/edge.models.ts

@ -77,9 +77,7 @@ export enum EdgeEventActionType {
UNASSIGNED_FROM_CUSTOMER = 'UNASSIGNED_FROM_CUSTOMER',
RELATION_ADD_OR_UPDATE = 'RELATION_ADD_OR_UPDATE',
RELATION_DELETED = 'RELATION_DELETED',
RPC_CALL = 'RPC_CALL', // deprecated - to be removed in 4.x
RPC_CALL_REQUEST = 'RPC_CALL_REQUEST',
RPC_CALL_RESPONSE = 'RPC_CALL_RESPONSE',
RPC_CALL = 'RPC_CALL',
ALARM_ACK = 'ALARM_ACK',
ALARM_CLEAR = 'ALARM_CLEAR',
ASSIGNED_TO_EDGE = 'ASSIGNED_TO_EDGE',
@ -130,8 +128,6 @@ export const edgeEventActionTypeTranslations = new Map<EdgeEventActionType, stri
[EdgeEventActionType.RELATION_ADD_OR_UPDATE, 'edge-event.action-type-relation-add-or-update'],
[EdgeEventActionType.RELATION_DELETED, 'edge-event.action-type-relation-deleted'],
[EdgeEventActionType.RPC_CALL, 'edge-event.action-type-rpc-call'],
[EdgeEventActionType.RPC_CALL_REQUEST, 'edge-event.action-type-rpc-call-request'],
[EdgeEventActionType.RPC_CALL_RESPONSE, 'edge-event.action-type-rpc-call-response'],
[EdgeEventActionType.ALARM_ACK, 'edge-event.action-type-alarm-ack'],
[EdgeEventActionType.ALARM_CLEAR, 'edge-event.action-type-alarm-clear'],
[EdgeEventActionType.ASSIGNED_TO_EDGE, 'edge-event.action-type-assigned-to-edge'],
@ -146,9 +142,7 @@ export const bodyContentEdgeEventActionTypes: EdgeEventActionType[] = [
EdgeEventActionType.ATTRIBUTES_UPDATED,
EdgeEventActionType.ATTRIBUTES_DELETED,
EdgeEventActionType.TIMESERIES_UPDATED,
EdgeEventActionType.RPC_CALL,
EdgeEventActionType.RPC_CALL_REQUEST,
EdgeEventActionType.RPC_CALL_RESPONSE,
EdgeEventActionType.RPC_CALL
];
export const edgeEventStatusColor = new Map<EdgeEventStatus, string>(

2
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -1791,8 +1791,6 @@
"action-type-relation-add-or-update": "Relation Add or Update",
"action-type-relation-deleted": "Relation Deleted",
"action-type-rpc-call": "RPC Call",
"action-type-rpc-call-request": "RPC Call Request",
"action-type-rpc-call-response": "RPC Call Response",
"action-type-alarm-ack": "Alarm Ack",
"action-type-alarm-clear": "Alarm Clear",
"action-type-assigned-to-edge": "Assigned to Edge",

Loading…
Cancel
Save