Browse Source

fix(rpc): preserve response and prevent row resurrection in batched upsert

Split the batched RPC persistence into a create path (INSERT ... ON CONFLICT)
and an update path (UPDATE ... WHERE id = ?). A status update for a row deleted
in the meantime (TTL cleanup / manual delete) is no longer resurrected,
restoring the old findById-null skip. Both paths COALESCE the response so a
status update that carries none no longer clobbers a previously stored one.

The create-vs-update intent is carried by a dedicated RpcQueueEntry record
(keeping RpcEntity pure data) and exposed via explicit createAsync/updateAsync
DAO methods instead of a boolean flag; insert/update batch scaffolding is
collapsed into a shared helper.

Adds DAO regression tests (re-queue keeps existing row, null-response
preservation, no resurrection of a deleted row) and TbRpcService create/update
wiring tests.
pull/15853/head
dshvaika 3 months ago
parent
commit
f8d59a25e4
  1. 2
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 9
      application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java
  3. 33
      application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java
  4. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java
  5. 12
      dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java
  6. 4
      dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
  7. 17
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
  8. 60
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java
  9. 34
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcQueueEntry.java
  10. 76
      dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java

2
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -258,7 +258,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
private void createRpc(ToDeviceRpcRequest request, RpcStatus status) { private void createRpc(ToDeviceRpcRequest request, RpcStatus status) {
systemContext.getTbRpcService().save(tenantId, buildRpc(request, status, null)); systemContext.getTbRpcService().create(tenantId, buildRpc(request, status, null));
} }
private Rpc buildRpc(ToDeviceRpcRequest request, RpcStatus status, JsonNode response) { private Rpc buildRpc(ToDeviceRpcRequest request, RpcStatus status, JsonNode response) {

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

@ -75,8 +75,15 @@ public class TbRpcService {
} }
} }
public void create(TenantId tenantId, Rpc rpc) {
persist(tenantId, rpc, rpcService.createAsync(rpc));
}
public void save(TenantId tenantId, Rpc rpc) { public void save(TenantId tenantId, Rpc rpc) {
ListenableFuture<Void> future = rpcService.saveAsync(rpc); persist(tenantId, rpc, rpcService.updateAsync(rpc));
}
private void persist(TenantId tenantId, Rpc rpc, ListenableFuture<Void> future) {
DonAsynchron.withCallback(future, DonAsynchron.withCallback(future,
v -> pushRpcMsgToRuleEngine(tenantId, rpc), v -> pushRpcMsgToRuleEngine(tenantId, rpc),
t -> log.error("[{}][{}][{}] Failed to persist RPC with status [{}]", t -> log.error("[{}][{}][{}] Failed to persist RPC with status [{}]",

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

@ -53,18 +53,35 @@ public class TbRpcServiceTest {
} }
@Test @Test
public void savePushesToRuleEngineAfterFlush() { public void savePersistsViaUpdateAsyncThenPushesToRuleEngine() {
Rpc rpc = new Rpc(new RpcId(UUID.randomUUID())); Rpc rpc = newRpc();
rpc.setTenantId(TenantId.SYS_TENANT_ID); when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(null));
rpc.setDeviceId(new DeviceId(UUID.randomUUID()));
rpc.setStatus(RpcStatus.QUEUED);
rpc.setRequest(JacksonUtil.toJsonNode("{}"));
when(rpcService.saveAsync(rpc)).thenReturn(Futures.immediateFuture(null));
tbRpcService.save(rpc.getTenantId(), rpc); tbRpcService.save(rpc.getTenantId(), rpc);
verify(rpcService).updateAsync(rpc);
verify(clusterService, timeout(5000))
.pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull());
}
@Test
public void createPersistsViaCreateAsyncThenPushesToRuleEngine() {
Rpc rpc = newRpc();
when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(null));
tbRpcService.create(rpc.getTenantId(), rpc);
verify(rpcService).createAsync(rpc);
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());
} }
private Rpc newRpc() {
Rpc rpc = new Rpc(new RpcId(UUID.randomUUID()));
rpc.setTenantId(TenantId.SYS_TENANT_ID);
rpc.setDeviceId(new DeviceId(UUID.randomUUID()));
rpc.setStatus(RpcStatus.QUEUED);
rpc.setRequest(JacksonUtil.toJsonNode("{}"));
return rpc;
}
} }

4
common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java

