Browse Source

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.
pull/15853/head
dshvaika 3 months ago
parent
commit
2757172ef3
  1. 31
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 3
      application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java
  3. 2
      application/src/main/resources/thingsboard.yml
  4. 40
      application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java

31
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()) {

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

2
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

40
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<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(clusterService, timeout(5000).times(2))
.pushMsgToRuleEngine(eq(TenantId.SYS_TENANT_ID), eq(deviceId), msgCaptor.capture(), isNull());
List<TbMsg> 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;
}

Loading…
Cancel
Save