From 2757172ef333df755cf8ce199bfa2c3e7438c086 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 29 Jun 2026 12:40:35 +0300 Subject: [PATCH] fix(rpc): preserve original createdTime in post-persist RE notifications On the update paths buildRpc rebuilt the Rpc with createdTime=now, which the post-persist rule-engine notification then serialized, so RPC_DELIVERED/etc. carried the update moment instead of the row's real creation time (the UPDATE never touches created_time). Capture createdTime once at create and thread it through ToDeviceRpcRequestMetadata so update-path notifications report the original value, restoring pre-async behavior. Also add a striped-executor ordering test (same rpcId -> RPC_QUEUED before RPC_DELIVERED) and document that sql.rpc.callback_threads is independent of batch_threads. --- .../device/DeviceActorMessageProcessor.java | 31 +++++++------- .../device/ToDeviceRpcRequestMetadata.java | 3 ++ .../src/main/resources/thingsboard.yml | 2 +- .../server/service/rpc/TbRpcServiceTest.java | 40 +++++++++++++++++-- 4 files changed, 57 insertions(+), 19 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 6d0f07e507..21defb9f6a 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -193,17 +193,18 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso log.debug("[{}][{}] Received RPC request to process ...", deviceId, rpcId); ToDeviceRpcRequestMsg rpcRequest = createToDeviceRpcRequestMsg(request); - long timeout = request.getExpirationTime() - System.currentTimeMillis(); + long createdTime = System.currentTimeMillis(); + long timeout = request.getExpirationTime() - createdTime; boolean persisted = request.isPersisted(); if (timeout <= 0) { log.debug("[{}][{}] Ignoring message due to exp time reached, {}", deviceId, rpcId, request.getExpirationTime()); if (persisted) { - createRpc(request, RpcStatus.EXPIRED); + createRpc(request, RpcStatus.EXPIRED, createdTime); } return; } else if (persisted) { - createRpc(request, RpcStatus.QUEUED); + createRpc(request, RpcStatus.QUEUED, createdTime); } boolean sent = false; @@ -243,7 +244,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso log.debug("[{}] RPC command response sent [{}][{}]!", deviceId, rpcId, requestId); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null)); } else { - registerPendingRpcRequest(context, msg, sent, rpcRequest, timeout); + registerPendingRpcRequest(context, msg, sent, rpcRequest, timeout, createdTime); } String rpcSent = sent ? "sent!" : "NOT sent!"; log.debug("[{}][{}][{}] RPC request is {}", deviceId, rpcId, requestId, rpcSent); @@ -257,13 +258,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso }; } - private void createRpc(ToDeviceRpcRequest request, RpcStatus status) { - systemContext.getTbRpcService().create(tenantId, buildRpc(request, status, null)); + private void createRpc(ToDeviceRpcRequest request, RpcStatus status, long createdTime) { + systemContext.getTbRpcService().create(tenantId, buildRpc(request, status, null, createdTime)); } - private Rpc buildRpc(ToDeviceRpcRequest request, RpcStatus status, JsonNode response) { + private Rpc buildRpc(ToDeviceRpcRequest request, RpcStatus status, JsonNode response, long createdTime) { Rpc rpc = new Rpc(new RpcId(request.getId())); - rpc.setCreatedTime(System.currentTimeMillis()); + rpc.setCreatedTime(createdTime); rpc.setTenantId(tenantId); rpc.setDeviceId(deviceId); rpc.setExpirationTime(request.getExpirationTime()); @@ -339,11 +340,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } - private void registerPendingRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { + private void registerPendingRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout, long createdTime) { int requestId = rpcRequest.getRequestId(); UUID rpcId = new UUID(rpcRequest.getRequestIdMSB(), rpcRequest.getRequestIdLSB()); log.debug("[{}][{}][{}] Registering pending RPC request...", deviceId, rpcId, requestId); - toDeviceRpcPendingMap.put(requestId, new ToDeviceRpcRequestMetadata(msg, sent)); + toDeviceRpcPendingMap.put(requestId, new ToDeviceRpcRequestMetadata(msg, sent, createdTime)); DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg(requestId, timeout); scheduleMsgWithDelay(context, timeoutMsg, timeoutMsg.getTimeout()); } @@ -356,7 +357,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso UUID rpcId = toDeviceRpcRequest.getId(); log.debug("[{}][{}][{}] RPC request timeout detected!", deviceId, rpcId, requestId); if (toDeviceRpcRequest.isPersisted()) { - systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, RpcStatus.EXPIRED, null)); + systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, RpcStatus.EXPIRED, null, requestMd.getCreatedTime())); } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); @@ -667,7 +668,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } catch (IllegalArgumentException e) { response = JacksonUtil.newObjectNode().put("error", payload); } - systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRequestMsg, status, response)); + systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRequestMsg, status, response, requestMd.getCreatedTime())); } } finally { if (rpcSubmitStrategy.equals(RpcSubmitStrategy.SEQUENTIAL_ON_RESPONSE_FROM_DEVICE)) { @@ -731,7 +732,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } if (persisted) { - systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, status, response)); + systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, status, response, md.getCreatedTime())); } if (rpcSubmitStrategy.equals(RpcSubmitStrategy.SEQUENTIAL_ON_RESPONSE_FROM_DEVICE) && status.equals(RpcStatus.DELIVERED) && !oneWayRpc) { @@ -824,7 +825,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso var toDeviceRpcRequest = md.getMsg().getMsg(); if (toDeviceRpcRequest.isPersisted()) { var responseAwaitTimeout = JacksonUtil.newObjectNode().put("error", "There was a timeout awaiting for RPC response from device."); - systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, RpcStatus.FAILED, responseAwaitTimeout)); + systemContext.getTbRpcService().update(tenantId, buildRpc(toDeviceRpcRequest, RpcStatus.FAILED, responseAwaitTimeout, md.getCreatedTime())); } }, systemContext.getRpcResponseTimeout(), TimeUnit.MILLISECONDS); } @@ -1037,7 +1038,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso rpc.setStatus(RpcStatus.EXPIRED); systemContext.getTbRpcService().update(tenantId, rpc); } else { - registerPendingRpcRequest(ctx, new ToDeviceRpcRequestActorMsg(systemContext.getServiceId(), msg), false, createToDeviceRpcRequestMsg(msg), timeout); + registerPendingRpcRequest(ctx, new ToDeviceRpcRequestActorMsg(systemContext.getServiceId(), msg), false, createToDeviceRpcRequestMsg(msg), timeout, rpc.getCreatedTime()); } }); if (pageData.hasNext()) { diff --git a/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java b/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java index 782fdd92ad..c502ce8baa 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java @@ -25,6 +25,9 @@ import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg; public class ToDeviceRpcRequestMetadata { private final ToDeviceRpcRequestActorMsg msg; private final boolean sent; + // Creation time of the persisted RPC row, captured once at create so post-persist rule-engine + // notifications on the update paths carry the original createdTime instead of the update moment. + private final long createdTime; private int retries; private boolean delivered; } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index af1b7d328d..5cda04a33d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -429,7 +429,7 @@ sql: batch_max_delay: "${SQL_RPC_BATCH_MAX_DELAY_MS:50}" # Max timeout for RPC entries queue polling, in milliseconds stats_print_interval_ms: "${SQL_RPC_BATCH_STATS_PRINT_MS:10000}" # Interval in milliseconds for printing RPC persistence statistic batch_threads: "${SQL_RPC_BATCH_THREADS:3}" # number of queue partitions for batched RPC persistence; writes for a given RPC always map to the same partition (by rpcId hash) so they persist in submission order. A prime value (e.g. 3 or 5) keeps the hash distribution even - callback_threads: "${SQL_RPC_CALLBACK_THREADS:3}" # striped threads for post-persist rule-engine notifications; each rpcId maps to a single stripe, preserving per-command notification order. A prime value (e.g. 3 or 5) keeps the hash distribution even + callback_threads: "${SQL_RPC_CALLBACK_THREADS:3}" # striped threads for post-persist rule-engine notifications; each rpcId maps to a single stripe, preserving per-command notification order (e.g. RPC_QUEUED before RPC_DELIVERED). Independent of batch_threads above: that pool partitions DB writes, this one partitions notifications, and each stripes by rpcId on its own - so the two can be sized independently without affecting ordering. A prime value (e.g. 3 or 5) keeps the hash distribution even events: batch_size: "${SQL_EVENTS_BATCH_SIZE:10000}" # Batch size for persisting latest telemetry updates batch_max_delay: "${SQL_EVENTS_BATCH_MAX_DELAY_MS:100}" # Max timeout for latest telemetry entries queue polling. The value set in milliseconds 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 cf11dfed48..7190e1cb69 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 @@ -18,18 +18,22 @@ package org.thingsboard.server.service.rpc; import com.google.common.util.concurrent.Futures; import org.junit.Before; import org.junit.Test; +import org.mockito.ArgumentCaptor; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.RpcId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.rpc.RpcService; +import java.util.List; import java.util.UUID; +import static org.junit.Assert.assertEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; @@ -88,11 +92,41 @@ public class TbRpcServiceTest { .pushMsgToRuleEngine(any(TenantId.class), any(DeviceId.class), any(TbMsg.class), isNull()); } + @Test + public void sameRpcIdNotificationsRunInSubmissionOrderOnStripedExecutor() { + // Use several stripes so executorFor actually partitions. Both writes share an rpcId, so they + // must map to the same stripe and the rule-engine notifications must arrive in submission order + // (RPC_QUEUED before RPC_DELIVERED) - this is the ordering guarantee executorFor exists for. + tbRpcService = new TbRpcService(rpcService, clusterService, 3); + + RpcId rpcId = new RpcId(UUID.randomUUID()); + DeviceId deviceId = new DeviceId(UUID.randomUUID()); + Rpc queued = newRpc(rpcId, deviceId, RpcStatus.QUEUED); + Rpc delivered = newRpc(rpcId, deviceId, RpcStatus.DELIVERED); + when(rpcService.createAsync(queued)).thenReturn(Futures.immediateFuture(true)); + when(rpcService.updateAsync(delivered)).thenReturn(Futures.immediateFuture(true)); + + tbRpcService.create(queued.getTenantId(), queued); + tbRpcService.update(delivered.getTenantId(), delivered); + + ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(clusterService, timeout(5000).times(2)) + .pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), msgCaptor.capture(), isNull()); + + List msgs = msgCaptor.getAllValues(); + assertEquals(TbMsgType.RPC_QUEUED, msgs.get(0).getInternalType()); + assertEquals(TbMsgType.RPC_DELIVERED, msgs.get(1).getInternalType()); + } + private Rpc newRpc() { - Rpc rpc = new Rpc(new RpcId(UUID.randomUUID())); + return newRpc(new RpcId(UUID.randomUUID()), new DeviceId(UUID.randomUUID()), RpcStatus.QUEUED); + } + + private Rpc newRpc(RpcId rpcId, DeviceId deviceId, RpcStatus status) { + Rpc rpc = new Rpc(rpcId); rpc.setTenantId(TenantId.SYS_TENANT_ID); - rpc.setDeviceId(new DeviceId(UUID.randomUUID())); - rpc.setStatus(RpcStatus.QUEUED); + rpc.setDeviceId(deviceId); + rpc.setStatus(status); rpc.setRequest(JacksonUtil.toJsonNode("{}")); return rpc; }