@ -29,7 +29,9 @@ public interface RpcService extends EntityDaoService {
Rpc save(Rpc rpc); Rpc save(Rpc rpc);
ListenableFuture<Void> saveAsync(Rpc rpc); ListenableFuture<Void> createAsync(Rpc rpc);
ListenableFuture<Void> updateAsync(Rpc rpc);
void deleteRpc(TenantId tenantId, RpcId id); void deleteRpc(TenantId tenantId, RpcId id);

12
dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java

@ -55,9 +55,15 @@ public class BaseRpcService implements RpcService {
} }
@Override @Override
public ListenableFuture<Void> saveAsync(Rpc rpc) { public ListenableFuture<Void> createAsync(Rpc rpc) {
log.trace("Executing saveAsync, [{}]", rpc); log.trace("Executing createAsync, [{}]", rpc);
return rpcDao.saveAsync(rpc.getTenantId(), rpc); return rpcDao.createAsync(rpc.getTenantId(), rpc);
}
@Override
public ListenableFuture<Void> updateAsync(Rpc rpc) {
log.trace("Executing updateAsync, [{}]", rpc);
return rpcDao.updateAsync(rpc.getTenantId(), rpc);
} }
@Override @Override

4
dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java

@ -26,7 +26,9 @@ import org.thingsboard.server.dao.Dao;
public interface RpcDao extends Dao<Rpc> { public interface RpcDao extends Dao<Rpc> {
ListenableFuture<Void> saveAsync(TenantId tenantId, Rpc rpc); ListenableFuture<Void> createAsync(TenantId tenantId, Rpc rpc);
ListenableFuture<Void> updateAsync(TenantId tenantId, Rpc rpc);
PageData<Rpc> findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink); PageData<Rpc> findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink);

17
dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java

@ -71,7 +71,7 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
@Value("${sql.batch_sort:true}") @Value("${sql.batch_sort:true}")
private boolean batchSortEnabled; private boolean batchSortEnabled;
private TbSqlBlockingQueueWrapper<RpcEntity, Void> queue; private TbSqlBlockingQueueWrapper<RpcQueueEntry, Void> queue;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -84,10 +84,10 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
.batchSortEnabled(batchSortEnabled) .batchSortEnabled(batchSortEnabled)
.withResponse(false) .withResponse(false)
.build(); .build();
Function<RpcEntity, Integer> hashcodeFunction = entity -> entity.getUuid().hashCode(); Function<RpcQueueEntry, Integer> hashcodeFunction = entry -> entry.entity().getUuid().hashCode();
queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory);
queue.init(logExecutor, entities -> rpcInsertRepository.saveOrUpdate(entities), queue.init(logExecutor, entries -> rpcInsertRepository.saveOrUpdate(entries),
Comparator.comparing(RpcEntity::getUuid)); Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid()));
} }
@PreDestroy @PreDestroy
@ -98,8 +98,13 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
} }
@Override @Override
public ListenableFuture<Void> saveAsync(TenantId tenantId, Rpc rpc) { public ListenableFuture<Void> createAsync(TenantId tenantId, Rpc rpc) {
return queue.add(new RpcEntity(rpc)); return queue.add(RpcQueueEntry.forInsert(new RpcEntity(rpc)));
}
@Override
public ListenableFuture<Void> updateAsync(TenantId tenantId, Rpc rpc) {
return queue.add(RpcQueueEntry.forUpdate(new RpcEntity(rpc)));
} }
@Override @Override

60
dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java

