Browse Source

fix(rpc): persist QUEUED RPC synchronously before sending to device

pull/15853/head
dshvaika 3 months ago
parent
commit
7e35f7b184
  1. 12
      application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java
  2. 9
      application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java

12
application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java

@ -73,7 +73,8 @@ public class TbRpcService {
} }
public void create(TenantId tenantId, Rpc rpc) { public void create(TenantId tenantId, Rpc rpc) {
persist(tenantId, rpc, rpcService.createAsync(rpc)); rpcService.save(rpc);
executorFor(rpc.getUuidId()).execute(() -> notifyRuleEngine(tenantId, rpc));
} }
public void update(TenantId tenantId, Rpc rpc) { public void update(TenantId tenantId, Rpc rpc) {
@ -99,6 +100,15 @@ public class TbRpcService {
return callbackExecutors[HashPartitioner.resolvePartition(rpcId.hashCode(), callbackExecutors.length)]; return callbackExecutors[HashPartitioner.resolvePartition(rpcId.hashCode(), callbackExecutors.length)];
} }
private void notifyRuleEngine(TenantId tenantId, Rpc rpc) {
try {
pushRpcMsgToRuleEngine(tenantId, rpc);
} catch (Throwable t) {
log.error("[{}][{}][{}] Failed to push RPC with status [{}] to rule engine",
tenantId, rpc.getDeviceId(), rpc.getId(), rpc.getStatus(), t);
}
}
private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) { private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) {
TbMsg msg = TbMsg.newMsg() TbMsg msg = TbMsg.newMsg()
.type(TbMsgType.valueOf("RPC_" + rpc.getStatus().name())) .type(TbMsgType.valueOf("RPC_" + rpc.getStatus().name()))

9
application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java

@ -43,6 +43,7 @@ import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.after; import static org.mockito.Mockito.after;
import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when; import static org.mockito.Mockito.when;
@ -71,13 +72,15 @@ public class TbRpcServiceTest {
} }
@Test @Test
public void createPersistsViaCreateAsyncThenPushesToRuleEngine() { public void createPersistsSynchronouslyThenPushesToRuleEngine() {
Rpc rpc = newRpc(); Rpc rpc = newRpc();
when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(true)); when(rpcService.save(rpc)).thenReturn(rpc);
tbRpcService.create(rpc.getTenantId(), rpc); tbRpcService.create(rpc.getTenantId(), rpc);
verify(rpcService).createAsync(rpc); // create must persist synchronously via save(...), never via the async batch path
verify(rpcService).save(rpc);
verify(rpcService, never()).createAsync(any());
verify(clusterService, timeout(5000)) verify(clusterService, timeout(5000))
.pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull()); .pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull());
} }

Loading…
Cancel
Save