Browse Source

Merge pull request #15662 from dashevchenko/fromDeviceRPCResponseProtoHandlingFix

Fixed RPC call request rule node returning null body
pull/15744/head
Viacheslav Klimov 4 months ago
committed by GitHub
parent
commit
d011b611f6
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  3. 32
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerServiceTest.java
  4. 78
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerServiceTest.java
  5. 12
      common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java
  6. 7
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  7. 2
      common/proto/src/main/proto/queue.proto
  8. 11
      common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -465,10 +465,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
return firmwareStateService.process(msg.getValue());
}
private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = RpcError.fromProtoErrorCode(proto.getError());
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
, proto.getResponse(), error);
, proto.hasResponse() ? proto.getResponse() : null, error);
tbCoreDeviceRpcService.processRpcResponseFromRuleEngine(response);
callback.onSuccess();
}

4
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -180,9 +180,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractPartitionBasedCo
callback.onSuccess();
} else if (nfMsg.hasFromDeviceRpcResponse()) {
TransportProtos.FromDeviceRPCResponseProto proto = nfMsg.getFromDeviceRpcResponse();
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
RpcError error = RpcError.fromProtoErrorCode(proto.getError());
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
, proto.getResponse(), error);
, proto.hasResponse() ? proto.getResponse() : null, error);
tbDeviceRpcService.processRpcResponseFromDevice(response);
callback.onSuccess();
} else if (nfMsg.getQueueUpdateMsgsCount() > 0) {

32
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerServiceTest.java

@ -27,8 +27,11 @@ import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.rpc.RpcError;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
import org.thingsboard.server.service.ruleengine.RuleEngineCallService;
import org.thingsboard.server.service.state.DeviceStateService;
@ -51,6 +54,8 @@ public class DefaultTbCoreConsumerServiceTest {
private TbCoreConsumerStats statsMock;
@Mock
private RuleEngineCallService ruleEngineCallServiceMock;
@Mock
private TbCoreDeviceRpcService tbCoreDeviceRpcServiceMock;
@Mock
private TbCallback tbCallbackMock;
@ -638,4 +643,31 @@ public class DefaultTbCoreConsumerServiceTest {
then(ruleEngineCallServiceMock).should().onQueueMsg(restApiCallResponseMsgProto, tbCallbackMock);
}
@Test
public void givenNotFoundErrorAndNoResponse_whenForwardToCoreRpcService_thenNotFoundAndNullResponseAreRecovered() {
// GIVEN
ReflectionTestUtils.setField(defaultTbCoreConsumerServiceMock, "tbCoreDeviceRpcService", tbCoreDeviceRpcServiceMock);
var requestId = UUID.randomUUID();
// error = NOT_FOUND.ordinal() (0) and response left unset: the previously broken combination
// ('error > 0' dropped NOT_FOUND, proto3 default collapsed a null response to "").
var proto = TransportProtos.FromDeviceRPCResponseProto.newBuilder()
.setRequestIdMSB(requestId.getMostSignificantBits())
.setRequestIdLSB(requestId.getLeastSignificantBits())
.setError(RpcError.NOT_FOUND.ordinal())
.build();
doCallRealMethod().when(defaultTbCoreConsumerServiceMock).forwardToCoreRpcService(proto, tbCallbackMock);
// WHEN
defaultTbCoreConsumerServiceMock.forwardToCoreRpcService(proto, tbCallbackMock);
// THEN
var responseCaptor = ArgumentCaptor.forClass(FromDeviceRpcResponse.class);
then(tbCoreDeviceRpcServiceMock).should().processRpcResponseFromRuleEngine(responseCaptor.capture());
var response = responseCaptor.getValue();
assertThat(response.getId()).isEqualTo(requestId);
assertThat(response.getError()).contains(RpcError.NOT_FOUND);
assertThat(response.getResponse()).isEmpty();
then(tbCallbackMock).should().onSuccess();
}
}

78
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerServiceTest.java

@ -0,0 +1,78 @@
/**
* Copyright © 2016-2026 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.queue;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.rpc.RpcError;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.doCallRealMethod;
@ExtendWith(MockitoExtension.class)
public class DefaultTbRuleEngineConsumerServiceTest {
@Mock
private TbRuleEngineDeviceRpcService tbDeviceRpcServiceMock;
@Mock
private TbCallback tbCallbackMock;
@Mock
private DefaultTbRuleEngineConsumerService defaultTbRuleEngineConsumerServiceMock;
@Test
public void givenNotFoundErrorAndNoResponse_whenHandleFromDeviceRpcResponse_thenNotFoundAndNullResponseAreRecovered() {
// GIVEN
ReflectionTestUtils.setField(defaultTbRuleEngineConsumerServiceMock, "tbDeviceRpcService", tbDeviceRpcServiceMock);
var requestId = UUID.randomUUID();
// error = NOT_FOUND.ordinal() (0) and response left unset: the previously broken combination
// ('error > 0' dropped NOT_FOUND, proto3 default collapsed a null response to "").
var proto = TransportProtos.FromDeviceRPCResponseProto.newBuilder()
.setRequestIdMSB(requestId.getMostSignificantBits())
.setRequestIdLSB(requestId.getLeastSignificantBits())
.setError(RpcError.NOT_FOUND.ordinal())
.build();
var nfMsg = ToRuleEngineNotificationMsg.newBuilder().setFromDeviceRpcResponse(proto).build();
var queueMsg = new TbProtoQueueMsg<>(requestId, nfMsg);
doCallRealMethod().when(defaultTbRuleEngineConsumerServiceMock).handleNotification(requestId, queueMsg, tbCallbackMock);
// WHEN
defaultTbRuleEngineConsumerServiceMock.handleNotification(requestId, queueMsg, tbCallbackMock);
// THEN
var responseCaptor = ArgumentCaptor.forClass(FromDeviceRpcResponse.class);
then(tbDeviceRpcServiceMock).should().processRpcResponseFromDevice(responseCaptor.capture());
var response = responseCaptor.getValue();
assertThat(response.getId()).isEqualTo(requestId);
assertThat(response.getError()).contains(RpcError.NOT_FOUND);
assertThat(response.getResponse()).isEmpty();
then(tbCallbackMock).should().onSuccess();
}
}

12
common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java

@ -20,4 +20,16 @@ package org.thingsboard.server.common.data.rpc;
*/
public enum RpcError {
NOT_FOUND, FORBIDDEN, NO_ACTIVE_CONNECTION, TIMEOUT, INTERNAL;
private static final RpcError[] VALUES = values();
/**
* Resolves an {@link RpcError} from the proto {@code error} ordinal.
* Returns {@code null} both for the "no error" sentinel (negative value) and for unknown ordinals
* that a newer node in a mixed-version cluster might emit, so callers never hit an
* {@link ArrayIndexOutOfBoundsException}.
*/
public static RpcError fromProtoErrorCode(int errorCode) {
return errorCode >= 0 && errorCode < VALUES.length ? VALUES[errorCode] : null;
}
}

7
common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java

@ -583,10 +583,11 @@ public class ProtoUtils {
}
private static ToDeviceActorNotificationMsg fromProto(TransportProtos.FromDeviceRpcResponseActorMsgProto proto) {
TransportProtos.FromDeviceRPCResponseProto rpcResponse = proto.getRpcResponse();
FromDeviceRpcResponse fromDeviceRpcResponse = new FromDeviceRpcResponse(
new UUID(proto.getRpcResponse().getRequestIdMSB(), proto.getRpcResponse().getRequestIdLSB()),
proto.getRpcResponse().getResponse(),
proto.getRpcResponse().getError() >= 0 ? RpcError.values()[proto.getRpcResponse().getError()] : null);
new UUID(rpcResponse.getRequestIdMSB(), rpcResponse.getRequestIdLSB()),
rpcResponse.hasResponse() ? rpcResponse.getResponse() : null,
RpcError.fromProtoErrorCode(rpcResponse.getError()));
return new FromDeviceRpcResponseActorMsg(
proto.getRequestId(),
TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())),