@ -29,17 +29,33 @@ import java.util.List;
@Repository @Repository
public class RpcInsertRepository extends AbstractInsertRepository { public class RpcInsertRepository extends AbstractInsertRepository {
private static final String INSERT_OR_UPDATE = // Used only for the initial persistence of a new RPC (RpcStatus.QUEUED / already-EXPIRED on arrival).
// ON CONFLICT keeps the create idempotent (e.g. on actor re-processing); COALESCE never clobbers an
// existing response with NULL.
private static final String INSERT =
"INSERT INTO rpc (id, created_time, tenant_id, device_id, expiration_time, request, response, additional_info, status) " + "INSERT INTO rpc (id, created_time, tenant_id, device_id, expiration_time, request, response, additional_info, status) " +
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " +
"ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, response = EXCLUDED.response;"; "ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, response = COALESCE(EXCLUDED.response, rpc.response);";
public void saveOrUpdate(List<RpcEntity> entities) { // Used for every subsequent status change. WHERE id = ? means a row deleted in the meantime
// (TTL cleanup / manual delete) is NOT resurrected - it matches the old findById-null skip.
// COALESCE preserves a previously stored response when a status update carries none.
private static final String UPDATE =
"UPDATE rpc SET status = ?, response = COALESCE(?, response) WHERE id = ?;";
@FunctionalInterface
private interface ColumnBinder {
void bind(PreparedStatement ps, RpcEntity rpc) throws SQLException;
}
public void saveOrUpdate(List<RpcQueueEntry> entries) {
List<RpcEntity> inserts = entries.stream().filter(RpcQueueEntry::insert).map(RpcQueueEntry::entity).toList();
List<RpcEntity> updates = entries.stream().filter(entry -> !entry.insert()).map(RpcQueueEntry::entity).toList();
transactionTemplate.execute(status -> { transactionTemplate.execute(status -> {
jdbcTemplate.batchUpdate(INSERT_OR_UPDATE, new BatchPreparedStatementSetter() { // Inserts run first so a create and a status update for the same rpcId coalesced into one
@Override // batch still apply in create -> update order.
public void setValues(PreparedStatement ps, int i) throws SQLException { if (!inserts.isEmpty()) {
RpcEntity rpc = entities.get(i); batch(INSERT, inserts, (ps, rpc) -> {
ps.setObject(1, rpc.getUuid()); ps.setObject(1, rpc.getUuid());
ps.setLong(2, rpc.getCreatedTime()); ps.setLong(2, rpc.getCreatedTime());
ps.setObject(3, rpc.getTenantId()); ps.setObject(3, rpc.getTenantId());
@ -49,17 +65,33 @@ public class RpcInsertRepository extends AbstractInsertRepository {
ps.setString(7, toJsonStr(rpc.getResponse())); ps.setString(7, toJsonStr(rpc.getResponse()));
ps.setString(8, toJsonStr(rpc.getAdditionalInfo())); ps.setString(8, toJsonStr(rpc.getAdditionalInfo()));
ps.setString(9, rpc.getStatus().name()); ps.setString(9, rpc.getStatus().name());
} });
}
@Override if (!updates.isEmpty()) {
public int getBatchSize() { batch(UPDATE, updates, (ps, rpc) -> {
return entities.size(); ps.setString(1, rpc.getStatus().name());
} ps.setString(2, toJsonStr(rpc.getResponse()));
}); ps.setObject(3, rpc.getUuid());
});
}
return null; return null;
}); });
} }
private void batch(String sql, List<RpcEntity> entities, ColumnBinder binder) {
jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
binder.bind(ps, entities.get(i));
}
@Override
public int getBatchSize() {
return entities.size();
}
});
}
private String toJsonStr(JsonNode node) { private String toJsonStr(JsonNode node) {
return node == null ? null : replaceNullChars(JacksonUtil.toString(node)); return node == null ? null : replaceNullChars(JacksonUtil.toString(node));
} }

34
dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcQueueEntry.java

@ -0,0 +1,34 @@
/**
* 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.dao.sql.rpc;
import org.thingsboard.server.dao.model.sql.RpcEntity;
/**
* Pairs an {@link RpcEntity} with the write intent for the batched persistence queue, so the entity
* itself stays pure data. {@code insert} selects INSERT-on-conflict (initial create) vs UPDATE-by-id
* (status change that must never resurrect a deleted row).
*/
record RpcQueueEntry(RpcEntity entity, boolean insert) {
static RpcQueueEntry forInsert(RpcEntity entity) {
return new RpcQueueEntry(entity, true);
}
static RpcQueueEntry forUpdate(RpcEntity entity) {
return new RpcQueueEntry(entity, false);
}
}

