Browse Source

fixed pr comments

pull/15662/head
dashevchenko 4 months ago
parent
commit
93df85a78c
  1. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  2. 2
      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. 2
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java

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

@ -465,8 +465,8 @@ 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.hasResponse() ? proto.getResponse() : null, error);
tbCoreDeviceRpcService.processRpcResponseFromRuleEngine(response);

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

@ -180,7 +180,7 @@ 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.hasResponse() ? proto.getResponse() : null, error);
tbDeviceRpcService.processRpcResponseFromDevice(response);

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;
}
}

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

@ -587,7 +587,7 @@ public class ProtoUtils {
FromDeviceRpcResponse fromDeviceRpcResponse = new FromDeviceRpcResponse(
new UUID(rpcResponse.getRequestIdMSB(), rpcResponse.getRequestIdLSB()),
rpcResponse.hasResponse() ? rpcResponse.getResponse() : null,
rpcResponse.getError() >= 0 ? RpcError.values()[rpcResponse.getError()] : null);
RpcError.fromProtoErrorCode(rpcResponse.getError()));
return new FromDeviceRpcResponseActorMsg(
proto.getRequestId(),
TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())),

Loading…
Cancel
Save