2
common/proto/src/main/proto/queue.proto

@ -1238,7 +1238,7 @@ message LocalSubscriptionServiceMsgProto {
message FromDeviceRPCResponseProto {
int64 requestIdMSB = 1;
int64 requestIdLSB = 2;
string response = 3;
optional string response = 3;
int32 error = 4;
}

11
common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java

@ -226,6 +226,17 @@ class ProtoUtilsTest {
assertThat(ProtoUtils.fromProto(serializedMsg)).as("deserialized").isEqualTo(msg);
}
@Test
void protoFromDeviceRpcResponseOnewaySerialization() {
// Oneway RPC success: response and error are both null. Relies on the proto
// 'optional string response' presence bit so the receiver round-trips null
// rather than seeing the proto3 default "".
FromDeviceRpcResponseActorMsg msg = new FromDeviceRpcResponseActorMsg(23, tenantId, deviceId, new FromDeviceRpcResponse(id, null, null));
TransportProtos.ToDeviceActorNotificationMsgProto serializedMsg = ProtoUtils.toProto(msg);
Assertions.assertNotNull(serializedMsg);
assertThat(ProtoUtils.fromProto(serializedMsg)).as("deserialized").isEqualTo(msg);
}
@Test
void protoRemoveRpcActorSerialization() {
RemoveRpcActorMsg msg = new RemoveRpcActorMsg(tenantId, deviceId, id);

Loading…
Cancel
Save