From ef571175f742ecad9ad6fc491891c9006e6ce504 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Fri, 26 Jun 2026 17:04:24 +0300 Subject: [PATCH] fix(rpc): skip rule-engine notify when async update matches no row Async RPC persistence now returns a per-write Boolean: an INSERT-on-conflict always persists (true), while an UPDATE-by-id reports false when its WHERE id = ? matched no row (RPC deleted via TTL/manual delete). TbRpcService notifies the rule engine only when the write actually persisted, restoring the findById-null skip the async refactor had dropped. Also: drop the unused tenantId param from createAsync/updateAsync; make TbRpcService's callback-thread count constructor-injectable (removes test reflection) with a comment on why callback striping exists alongside the queue partitioning; correct the RPC batch_threads yml comment; assert the Boolean contract in JpaRpcDaoTest and add a notification-suppression unit test. --- .../server/service/rpc/TbRpcService.java | 41 +++++++++------- .../src/main/resources/thingsboard.yml | 2 +- .../server/service/rpc/TbRpcServiceTest.java | 24 +++++++--- .../server/dao/rpc/RpcService.java | 4 +- .../server/dao/rpc/BaseRpcService.java | 8 ++-- .../thingsboard/server/dao/rpc/RpcDao.java | 4 +- .../server/dao/sql/rpc/JpaRpcDao.java | 14 ++++-- .../dao/sql/rpc/RpcInsertRepository.java | 24 +++++++--- .../server/dao/sql/rpc/JpaRpcDaoTest.java | 47 +++++++------------ 9 files changed, 96 insertions(+), 72 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 486feab2e2..da90db7846 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 @@ -16,9 +16,7 @@ package org.thingsboard.server.service.rpc; import com.google.common.util.concurrent.ListenableFuture; -import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; @@ -46,20 +44,24 @@ import java.util.concurrent.Executors; @TbCoreComponent @Service -@RequiredArgsConstructor @Slf4j public class TbRpcService { private final RpcService rpcService; private final TbClusterService tbClusterService; - @Value("${sql.rpc.callback_threads:3}") - private int callbackThreads; + // Post-persist rule-engine notifications run on these striped single-thread executors, keyed by + // rpcId, instead of inline on the SQL persist threads. A flush batch can carry up to batch_size + // entries, and each notification's pushMsgToRuleEngine may do a (potentially blocking) device-profile + // cache lookup; running them inline would serialize that work on the persist thread and stall the + // next DB flush. Striping by rpcId (rather than a shared multi-threaded pool) also preserves + // per-command notification order, e.g. RPC_QUEUED before RPC_DELIVERED. + private final ExecutorService[] callbackExecutors; - private ExecutorService[] callbackExecutors; - - @PostConstruct - private void init() { - callbackExecutors = new ExecutorService[callbackThreads]; + public TbRpcService(RpcService rpcService, TbClusterService tbClusterService, + @Value("${sql.rpc.callback_threads:3}") int callbackThreads) { + this.rpcService = rpcService; + this.tbClusterService = tbClusterService; + this.callbackExecutors = new ExecutorService[callbackThreads]; for (int i = 0; i < callbackThreads; i++) { callbackExecutors[i] = Executors.newSingleThreadExecutor( ThingsBoardThreadFactory.forName("rpc-persist-callback-" + i)); @@ -68,10 +70,8 @@ public class TbRpcService { @PreDestroy private void destroy() { - if (callbackExecutors != null) { - for (ExecutorService executor : callbackExecutors) { - executor.shutdownNow(); - } + for (ExecutorService executor : callbackExecutors) { + executor.shutdownNow(); } } @@ -83,16 +83,23 @@ public class TbRpcService { persist(tenantId, rpc, rpcService.updateAsync(rpc)); } - private void persist(TenantId tenantId, Rpc rpc, ListenableFuture future) { + private void persist(TenantId tenantId, Rpc rpc, ListenableFuture future) { DonAsynchron.withCallback(future, - v -> pushRpcMsgToRuleEngine(tenantId, rpc), + persisted -> { + if (Boolean.TRUE.equals(persisted)) { + pushRpcMsgToRuleEngine(tenantId, rpc); + } else { + log.debug("[{}][{}][{}] Skipping rule engine notification for status [{}] - RPC row no longer exists", + tenantId, rpc.getDeviceId(), rpc.getId(), rpc.getStatus()); + } + }, t -> log.error("[{}][{}][{}] Failed to persist RPC with status [{}]", tenantId, rpc.getDeviceId(), rpc.getId(), rpc.getStatus(), t), executorFor(rpc.getUuidId())); } private Executor executorFor(UUID rpcId) { - return callbackExecutors[(rpcId.hashCode() & 0x7FFFFFFF) % callbackThreads]; + return callbackExecutors[(rpcId.hashCode() & 0x7FFFFFFF) % callbackExecutors.length]; } private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f498d222ac..19075e7305 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -428,7 +428,7 @@ sql: batch_size: "${SQL_RPC_BATCH_SIZE:1000}" # Batch size for persisting RPC inserts/updates 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}" # batch thread count has to be a prime number like 3 or 5 to gain perfect hash distribution + 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 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 events: batch_size: "${SQL_EVENTS_BATCH_SIZE:10000}" # Batch size for persisting latest telemetry updates 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 014ab256af..87162c1123 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,7 +18,6 @@ package org.thingsboard.server.service.rpc; import com.google.common.util.concurrent.Futures; import org.junit.Before; import org.junit.Test; -import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.id.DeviceId; @@ -34,6 +33,7 @@ import java.util.UUID; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.after; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; @@ -47,15 +47,13 @@ public class TbRpcServiceTest { @Before public void setUp() { - tbRpcService = new TbRpcService(rpcService, clusterService); - ReflectionTestUtils.setField(tbRpcService, "callbackThreads", 1); - ReflectionTestUtils.invokeMethod(tbRpcService, "init"); + tbRpcService = new TbRpcService(rpcService, clusterService, 1); } @Test public void savePersistsViaUpdateAsyncThenPushesToRuleEngine() { Rpc rpc = newRpc(); - when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(null)); + when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(true)); tbRpcService.save(rpc.getTenantId(), rpc); @@ -67,7 +65,7 @@ public class TbRpcServiceTest { @Test public void createPersistsViaCreateAsyncThenPushesToRuleEngine() { Rpc rpc = newRpc(); - when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(null)); + when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(true)); tbRpcService.create(rpc.getTenantId(), rpc); @@ -76,6 +74,20 @@ public class TbRpcServiceTest { .pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull()); } + @Test + public void saveDoesNotNotifyRuleEngineWhenRowMissing() { + Rpc rpc = newRpc(); + // updateAsync resolves false: the UPDATE matched no row (RPC was deleted), so the rule engine + // must not be notified for a status change that never persisted. + when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(false)); + + tbRpcService.save(rpc.getTenantId(), rpc); + + verify(rpcService).updateAsync(rpc); + verify(clusterService, after(500).never()) + .pushMsgToRuleEngine(any(TenantId.class), any(DeviceId.class), any(TbMsg.class), isNull()); + } + private Rpc newRpc() { Rpc rpc = new Rpc(new RpcId(UUID.randomUUID())); rpc.setTenantId(TenantId.SYS_TENANT_ID); 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 9dc821fcc7..aaac66c86b 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,9 +29,9 @@ public interface RpcService extends EntityDaoService { Rpc save(Rpc rpc); - ListenableFuture createAsync(Rpc rpc); + ListenableFuture createAsync(Rpc rpc); - ListenableFuture updateAsync(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 e374aaf9ce..759b90f3a2 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,15 +55,15 @@ public class BaseRpcService implements RpcService { } @Override - public ListenableFuture createAsync(Rpc rpc) { + public ListenableFuture createAsync(Rpc rpc) { log.trace("Executing createAsync, [{}]", rpc); - return rpcDao.createAsync(rpc.getTenantId(), rpc); + return rpcDao.createAsync(rpc); } @Override - public ListenableFuture updateAsync(Rpc rpc) { + public ListenableFuture updateAsync(Rpc rpc) { log.trace("Executing updateAsync, [{}]", rpc); - return rpcDao.updateAsync(rpc.getTenantId(), rpc); + return rpcDao.updateAsync(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 7bf4f5f8ac..b0dde0d969 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,9 +26,9 @@ import org.thingsboard.server.dao.Dao; public interface RpcDao extends Dao { - ListenableFuture createAsync(TenantId tenantId, Rpc rpc); + ListenableFuture createAsync(Rpc rpc); - ListenableFuture updateAsync(TenantId tenantId, Rpc rpc); + ListenableFuture updateAsync(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 45c765e42f..d7fe95b3ca 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,10 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, @Value("${sql.batch_sort:true}") private boolean batchSortEnabled; - private TbSqlBlockingQueueWrapper queue; + // Boolean response per queued write: an INSERT-on-conflict always persists (true), while an + // UPDATE-by-id reports false when it matched no row (the RPC was deleted in the meantime), so the + // service layer can skip the rule-engine notification for a row that no longer exists. + private TbSqlBlockingQueueWrapper queue; @PostConstruct private void init() { @@ -82,12 +85,13 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, .statsPrintIntervalMs(statsPrintIntervalMs) .statsNamePrefix("rpc") .batchSortEnabled(batchSortEnabled) - .withResponse(false) + .withResponse(true) .build(); Function hashcodeFunction = entry -> entry.entity().getUuid().hashCode(); queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); queue.init(logExecutor, entries -> rpcInsertRepository.saveOrUpdate(entries), - Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid())); + Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid()), + Function.identity()); } @PreDestroy @@ -98,12 +102,12 @@ public class JpaRpcDao extends JpaAbstractDao implements RpcDao, } @Override - public ListenableFuture createAsync(TenantId tenantId, Rpc rpc) { + public ListenableFuture createAsync(Rpc rpc) { return queue.add(RpcQueueEntry.forInsert(new RpcEntity(rpc))); } @Override - public ListenableFuture updateAsync(TenantId tenantId, Rpc rpc) { + public ListenableFuture updateAsync(Rpc rpc) { return queue.add(RpcQueueEntry.forUpdate(new RpcEntity(rpc))); } 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 f3023b0c7c..9d0eb45464 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 @@ -24,6 +24,7 @@ import org.thingsboard.server.dao.sqlts.insert.AbstractInsertRepository; import java.sql.PreparedStatement; import java.sql.SQLException; +import java.util.ArrayList; import java.util.List; @Repository @@ -48,10 +49,10 @@ public class RpcInsertRepository extends AbstractInsertRepository { void bind(PreparedStatement ps, RpcEntity rpc) throws SQLException; } - public void saveOrUpdate(List entries) { + public List 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 -> { + int[] updateCounts = transactionTemplate.execute(status -> { // 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()) { @@ -68,18 +69,29 @@ public class RpcInsertRepository extends AbstractInsertRepository { }); } if (!updates.isEmpty()) { - batch(UPDATE, updates, (ps, rpc) -> { + return batch(UPDATE, updates, (ps, rpc) -> { ps.setString(1, rpc.getStatus().name()); ps.setString(2, toJsonStr(rpc.getResponse())); ps.setObject(3, rpc.getUuid()); }); } - return null; + return new int[0]; }); + + // Result is aligned to the submission order of entries: an insert always persists + // (INSERT ... ON CONFLICT), an update persists only if its WHERE id = ? matched a still-existing + // row. updateCounts keeps the same relative order as the filtered updates, so a single cursor + // walks it as we encounter update entries; insert entries short-circuit and don't advance it. + List persisted = new ArrayList<>(entries.size()); + int updateIdx = 0; + for (RpcQueueEntry entry : entries) { + persisted.add(entry.insert() || updateCounts[updateIdx++] > 0); + } + return persisted; } - private void batch(String sql, List entities, ColumnBinder binder) { - jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() { + private int[] batch(String sql, List entities, ColumnBinder binder) { + return jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() { @Override public void setValues(PreparedStatement ps, int i) throws SQLException { binder.bind(ps, entities.get(i)); 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 d49e0c2355..ad3f82c5bd 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 @@ -63,31 +63,19 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { @Test public void saveAsyncInsertThenUpsert() throws Exception { UUID id = UUID.randomUUID(); - Rpc rpc = new Rpc(new RpcId(id)); - rpc.setCreatedTime(System.currentTimeMillis()); - rpc.setTenantId(TenantId.SYS_TENANT_ID); - rpc.setDeviceId(new DeviceId(UUID.randomUUID())); - rpc.setExpirationTime(System.currentTimeMillis() + 60_000); - rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}")); - rpc.setStatus(RpcStatus.QUEUED); + DeviceId deviceId = new DeviceId(UUID.randomUUID()); - rpcDao.createAsync(rpc.getTenantId(), rpc).get(5, TimeUnit.SECONDS); + // A create always persists (INSERT ... ON CONFLICT), so the future resolves true. + assertThat(rpcDao.createAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS)).isTrue(); Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id); assertThat(stored).isNotNull(); assertThat(stored.getStatus()).isEqualTo(RpcStatus.QUEUED); assertThat(stored.getResponse()).isNull(); - Rpc update = new Rpc(new RpcId(id)); - update.setCreatedTime(rpc.getCreatedTime()); - update.setTenantId(TenantId.SYS_TENANT_ID); - update.setDeviceId(rpc.getDeviceId()); - update.setExpirationTime(rpc.getExpirationTime()); - update.setRequest(rpc.getRequest()); - update.setStatus(RpcStatus.DELIVERED); - update.setResponse(JacksonUtil.toJsonNode("{\"ok\":true}")); - - rpcDao.updateAsync(update.getTenantId(), update).get(5, TimeUnit.SECONDS); + // The update matches the existing row, so the future resolves true. + assertThat(rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.DELIVERED, JacksonUtil.toJsonNode("{\"ok\":true}"))) + .get(5, TimeUnit.SECONDS)).isTrue(); Rpc afterUpdate = rpcDao.findById(TenantId.SYS_TENANT_ID, id); assertThat(afterUpdate.getStatus()).isEqualTo(RpcStatus.DELIVERED); @@ -121,8 +109,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.createAsync(TenantId.SYS_TENANT_ID, queued); - var deliveredFuture = rpcDao.updateAsync(TenantId.SYS_TENANT_ID, delivered); + var queuedFuture = rpcDao.createAsync(queued); + var deliveredFuture = rpcDao.updateAsync(delivered); queuedFuture.get(5, TimeUnit.SECONDS); deliveredFuture.get(5, TimeUnit.SECONDS); @@ -138,12 +126,12 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { DeviceId deviceId = new DeviceId(UUID.randomUUID()); // Initial create. - rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS); + rpcDao.createAsync(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); + rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS); Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id); assertThat(stored).isNotNull(); @@ -155,14 +143,14 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { 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); + rpcDao.createAsync(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}"))) + rpcDao.updateAsync(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); + rpcDao.updateAsync(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); @@ -174,15 +162,16 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { 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); + rpcDao.createAsync(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); + // A late status update must not re-create the deleted row, and the future must resolve false + // (no row matched) so the service layer can skip the rule-engine notification. + assertThat(rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}"))) + .get(5, TimeUnit.SECONDS)).isFalse(); assertThat(rpcDao.findById(TenantId.SYS_TENANT_ID, id)).isNull(); }