Skip to content

fix missing RPC calls via coap transport - #15

Closed
JcBernack wants to merge 2 commits into
release-4.3from
fix/server-side-rpc-via-coap
Closed

JcBernack wants to merge 2 commits into
release-4.3from
fix/server-side-rpc-via-coap

Conversation

@JcBernack

Copy link
Copy Markdown
Member

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.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +736 to +737
}
processNextQueuedRpc(state);
Comment on lines +668 to +695
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;
}
Comment on lines 34 to +38
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;

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

the solution is acceptible for now

Comment on lines +195 to +201
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);
}

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

during tests the IDs are deterministic, this is acceptible for now

Comment on lines +330 to +334
try {
rpcPayloads.add(JacksonUtil.toString(JacksonUtil.fromBytes(payload)));
} catch (Exception e) {
fail("Failed to decode CoAP RPC payload: " + e.getMessage());
}

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 5 out of 5 changed files in this pull request and generated 5 comments.

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);
@JcBernack

Copy link
Copy Markdown
Member Author

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.

@JcBernack JcBernack closed this Jun 29, 2026
@JcBernack
JcBernack deleted the fix/server-side-rpc-via-coap branch June 29, 2026 21:32
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants