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 ef48e6791b..6f257405af 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 @@ -258,7 +258,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } 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) { 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 2f5c5ea6a8..486feab2e2 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 @@ -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) { - ListenableFuture future = rpcService.saveAsync(rpc); + persist(tenantId, rpc, rpcService.updateAsync(rpc)); + } + + private void persist(TenantId tenantId, Rpc rpc, ListenableFuture future) { DonAsynchron.withCallback(future, v -> pushRpcMsgToRuleEngine(tenantId, rpc), t -> log.error("[{}][{}][{}] Failed to persist RPC with status [{}]", 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 d0736ada26..014ab256af 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 @@ -53,18 +53,35 @@ public class TbRpcServiceTest { } @Test - public void savePushesToRuleEngineAfterFlush() { - 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("{}")); - - when(rpcService.saveAsync(rpc)).thenReturn(Futures.immediateFuture(null)); + public void savePersistsViaUpdateAsyncThenPushesToRuleEngine() { + Rpc rpc = newRpc(); + when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(null)); 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)) .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; + } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java index deb3bf8cf6..9dc821fcc7 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java +++ b/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); - ListenableFuture saveAsync(Rpc rpc); + ListenableFuture createAsync(Rpc rpc); + + ListenableFuture updateAsync(Rpc rpc); void deleteRpc(TenantId tenantId, RpcId id); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java index 4162422e69..e374aaf9ce 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java @@ -55,9 +55,15 @@ public class BaseRpcService implements RpcService { } @Override - public ListenableFuture saveAsync(Rpc rpc) { - log.trace("Executing saveAsync, [{}]", rpc); - return rpcDao.saveAsync(rpc.getTenantId(), rpc); + public ListenableFuture createAsync(Rpc rpc) { + log.trace("Executing createAsync, [{}]", rpc); + return rpcDao.createAsync(rpc.getTenantId(), rpc); + } + + @Override + public ListenableFuture updateAsync(Rpc rpc) { + log.trace("Executing updateAsync, [{}]", rpc); + return rpcDao.updateAsync(rpc.getTenantId(), rpc); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java index a823b05507..7bf4f5f8ac 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java +++ b/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 { - ListenableFuture saveAsync(TenantId tenantId, Rpc rpc); + ListenableFuture createAsync(TenantId tenantId, Rpc rpc); + + ListenableFuture updateAsync(TenantId tenantId, Rpc rpc); PageData findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java index e46dc2a704..45c765e42f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java @@ -71,7 +71,7 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, @Value("${sql.batch_sort:true}") private boolean batchSortEnabled; - private TbSqlBlockingQueueWrapper queue; + private TbSqlBlockingQueueWrapper queue; @PostConstruct private void init() { @@ -84,10 +84,10 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, .batchSortEnabled(batchSortEnabled) .withResponse(false) .build(); - Function hashcodeFunction = entity -> entity.getUuid().hashCode(); + Function hashcodeFunction = entry -> entry.entity().getUuid().hashCode(); queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); - queue.init(logExecutor, entities -> rpcInsertRepository.saveOrUpdate(entities), - Comparator.comparing(RpcEntity::getUuid)); + queue.init(logExecutor, entries -> rpcInsertRepository.saveOrUpdate(entries), + Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid())); } @PreDestroy @@ -98,8 +98,13 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, } @Override - public ListenableFuture saveAsync(TenantId tenantId, Rpc rpc) { - return queue.add(new RpcEntity(rpc)); + public ListenableFuture createAsync(TenantId tenantId, Rpc rpc) { + return queue.add(RpcQueueEntry.forInsert(new RpcEntity(rpc))); + } + + @Override + public ListenableFuture updateAsync(TenantId tenantId, Rpc rpc) { + return queue.add(RpcQueueEntry.forUpdate(new RpcEntity(rpc))); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java index 5073702ef4..f3023b0c7c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java @@ -29,17 +29,33 @@ import java.util.List; @Repository 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) " + "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 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 entries) { + List inserts = entries.stream().filter(RpcQueueEntry::insert).map(RpcQueueEntry::entity).toList(); + List updates = entries.stream().filter(entry -> !entry.insert()).map(RpcQueueEntry::entity).toList(); transactionTemplate.execute(status -> { - jdbcTemplate.batchUpdate(INSERT_OR_UPDATE, new BatchPreparedStatementSetter() { - @Override - public void setValues(PreparedStatement ps, int i) throws SQLException { - RpcEntity rpc = entities.get(i); + // Inserts run first so a create and a status update for the same rpcId coalesced into one + // batch still apply in create -> update order. + if (!inserts.isEmpty()) { + batch(INSERT, inserts, (ps, rpc) -> { ps.setObject(1, rpc.getUuid()); ps.setLong(2, rpc.getCreatedTime()); ps.setObject(3, rpc.getTenantId()); @@ -49,17 +65,33 @@ public class RpcInsertRepository extends AbstractInsertRepository { ps.setString(7, toJsonStr(rpc.getResponse())); ps.setString(8, toJsonStr(rpc.getAdditionalInfo())); ps.setString(9, rpc.getStatus().name()); - } - - @Override - public int getBatchSize() { - return entities.size(); - } - }); + }); + } + if (!updates.isEmpty()) { + batch(UPDATE, updates, (ps, rpc) -> { + ps.setString(1, rpc.getStatus().name()); + ps.setString(2, toJsonStr(rpc.getResponse())); + ps.setObject(3, rpc.getUuid()); + }); + } return null; }); } + private void batch(String sql, List 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) { return node == null ? null : replaceNullChars(JacksonUtil.toString(node)); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcQueueEntry.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcQueueEntry.java new file mode 100644 index 0000000000..8b5a5f96cf --- /dev/null +++ b/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); + } +} diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java index f631ff2daf..d49e0c2355 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.sql.rpc; +import com.fasterxml.jackson.databind.JsonNode; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.JacksonUtil; @@ -70,7 +71,7 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}")); 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); assertThat(stored).isNotNull(); @@ -86,7 +87,7 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { update.setStatus(RpcStatus.DELIVERED); 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); 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. // 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. - var queuedFuture = rpcDao.saveAsync(TenantId.SYS_TENANT_ID, queued); - var deliveredFuture = rpcDao.saveAsync(TenantId.SYS_TENANT_ID, delivered); + var queuedFuture = rpcDao.createAsync(TenantId.SYS_TENANT_ID, queued); + var deliveredFuture = rpcDao.updateAsync(TenantId.SYS_TENANT_ID, delivered); queuedFuture.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}")); } + @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; + } + }