76
dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.sql.rpc; package org.thingsboard.server.dao.sql.rpc;
import com.fasterxml.jackson.databind.JsonNode;
import org.junit.Test; import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
@ -70,7 +71,7 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}")); rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}"));
rpc.setStatus(RpcStatus.QUEUED); rpc.setStatus(RpcStatus.QUEUED);
rpcDao.saveAsync(rpc.getTenantId(), rpc).get(5, TimeUnit.SECONDS); rpcDao.createAsync(rpc.getTenantId(), rpc).get(5, TimeUnit.SECONDS);
Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id); Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(stored).isNotNull(); assertThat(stored).isNotNull();
@ -86,7 +87,7 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
update.setStatus(RpcStatus.DELIVERED); update.setStatus(RpcStatus.DELIVERED);
update.setResponse(JacksonUtil.toJsonNode("{\"ok\":true}")); update.setResponse(JacksonUtil.toJsonNode("{\"ok\":true}"));
rpcDao.saveAsync(update.getTenantId(), update).get(5, TimeUnit.SECONDS); rpcDao.updateAsync(update.getTenantId(), update).get(5, TimeUnit.SECONDS);
Rpc afterUpdate = rpcDao.findById(TenantId.SYS_TENANT_ID, id); Rpc afterUpdate = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(afterUpdate.getStatus()).isEqualTo(RpcStatus.DELIVERED); assertThat(afterUpdate.getStatus()).isEqualTo(RpcStatus.DELIVERED);
@ -120,8 +121,8 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
// Enqueue both writes for the same rpcId back-to-back so they coalesce into one flush batch. // Enqueue both writes for the same rpcId back-to-back so they coalesce into one flush batch.
// Same rpcId -> same partition; the queue's stable sort must keep submission order // Same rpcId -> same partition; the queue's stable sort must keep submission order
// (QUEUED before DELIVERED), so the final persisted row must be DELIVERED, never QUEUED. // (QUEUED before DELIVERED), so the final persisted row must be DELIVERED, never QUEUED.
var queuedFuture = rpcDao.saveAsync(TenantId.SYS_TENANT_ID, queued); var queuedFuture = rpcDao.createAsync(TenantId.SYS_TENANT_ID, queued);
var deliveredFuture = rpcDao.saveAsync(TenantId.SYS_TENANT_ID, delivered); var deliveredFuture = rpcDao.updateAsync(TenantId.SYS_TENANT_ID, delivered);
queuedFuture.get(5, TimeUnit.SECONDS); queuedFuture.get(5, TimeUnit.SECONDS);
deliveredFuture.get(5, TimeUnit.SECONDS); deliveredFuture.get(5, TimeUnit.SECONDS);
@ -131,4 +132,71 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
assertThat(stored.getResponse()).isEqualTo(JacksonUtil.toJsonNode("{\"ok\":true}")); assertThat(stored.getResponse()).isEqualTo(JacksonUtil.toJsonNode("{\"ok\":true}"));
} }
@Test
public void saveAsyncRequeueUpdatesExistingRowToQueued() throws Exception {
UUID id = UUID.randomUUID();
DeviceId deviceId = new DeviceId(UUID.randomUUID());
// Initial create.
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// Delivery timeout with closeTransportSessionOnRpcDeliveryTimeout=true re-queues the RPC: the
// device actor persists status=QUEUED again as a status update so init() can re-pick it up.
// The update must land on the existing row, never be dropped.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(stored).isNotNull();
assertThat(stored.getStatus()).isEqualTo(RpcStatus.QUEUED);
}
@Test
public void saveAsyncNullResponseUpdateKeepsStoredResponse() throws Exception {
UUID id = UUID.randomUUID();
DeviceId deviceId = new DeviceId(UUID.randomUUID());
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// A successful response is stored.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS);
// A later status update carries no response - it must NOT clobber the stored one.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.EXPIRED, null)).get(5, TimeUnit.SECONDS);
Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(stored.getStatus()).isEqualTo(RpcStatus.EXPIRED);
assertThat(stored.getResponse()).isEqualTo(JacksonUtil.toJsonNode("{\"ok\":true}"));
}
@Test
public void saveAsyncUpdateForDeletedRpcDoesNotResurrect() throws Exception {
UUID id = UUID.randomUUID();
DeviceId deviceId = new DeviceId(UUID.randomUUID());
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// RPC is removed (TTL cleanup / manual delete) while a response is still in flight.
rpcDao.removeById(TenantId.SYS_TENANT_ID, id);
assertThat(rpcDao.findById(TenantId.SYS_TENANT_ID, id)).isNull();
// A late status update must not re-create the deleted row.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS);
assertThat(rpcDao.findById(TenantId.SYS_TENANT_ID, id)).isNull();
}
private Rpc rpc(UUID id, DeviceId deviceId, RpcStatus status, JsonNode response) {
Rpc rpc = new Rpc(new RpcId(id));
rpc.setCreatedTime(System.currentTimeMillis());
rpc.setTenantId(TenantId.SYS_TENANT_ID);
rpc.setDeviceId(deviceId);
rpc.setExpirationTime(System.currentTimeMillis() + 60_000);
rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}"));
rpc.setStatus(status);
rpc.setResponse(response);
return rpc;
}
} }

Loading…
Cancel
Save