Repository navigation
Conversation
There was a problem hiding this comment.
Pull request overview
This PR addresses lost/dropped server-side RPC deliveries over CoAP by serializing replayed (and newly arriving) RPC messages per RPC observation subscription, avoiding concurrent downlink sends that can race and get dropped.
Changes:
- Adds per-RPC-observation buffering (
pendingRpcRequests) plus readiness/in-flight tracking to serialize CoAP RPC deliveries. - Updates RPC ACK tracking to use a session-scoped key (
sessionId:mid) rather than only MID, reducing cross-session collisions. - Adds an integration test covering buffered RPC replay after re-subscribing to the RPC observe resource.
Reviewed changes
Copilot reviewed 5 out of 5 changed files in this pull request and generated 5 comments.
Show a summary per file
| File | Description |
|---|---|
| common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportContext.java | Changes rpcAwaitingAck keying to a session-scoped String key. |
| common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/TbCoapObservationState.java | Adds per-observation RPC queue + readiness/in-flight flags to support serialized delivery. |
| common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java | Enqueues incoming RPCs and drains them one-at-a-time once RPC observe subscription is ready; updates awaiting-ack bookkeeping. |
| application/src/test/java/org/thingsboard/server/transport/coap/rpc/CoapServerSideRpcDefaultIntegrationTest.java | Adds a new test entrypoint for buffered RPC replay after resubscribe. |
| application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java | Implements the buffered replay test + a recording observe callback. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| } | ||
| processNextQueuedRpc(state); |
| String awaitingAckKey = getAwaitingAckKey(sessionId, requestId); | ||
| transportContext.getRpcAwaitingAck().put(awaitingAckKey, msg); | ||
| transportContext.getScheduler().schedule(() -> { | ||
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(awaitingAckKey); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update due to server timeout ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); | ||
| onRpcDeliveryFinished(state, rpcState); | ||
| } | ||
| }, Math.min(getTimeout(state, powerMode, profileSettings), msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); | ||
|
|
||
| response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { | ||
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(awaitingAckKey); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request DELIVERED status update ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, true, TransportServiceCallback.EMPTY); | ||
| onRpcDeliveryFinished(state, rpcState); | ||
| } | ||
| }, id -> { | ||
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(awaitingAckKey); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); | ||
| onRpcDeliveryFinished(state, rpcState); | ||
| } | ||
| })); | ||
| markInFlightCompleted = false; | ||
| } |
| private final AtomicInteger observeCounter = new AtomicInteger(0); | ||
| private final Queue<TransportProtos.ToDeviceRpcRequestMsg> pendingRpcRequests = new ArrayDeque<>(); | ||
| private volatile ObserveRelation observeRelation; | ||
| private volatile boolean ready; | ||
| private volatile boolean rpcInFlight; |
There was a problem hiding this comment.
the solution is acceptible for now
| private void assertRequestsMatch(List<String> expectedRequests, List<String> actualRequests) { | ||
| List<String> expectedSorted = new ArrayList<>(expectedRequests); | ||
| List<String> actualSorted = new ArrayList<>(actualRequests); | ||
| expectedSorted.sort(String::compareTo); | ||
| actualSorted.sort(String::compareTo); | ||
| assertEquals(expectedSorted, actualSorted); | ||
| } |
There was a problem hiding this comment.
during tests the IDs are deterministic, this is acceptible for now
| try { | ||
| rpcPayloads.add(JacksonUtil.toString(JacksonUtil.fromBytes(payload))); | ||
| } catch (Exception e) { | ||
| fail("Failed to decode CoAP RPC payload: " + e.getMessage()); | ||
| } |
acf3893 to
5ef52ed
Compare
| import static org.junit.Assert.assertEquals; | ||
| import static org.junit.Assert.assertNotNull; | ||
| import static org.junit.Assert.assertTrue; | ||
| import static org.junit.Assert.fail; |
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(key); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update due to server timeout ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); |
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(key); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request DELIVERED status update ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, true, TransportServiceCallback.EMPTY); |
| TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(key); | ||
| if (rpcRequestMsg != null) { | ||
| log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update ...", deviceId, sessionId, requestId); | ||
| transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); |
| error = "Failed to convert device RPC command to CoAP msg"; | ||
| } catch (Exception e) { | ||
| error = "Internal error: " + e.getMessage(); | ||
| state.getRpc().getPendingRpcRequests().offer(msg); |
|
Closing this PR because the code change was solving delivery sequencing inside the CoAP transport layer, but that turns out to be the wrong layer for the policy we actually want. The desired behavior is already supported by the actor-level RPC submit strategy: SEQUENTIAL_ON_ACK_FROM_DEVICE. After enabling that on the core service, RPC delivery became consistent and matched the documented semantics without transport-layer changes. We are also explicitly choosing not to replay RPCs after successful delivery ACK, because automatic replay can duplicate already-executed device actions if the response is lost. So instead of adding transport-side queuing and replay behavior, we will rely on the existing actor policy and treat retries as an explicit higher-level decision. |
When subscribing to the RPC resource via CoAP, currently pending RPC calls would all be emitted at once and sent concurrently, causing race conditions. Some (or sometimes most) of the RPC calls would end up being dropped and never reached the device.
This adds a queue per subscription that gets drained to send out messages in a serialized and controlled, race-free manner.