From 7e35f7b184ad20ddb572d1b9b4fff03cf1f01835 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 30 Jun 2026 16:25:44 +0300 Subject: [PATCH] fix(rpc): persist QUEUED RPC synchronously before sending to device --- .../thingsboard/server/service/rpc/TbRpcService.java | 12 +++++++++++- .../server/service/rpc/TbRpcServiceTest.java | 9 ++++++--- 2 files changed, 17 insertions(+), 4 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java index e6d7f06c48..0e11658451 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java +++ b/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) { - persist(tenantId, rpc, rpcService.createAsync(rpc)); + rpcService.save(rpc); + executorFor(rpc.getUuidId()).execute(() -> notifyRuleEngine(tenantId, rpc)); } public void update(TenantId tenantId, Rpc rpc) { @@ -99,6 +100,15 @@ public class TbRpcService { 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) { TbMsg msg = TbMsg.newMsg() .type(TbMsgType.valueOf("RPC_" + rpc.getStatus().name())) diff --git a/application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java b/application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java index 300a9ed8d3..585da09b43 100644 --- a/application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java +++ b/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.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -71,13 +72,15 @@ public class TbRpcServiceTest { } @Test - public void createPersistsViaCreateAsyncThenPushesToRuleEngine() { + public void createPersistsSynchronouslyThenPushesToRuleEngine() { Rpc rpc = newRpc(); - when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(true)); + when(rpcService.save(rpc)).thenReturn(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)) .pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull()); }