From 0a987bb55792e9d7bedfe4c873f5e32b8bbe4d40 Mon Sep 17 00:00:00 2001 From: b14ckyy <33039058+b14ckyy@users.noreply.github.com> Date: Fri, 25 Sep 2026 22:42:07 +0200 Subject: [PATCH 1/2] MAVLink tunnel: resume MSP replies that do not fit the TX buffer mavlinkSendTunnelMspReply() wrote all TUNNEL chunks of a reply back to back, and mavlinkSendMessage() drops any frame that does not fit the port's free TX space. Hardware UART TX rings are 256 bytes and a full chunk is 145 bytes on the wire, so the second chunk of every reply above about 221 framed bytes was dropped deterministically. MSP_BOXNAMES never completed over a UART; USB VCP (4 KB buffer) and SITL hid it. Keep the encoded reply and a resume offset as a single pending state (the reply payload buffer is already shared, so one reply at a time), send each chunk only after checking the ingress port's free TX space, and continue in the following telemetry cycles before RX and the periodic stream. Nothing blocks. A request that arrives while a reply is pending is dropped; a reply with no progress for one second is abandoned so a stalled port cannot lock the tunnel. The pending state is cleared on port re-init, the flush respects the half-duplex backoff, and the MSP_REBOOT reply is flushed before the reboot. Unit tests get a settable TX budget to reproduce the 256-byte ring; seven tunnel tests cover the resume, the held-back chunk, busy drops on the same and on a second port, the stall, the re-init and the reboot. Co-Authored-By: Claude Fable 5.1 --- docs/Mavlink.md | 1 + src/main/fc/fc_mavlink.c | 100 ++++++--- src/main/fc/fc_mavlink.h | 3 + src/main/mavlink/mavlink_internal.h | 11 + src/main/mavlink/mavlink_ports.c | 3 + src/main/mavlink/mavlink_runtime.c | 80 ++++--- src/main/mavlink/mavlink_runtime.h | 1 + src/test/unit/mavlink_unittest.cc | 319 ++++++++++++++++++++++++++-- 8 files changed, 454 insertions(+), 64 deletions(-) diff --git a/docs/Mavlink.md b/docs/Mavlink.md index 1e27cb6a3aa..70899432a32 100644 --- a/docs/Mavlink.md +++ b/docs/Mavlink.md @@ -264,6 +264,7 @@ CLI mode is unavailable in MSP-over-MAVLink. - `target_component` may be `0` or `MAV_COMP_ID_AUTOPILOT1`. - `target_component = 0` is handled on the ingress MAVLink port only; it is not fanned out to other local MAVLink ports. - MSP replies are sent back to the requester as one or more `TUNNEL` messages on that same ingress port. +- A reply that does not fit into the port's free TX buffer at once is sent over the following telemetry cycles instead of being dropped. The tunnel serves one request at a time: a request that arrives while a reply is still being sent is discarded, so wait for the complete reply before sending the next request; a reply that makes no progress for one second is abandoned. - MSP framing is preserved end-to-end: MSPv1 requests get MSPv1 replies, and MSPv2 requests get MSPv2 replies. - Reboot (`MSP_REBOOT`) is supported over the tunnel. Serial passthrough and ESC 4way passthrough are rejected before execution. diff --git a/src/main/fc/fc_mavlink.c b/src/main/fc/fc_mavlink.c index 5ef8b690304..99bf4f34e57 100644 --- a/src/main/fc/fc_mavlink.c +++ b/src/main/fc/fc_mavlink.c @@ -20,43 +20,82 @@ static void mavlinkReleaseTunnelOwnerIfIdle(uint8_t ingressPortIndex) } } -static void mavlinkSendTunnelReply(uint8_t targetSystem, uint8_t targetComponent, const uint8_t *payload, uint8_t payloadLength) +static bool mavlinkTunnelMspReplyIsPending(void) { - uint8_t tunnelPayload[MAVLINK_MSG_TUNNEL_FIELD_PAYLOAD_LEN] = { 0 }; - memcpy(tunnelPayload, payload, payloadLength); - - mavlink_msg_tunnel_pack( - mavSystemId, - mavComponentId, - &mavSendMsg, - targetSystem, - targetComponent, - MAVLINK_TUNNEL_PAYLOAD_TYPE_INAV_MSP, - payloadLength, - tunnelPayload); - mavlinkSendMessage(); + return mavTunnelPendingReply.offset < mavTunnelPendingReply.frame.totalLength; } -static bool mavlinkSendTunnelMspReply(uint8_t targetSystem, uint8_t targetComponent, mspPacket_t *reply, uint8_t *replyPayloadHead, mspVersion_e mspVersion) +// Must not block on TX drain; the telemetry cycle resumes whatever does not fit now. +void mavlinkFlushTunnelMspReply(uint8_t portIndex) { - sbufSwitchToReader(&reply->buf, replyPayloadHead); + mavlinkTunnelPendingReply_t *pending = &mavTunnelPendingReply; + if (!mavlinkTunnelMspReplyIsPending()) { + return; + } - mspEncodedFrame_t frame; - if (!mspSerialPrepareFrame(reply, mspVersion, &frame)) { - return false; + // Abandoned on any port's call, so a reply stuck in half-duplex backoff or a stalled TX cannot lock the tunnel. + if ((millis() - pending->lastProgressMs) >= MAVLINK_TUNNEL_MSP_TIMEOUT_MS) { + memset(pending, 0, sizeof(*pending)); + return; } - uint8_t chunk[MAVLINK_MSG_TUNNEL_FIELD_PAYLOAD_LEN]; - for (int offset = 0; offset < frame.totalLength; offset += MAVLINK_MSG_TUNNEL_FIELD_PAYLOAD_LEN) { - const int chunkLength = mspSerialFrameCopyRange(&frame, offset, chunk, sizeof(chunk)); + if (pending->portIndex != portIndex) { + return; + } + + while (mavlinkTunnelMspReplyIsPending()) { + uint8_t chunk[MAVLINK_MSG_TUNNEL_FIELD_PAYLOAD_LEN] = { 0 }; + const int chunkLength = mspSerialFrameCopyRange(&pending->frame, pending->offset, chunk, sizeof(chunk)); if (chunkLength <= 0) { - return false; + break; } - mavlinkSendTunnelReply(targetSystem, targetComponent, chunk, chunkLength); + + mavlink_msg_tunnel_pack( + mavSystemId, + mavComponentId, + &mavSendMsg, + pending->targetSystem, + pending->targetComponent, + MAVLINK_TUNNEL_PAYLOAD_TYPE_INAV_MSP, + chunkLength, + chunk); + if (!mavlinkSendMessageToPortIfRoom(portIndex)) { + return; + } + + pending->offset += chunkLength; + pending->lastProgressMs = millis(); } + + memset(pending, 0, sizeof(*pending)); +} + +static bool mavlinkSendTunnelMspReply(uint8_t ingressPortIndex, uint8_t targetSystem, uint8_t targetComponent, mspPacket_t *reply, uint8_t *replyPayloadHead, mspVersion_e mspVersion) +{ + sbufSwitchToReader(&reply->buf, replyPayloadHead); + + mavlinkTunnelPendingReply_t *pending = &mavTunnelPendingReply; + if (!mspSerialPrepareFrame(reply, mspVersion, &pending->frame)) { + memset(pending, 0, sizeof(*pending)); + return false; + } + + pending->offset = 0; + pending->portIndex = ingressPortIndex; + pending->targetSystem = targetSystem; + pending->targetComponent = targetComponent; + pending->lastProgressMs = millis(); + mavlinkFlushTunnelMspReply(ingressPortIndex); return true; } +static bool mavlinkTunnelMspReplyIsBusy(uint8_t ingressPortIndex) +{ + // Flushes only the ingress port; writing another port here would bypass its half-duplex backoff. + mavlinkFlushTunnelMspReply(ingressPortIndex); + return mavlinkTunnelMspReplyIsPending(); +} + static bool mavlinkTunnelMessageTargetsLocalFc(const mavlink_tunnel_t *msg) { return msg->payload_type == MAVLINK_TUNNEL_PAYLOAD_TYPE_INAV_MSP && @@ -94,11 +133,18 @@ static bool mavlinkProcessCompletedTunnelCommand(uint8_t ingressPortIndex) }; uint8_t *replyPayloadHead = reply.buf.ptr; + // Clients wait for the full reply before the next request, so only pipelining clients or a second port hit this. + if (mavlinkTunnelMspReplyIsBusy(ingressPortIndex)) { + mspPort->c_state = MSP_IDLE; + return false; + } + if (mspPort->cmdMSP == MSP_SET_PASSTHROUGH) { reply.cmd = MSP_SET_PASSTHROUGH; reply.result = MSP_RESULT_ERROR; mspPort->c_state = MSP_IDLE; mavlinkSendTunnelMspReply( + ingressPortIndex, mavlinkContext.recvMsg.sysid, mavlinkContext.recvMsg.compid, &reply, @@ -120,6 +166,7 @@ static bool mavlinkProcessCompletedTunnelCommand(uint8_t ingressPortIndex) if (status != MSP_RESULT_NO_REPLY) { mavlinkSendTunnelMspReply( + ingressPortIndex, mavlinkContext.recvMsg.sysid, mavlinkContext.recvMsg.compid, &reply, @@ -132,6 +179,11 @@ static bool mavlinkProcessCompletedTunnelCommand(uint8_t ingressPortIndex) } waitForSerialPortToFinishTransmitting(mavPortStates[ingressPortIndex].port); + // A reply left pending by a full TX ring would otherwise be lost to the reboot. + if (mavlinkTunnelMspReplyIsPending()) { + mavlinkFlushTunnelMspReply(ingressPortIndex); + waitForSerialPortToFinishTransmitting(mavPortStates[ingressPortIndex].port); + } mspPostProcessFn(mavPortStates[ingressPortIndex].port); return true; } diff --git a/src/main/fc/fc_mavlink.h b/src/main/fc/fc_mavlink.h index dd4543318c9..0e0993241c5 100644 --- a/src/main/fc/fc_mavlink.h +++ b/src/main/fc/fc_mavlink.h @@ -9,3 +9,6 @@ typedef enum { } mavlinkFcDispatchResult_e; mavlinkFcDispatchResult_e mavlinkFcDispatchIncomingMessage(uint8_t ingressPortIndex); +#ifdef USE_MAVLINK_MSP_TUNNEL +void mavlinkFlushTunnelMspReply(uint8_t portIndex); +#endif diff --git a/src/main/mavlink/mavlink_internal.h b/src/main/mavlink/mavlink_internal.h index 513a570f9a1..47a0fd842de 100644 --- a/src/main/mavlink/mavlink_internal.h +++ b/src/main/mavlink/mavlink_internal.h @@ -85,6 +85,15 @@ // rather than overflows, but sbufWrite*() is not bounds-checked in general, so // re-audit the handlers in fc_msp.c before shrinking this. #define MAVLINK_TUNNEL_MSP_REPLY_BUF_SIZE 768 + +typedef struct mavlinkTunnelPendingReply_s { + mspEncodedFrame_t frame; // payload points into the shared tunnelReplyPayloadBuf, so one reply at a time + uint16_t offset; + uint8_t portIndex; + uint8_t targetSystem; + uint8_t targetComponent; + timeMs_t lastProgressMs; +} mavlinkTunnelPendingReply_t; #endif #define MAVLINK_MISSION_UPLOAD_RETRY_MS 1500 #define MAVLINK_MISSION_UPLOAD_MAX_RETRIES 5 @@ -102,6 +111,7 @@ typedef struct mavlinkContext_s { uint8_t tunnelRemoteSystemIds[MAX_MAVLINK_PORTS]; uint8_t tunnelRemoteComponentIds[MAX_MAVLINK_PORTS]; uint8_t tunnelReplyPayloadBuf[MAVLINK_TUNNEL_MSP_REPLY_BUF_SIZE]; + mavlinkTunnelPendingReply_t tunnelPendingReply; #endif uint8_t sendMask; mavlinkPortRuntime_t *activePort; @@ -133,6 +143,7 @@ extern mavlinkContext_t mavlinkContext; #define mavTunnelRemoteSystemIds (mavlinkContext.tunnelRemoteSystemIds) #define mavTunnelRemoteComponentIds (mavlinkContext.tunnelRemoteComponentIds) #define mavTunnelReplyPayloadBuf (mavlinkContext.tunnelReplyPayloadBuf) +#define mavTunnelPendingReply (mavlinkContext.tunnelPendingReply) #endif #define mavSendMask (mavlinkContext.sendMask) #define mavActivePort (mavlinkContext.activePort) diff --git a/src/main/mavlink/mavlink_ports.c b/src/main/mavlink/mavlink_ports.c index 5a9cbf6d224..498485ba559 100644 --- a/src/main/mavlink/mavlink_ports.c +++ b/src/main/mavlink/mavlink_ports.c @@ -40,6 +40,9 @@ static void resetMAVLinkPortRuntimeState(uint8_t portIndex) memset(&state->mlrs, 0, sizeof(state->mlrs)); #ifdef USE_MAVLINK_MSP_TUNNEL mavlinkResetTunnelPortState(portIndex); + if (mavTunnelPendingReply.portIndex == portIndex) { + memset(&mavTunnelPendingReply, 0, sizeof(mavTunnelPendingReply)); + } #endif } diff --git a/src/main/mavlink/mavlink_runtime.c b/src/main/mavlink/mavlink_runtime.c index 6fce0de2422..015d46b04aa 100644 --- a/src/main/mavlink/mavlink_runtime.c +++ b/src/main/mavlink/mavlink_runtime.c @@ -133,6 +133,28 @@ void mavlinkRuntimeCheckState(void) } } +static int mavlinkEncodeMessageForPort(const mavlinkPortRuntime_t *state, const mavlink_msg_entry_t *msgEntry, uint8_t *mavBuffer) +{ + mavlink_status_t txStatus = { 0 }; + txStatus.current_tx_seq = state->txSeq; + if (mavlinkGetProtocolVersion() == 1) { + txStatus.flags |= MAVLINK_STATUS_FLAG_OUT_MAVLINK1; + } + + mavlink_message_t txMsg = mavSendMsg; + mavlink_finalize_message_buffer( + &txMsg, + txMsg.sysid, + txMsg.compid, + &txStatus, + msgEntry->min_msg_len, + txMsg.len, + msgEntry->crc_extra + ); + + return mavlink_msg_to_send_buffer(mavBuffer, &txMsg); +} + void mavlinkSendMessage(void) { const mavlink_msg_entry_t *msgEntry = mavlink_get_msg_entry(mavSendMsg.msgid); @@ -164,26 +186,9 @@ void mavlinkSendMessage(void) continue; } - mavlink_status_t txStatus = { 0 }; - txStatus.current_tx_seq = state->txSeq; - if (mavlinkGetProtocolVersion() == 1) { - txStatus.flags |= MAVLINK_STATUS_FLAG_OUT_MAVLINK1; - } - - mavlink_message_t txMsg = mavSendMsg; - mavlink_finalize_message_buffer( - &txMsg, - txMsg.sysid, - txMsg.compid, - &txStatus, - msgEntry->min_msg_len, - txMsg.len, - msgEntry->crc_extra - ); - state->txSeq = txStatus.current_tx_seq; - uint8_t mavBuffer[MAVLINK_MAX_PACKET_LEN]; - const int msgLength = mavlink_msg_to_send_buffer(mavBuffer, &txMsg); + const int msgLength = mavlinkEncodeMessageForPort(state, msgEntry, mavBuffer); + state->txSeq++; if (msgLength <= 0) { continue; } @@ -200,6 +205,27 @@ void mavlinkSendMessage(void) } } +bool mavlinkSendMessageToPortIfRoom(uint8_t portIndex) +{ + const mavlink_msg_entry_t *msgEntry = mavlink_get_msg_entry(mavSendMsg.msgid); + mavlinkPortRuntime_t *state = &mavPortStates[portIndex]; + if (!msgEntry || !state->telemetryEnabled || !state->port) { + return false; + } + + uint8_t mavBuffer[MAVLINK_MAX_PACKET_LEN]; + const int msgLength = mavlinkEncodeMessageForPort(state, msgEntry, mavBuffer); + if (msgLength <= 0 || serialTxBytesFree(state->port) < (uint32_t)msgLength) { + return false; + } + + state->txSeq++; + serialBeginWrite(state->port); + serialWriteBuf(state->port, mavBuffer, msgLength); + serialEndWrite(state->port); + return true; +} + static bool processMAVLinkIncomingTelemetry(uint8_t ingressPortIndex, timeUs_t currentTimeUs) { mavlinkPortRuntime_t *state = &mavPortStates[ingressPortIndex]; @@ -266,7 +292,7 @@ static bool processMAVLinkIncomingTelemetry(uint8_t ingressPortIndex, timeUs_t c return false; } -static bool isMAVLinkTelemetryHalfDuplex(uint8_t portIndex) +static bool isMAVLinkTelemetryHalfDuplexBackoff(uint8_t portIndex, timeUs_t currentTimeUs) { const mavlinkPortRuntime_t *state = &mavPortStates[portIndex]; @@ -274,7 +300,8 @@ static bool isMAVLinkTelemetryHalfDuplex(uint8_t portIndex) (state->portConfig->functionMask & FUNCTION_RX_SERIAL) && rxConfig()->receiverType == RX_TYPE_SERIAL && rxConfig()->serialrx_provider == SERIALRX_MAVLINK && - tristateWithDefaultOffIsActive(rxConfig()->halfDuplex); + tristateWithDefaultOffIsActive(rxConfig()->halfDuplex) && + ((currentTimeUs - state->lastRxFrameUs) < TELEMETRY_MAVLINK_DELAY); } void mavlinkRuntimeHandle(timeUs_t currentTimeUs) @@ -292,16 +319,21 @@ void mavlinkRuntimeHandle(timeUs_t currentTimeUs) mavlinkSetActivePortContext(portIndex); +#ifdef USE_MAVLINK_MSP_TUNNEL + // Before RX and the periodic stream, so a pending reply gets TX space first. + if (!isMAVLinkTelemetryHalfDuplexBackoff(portIndex, currentTimeUs)) { + mavlinkFlushTunnelMspReply(portIndex); + } +#endif + // Process incoming MAVLink on this port and forward when needed. processMAVLinkIncomingTelemetry(portIndex, currentTimeUs); // Restore context back to this port before periodic send decisions. mavlinkSetActivePortContext(portIndex); bool shouldSendTelemetry = false; - const bool halfDuplexBackoff = isMAVLinkTelemetryHalfDuplex(portIndex) && - ((currentTimeUs - state->lastRxFrameUs) < TELEMETRY_MAVLINK_DELAY); - if (halfDuplexBackoff) { + if (isMAVLinkTelemetryHalfDuplexBackoff(portIndex, currentTimeUs)) { continue; } diff --git a/src/main/mavlink/mavlink_runtime.h b/src/main/mavlink/mavlink_runtime.h index 988b7b552b7..37d1c37d47b 100644 --- a/src/main/mavlink/mavlink_runtime.h +++ b/src/main/mavlink/mavlink_runtime.h @@ -19,3 +19,4 @@ bool mavlinkPortTxBufferIsValid(uint8_t portIndex); uint8_t mavlinkPortTxBufferFree(uint8_t portIndex); void mavlinkSetActivePortContext(uint8_t portIndex); void mavlinkSendMessage(void); +bool mavlinkSendMessageToPortIfRoom(uint8_t portIndex); diff --git a/src/test/unit/mavlink_unittest.cc b/src/test/unit/mavlink_unittest.cc index a22461a8fa3..b1c76f0b938 100644 --- a/src/test/unit/mavlink_unittest.cc +++ b/src/test/unit/mavlink_unittest.cc @@ -80,6 +80,7 @@ extern "C" { #include "sensors/temperature.h" #include "mavlink/mavlink_types.h" + #include "mavlink/mavlink_internal.h" #include "telemetry/mavlink.h" #include "telemetry/telemetry.h" @@ -106,6 +107,19 @@ static uint8_t serialTxBuffer[2048]; static size_t serialRxLen; static size_t serialRxPos; static size_t serialTxLen; +// Negative: unlimited (stub reports 1024). Otherwise free TX bytes, consumed by writes, refilled by tests per cycle. +static int32_t serialTxBudget; +static bool serialTxOverrun; +static const int32_t testUartTxBufferFree = 255; +static bool testSecondPortEnabled; +static bool testSecondPortServed; +static serialPort_t testSerialPort2; +static serialPortConfig_t testPortConfig2; +static uint8_t serialRxBuffer2[512]; +static uint8_t serialTxBuffer2[2048]; +static size_t serialRxLen2; +static size_t serialRxPos2; +static size_t serialTxLen2; static const uint8_t testTargetComponent = MAV_COMP_ID_AUTOPILOT1; static const uint8_t testTunnelSourceSystem = 42; static const uint8_t testTunnelSourceComponent = 200; @@ -165,6 +179,9 @@ static void resetSerialBuffers(void) serialRxLen = 0; serialRxPos = 0; serialTxLen = 0; + serialRxLen2 = 0; + serialRxPos2 = 0; + serialTxLen2 = 0; } static std::vector makeMspV1Request(uint8_t cmd, const std::vector &payload = {}) @@ -210,10 +227,15 @@ static std::vector encodeMspV1Reply(uint8_t cmd, int16_t result, const return encodeMspReply(cmd, result, MSP_V1, payload); } -static void pushRxMessage(const mavlink_message_t *msg) +static void pushRxMessage(const mavlink_message_t *msg, bool secondPort = false) { uint8_t buffer[MAVLINK_MAX_PACKET_LEN]; int length = mavlink_msg_to_send_buffer(buffer, msg); + if (secondPort) { + memcpy(&serialRxBuffer2[serialRxLen2], buffer, (size_t)length); + serialRxLen2 += (size_t)length; + return; + } memcpy(&serialRxBuffer[serialRxLen], buffer, (size_t)length); serialRxLen += (size_t)length; } @@ -230,7 +252,8 @@ static void pushTunnelPayload( const std::vector &payload, uint8_t targetComponent = testTargetComponent, uint8_t sourceSystem = testTunnelSourceSystem, - uint8_t sourceComponent = testTunnelSourceComponent) + uint8_t sourceComponent = testTunnelSourceComponent, + bool secondPort = false) { uint8_t tunnelPayload[MAVLINK_MSG_TUNNEL_FIELD_PAYLOAD_LEN] = { 0 }; size_t copyLength = payload.size(); @@ -251,7 +274,7 @@ static void pushTunnelPayload( 0x8001, payloadLength, tunnelPayload); - pushRxMessage(&msg); + pushRxMessage(&msg, secondPort); } static bool popTxMessage(mavlink_message_t *msg) @@ -282,15 +305,15 @@ static bool findTxMessageById(uint32_t msgid, mavlink_message_t *match) return false; } -static std::vector parseTxMessages(void) +static std::vector parseTxMessagesFrom(const uint8_t *buffer, size_t length) { std::vector messages; mavlink_status_t status; memset(&status, 0, sizeof(status)); mavlink_message_t msg; - for (size_t i = 0; i < serialTxLen; i++) { - if (mavlink_parse_char(0, serialTxBuffer[i], &msg, &status) == MAVLINK_FRAMING_OK) { + for (size_t i = 0; i < length; i++) { + if (mavlink_parse_char(0, buffer[i], &msg, &status) == MAVLINK_FRAMING_OK) { messages.push_back(msg); } } @@ -298,6 +321,22 @@ static std::vector parseTxMessages(void) return messages; } +static std::vector parseTxMessages(void) +{ + return parseTxMessagesFrom(serialTxBuffer, serialTxLen); +} + +static std::vector filterTunnelMessages(const std::vector &messages) +{ + std::vector tunnelMessages; + for (const mavlink_message_t &msg : messages) { + if (msg.msgid == MAVLINK_MSG_ID_TUNNEL) { + tunnelMessages.push_back(msg); + } + } + return tunnelMessages; +} + static std::vector collectTunnelPayload( const std::vector &messages, uint8_t expectedTargetSystem = testTunnelSourceSystem, @@ -321,9 +360,12 @@ static std::vector collectTunnelPayload( return payload; } -static void initMavlinkTestState(void) +static void initMavlinkTestState(bool secondPort = false) { + testSecondPortEnabled = secondPort; resetSerialBuffers(); + serialTxBudget = -1; + serialTxOverrun = false; fakeMillis = 0; fakeMicros = 0; setWaypointCalls = 0; @@ -600,6 +642,213 @@ TEST(MavlinkTelemetryTest, TunnelReplyFragmentationIsExactAt127128And129Bytes) } } +static std::vector makeTestReplyPayload(uint16_t length) +{ + std::vector payload(length); + for (size_t i = 0; i < payload.size(); i++) { + payload[i] = (uint8_t)i; + } + return payload; +} + +static void runTelemetryCycleWithTxBudget(int32_t txBudget) +{ + serialTxBudget = txBudget; + handleMAVLinkTelemetry(1000); +} + +TEST(MavlinkTelemetryTest, TunnelLargeReplyResumesAcrossCyclesOnSmallTxBuffer) +{ + initMavlinkTestState(); + serialTxBudget = testUartTxBufferFree; + + const std::vector request = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)request.size(), request); + handleMAVLinkTelemetry(1000); + + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(parseTxMessages().size(), 1U); + + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + + const std::vector messages = parseTxMessages(); + ASSERT_EQ(messages.size(), 3U); + for (size_t i = 1; i < messages.size(); i++) { + EXPECT_EQ(messages[i].seq, (uint8_t)(messages[i - 1].seq + 1)); + } + EXPECT_EQ(collectTunnelPayload(messages), encodeMspV1Reply(testLargeReplyMspCommand, MSP_RESULT_ACK, makeTestReplyPayload(300))); + EXPECT_FALSE(serialTxOverrun); + EXPECT_EQ(mavlinkContext.portStates[0].txDroppedFrames, 0U); + + const size_t txLenAfterReply = serialTxLen; + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + EXPECT_EQ(serialTxLen, txLenAfterReply); +} + +TEST(MavlinkTelemetryTest, TunnelChunkIsHeldBackWhileItDoesNotFit) +{ + initMavlinkTestState(); + testReplyPayloadLength = 500; + serialTxBudget = 100; + + const std::vector request = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)request.size(), request); + handleMAVLinkTelemetry(1000); + + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(serialTxLen, 0U); + + runTelemetryCycleWithTxBudget(0); + runTelemetryCycleWithTxBudget(144); + EXPECT_EQ(serialTxLen, 0U); + + int cycles = 0; + while (parseTxMessages().size() < 4U && cycles < 20) { + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + cycles++; + } + + const std::vector messages = parseTxMessages(); + ASSERT_EQ(messages.size(), 4U); + EXPECT_EQ(cycles, 4); + for (size_t i = 1; i < messages.size(); i++) { + EXPECT_EQ(messages[i].seq, (uint8_t)(messages[i - 1].seq + 1)); + } + EXPECT_EQ(collectTunnelPayload(messages), encodeMspV1Reply(testLargeReplyMspCommand, MSP_RESULT_ACK, makeTestReplyPayload(500))); + EXPECT_FALSE(serialTxOverrun); + EXPECT_EQ(mavlinkContext.portStates[0].txDroppedFrames, 0U); +} + +TEST(MavlinkTelemetryTest, TunnelRequestWhileReplyPendingIsDropped) +{ + initMavlinkTestState(); + serialTxBudget = testUartTxBufferFree; + + const std::vector largeRequest = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)largeRequest.size(), largeRequest); + handleMAVLinkTelemetry(1000); + EXPECT_EQ(mspCommandCallCount, 1); + + const std::vector simpleRequest = makeMspV1Request(testSimpleMspCommand); + pushTunnelPayload((uint8_t)simpleRequest.size(), simpleRequest); + serialTxBudget = 0; + handleMavlinkUntilRxEmpty(1000); + + EXPECT_EQ(mspCommandCallCount, 1); + + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + + EXPECT_EQ( + collectTunnelPayload(parseTxMessages()), + encodeMspV1Reply(testLargeReplyMspCommand, MSP_RESULT_ACK, makeTestReplyPayload(300))); + EXPECT_FALSE(serialTxOverrun); + + resetSerialBuffers(); + pushTunnelPayload((uint8_t)simpleRequest.size(), simpleRequest); + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + + EXPECT_EQ(mspCommandCallCount, 2); + EXPECT_EQ(collectTunnelPayload(parseTxMessages()), encodeMspV1Reply(testSimpleMspCommand, MSP_RESULT_ACK)); +} + +TEST(MavlinkTelemetryTest, TunnelStalledReplyIsAbandonedAfterClientTimeout) +{ + initMavlinkTestState(); + serialTxBudget = testUartTxBufferFree; + + const std::vector largeRequest = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)largeRequest.size(), largeRequest); + handleMAVLinkTelemetry(1000); + + resetSerialBuffers(); + fakeMillis += 1000; + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + EXPECT_EQ(serialTxLen, 0U); + + const std::vector simpleRequest = makeMspV1Request(testSimpleMspCommand); + pushTunnelPayload((uint8_t)simpleRequest.size(), simpleRequest); + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + + EXPECT_EQ(mspCommandCallCount, 2); + EXPECT_EQ(collectTunnelPayload(parseTxMessages()), encodeMspV1Reply(testSimpleMspCommand, MSP_RESULT_ACK)); +} + +TEST(MavlinkTelemetryTest, TunnelPendingReplyIsDiscardedOnPortReinit) +{ + initMavlinkTestState(); + serialTxBudget = testUartTxBufferFree; + + const std::vector request = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)request.size(), request); + handleMAVLinkTelemetry(1000); + EXPECT_EQ(parseTxMessages().size(), 1U); + + freeMAVLinkTelemetryPort(); + checkMAVLinkTelemetryState(); + + resetSerialBuffers(); + runTelemetryCycleWithTxBudget(testUartTxBufferFree); + EXPECT_EQ(serialTxLen, 0U); +} + +TEST(MavlinkTelemetryTest, TunnelRequestOnOtherPortWhileReplyPendingIsDroppedWithoutWritingPendingPort) +{ + initMavlinkTestState(true); + ASSERT_EQ(mavlinkContext.portCount, 2); + + // Port 0 is a half-duplex MAVLink RX port, so its pending reply must wait out the RX backoff. + rxConfigMutable()->receiverType = RX_TYPE_SERIAL; + rxConfigMutable()->serialrx_provider = SERIALRX_MAVLINK; + rxConfigMutable()->halfDuplex = TRISTATE_ON; + testPortConfig.functionMask |= FUNCTION_RX_SERIAL; + serialTxBudget = testUartTxBufferFree; + + const std::vector largeRequest = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)largeRequest.size(), largeRequest); + handleMAVLinkTelemetry(1000); + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(filterTunnelMessages(parseTxMessages()).size(), 1U); + + const size_t pendingPortTxLen = serialTxLen; + const uint8_t otherSystem = testTunnelSourceSystem + 1; + const std::vector simpleRequest = makeMspV1Request(testSimpleMspCommand); + pushTunnelPayload((uint8_t)simpleRequest.size(), simpleRequest, testTargetComponent, otherSystem, testTunnelSourceComponent, true); + serialTxBudget = testUartTxBufferFree; + handleMAVLinkTelemetry(1000); + + EXPECT_EQ(serialRxPos2, serialRxLen2); + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(serialTxLen, pendingPortTxLen); + EXPECT_TRUE(filterTunnelMessages(parseTxMessagesFrom(serialTxBuffer2, serialTxLen2)).empty()); + + const timeUs_t afterBackoffUs = 1000 + TELEMETRY_MAVLINK_DELAY; + for (int cycle = 0; cycle < 5; cycle++) { + serialTxBudget = testUartTxBufferFree; + handleMAVLinkTelemetry(afterBackoffUs); + } + + EXPECT_EQ( + collectTunnelPayload(filterTunnelMessages(parseTxMessages())), + encodeMspV1Reply(testLargeReplyMspCommand, MSP_RESULT_ACK, makeTestReplyPayload(300))); + EXPECT_TRUE(filterTunnelMessages(parseTxMessagesFrom(serialTxBuffer2, serialTxLen2)).empty()); + EXPECT_FALSE(serialTxOverrun); +} + +TEST(MavlinkTelemetryTest, TunnelRebootReplyIsFlushedBeforeRebootWhenTxRingIsFull) +{ + initMavlinkTestState(); + serialTxBudget = 0; + + const std::vector request = makeMspV1Request((uint8_t)MSP_REBOOT); + pushTunnelPayload((uint8_t)request.size(), request); + handleMAVLinkTelemetry(1000); + + EXPECT_EQ(waitForSerialPortToFinishTransmittingCalls, 2); + EXPECT_EQ(mspRebootPostProcessCount, 1); + EXPECT_EQ(collectTunnelPayload(parseTxMessages()), encodeMspV1Reply((uint8_t)MSP_REBOOT, MSP_RESULT_ACK)); + EXPECT_FALSE(serialTxOverrun); +} + TEST(MavlinkTelemetryTest, PreparedFrameMatchesExactV1V2AndV2OverV1Bytes) { const std::vector payload = { 0x10, 0x20 }; @@ -3497,13 +3746,21 @@ serialPortConfig_t *findSerialPortConfig(serialPortFunction_e function) testPortConfig.functionMask = FUNCTION_TELEMETRY_MAVLINK; testPortConfig.identifier = SERIAL_PORT_USART1; testPortConfig.telemetry_baudrateIndex = BAUD_115200; + testPortConfig2.functionMask = FUNCTION_TELEMETRY_MAVLINK; + testPortConfig2.identifier = SERIAL_PORT_USART2; + testPortConfig2.telemetry_baudrateIndex = BAUD_115200; + testSecondPortServed = false; return &testPortConfig; } serialPortConfig_t *findNextSerialPortConfig(serialPortFunction_e function) { UNUSED(function); - return NULL; + if (!testSecondPortEnabled || testSecondPortServed) { + return NULL; + } + testSecondPortServed = true; + return &testPortConfig2; } // No mixer-profile switching is configured in these tests, so nothing here is a VTOL and @@ -3525,14 +3782,13 @@ serialPort_t *openSerialPort(serialPortIdentifier_e identifier, serialPortFuncti serialReceiveCallbackPtr rxCallback, void *rxCallbackData, uint32_t baudRate, portMode_t mode, portOptions_t options) { - UNUSED(identifier); UNUSED(function); UNUSED(rxCallback); UNUSED(rxCallbackData); UNUSED(baudRate); UNUSED(mode); UNUSED(options); - return &testSerialPort; + return identifier == testPortConfig2.identifier ? &testSerialPort2 : &testSerialPort; } void closeSerialPort(serialPort_t *serialPort) @@ -3542,31 +3798,59 @@ void closeSerialPort(serialPort_t *serialPort) uint32_t serialRxBytesWaiting(const serialPort_t *instance) { - UNUSED(instance); + if (instance == &testSerialPort2) { + return (uint32_t)(serialRxLen2 - serialRxPos2); + } return (uint32_t)(serialRxLen - serialRxPos); } uint32_t serialTxBytesFree(const serialPort_t *instance) { - UNUSED(instance); - return 1024; + if (instance == &testSerialPort2) { + return 1024; + } + return serialTxBudget < 0 ? 1024 : (uint32_t)serialTxBudget; +} + +static void consumeSerialTxBudget(int count) +{ + if (serialTxBudget < 0) { + return; + } + if (count > serialTxBudget) { + serialTxOverrun = true; + serialTxBudget = 0; + return; + } + serialTxBudget -= count; } uint8_t serialRead(serialPort_t *instance) { - UNUSED(instance); + if (instance == &testSerialPort2) { + return serialRxBuffer2[serialRxPos2++]; + } return serialRxBuffer[serialRxPos++]; } void serialWrite(serialPort_t *instance, uint8_t ch) { - UNUSED(instance); + if (instance == &testSerialPort2) { + serialTxBuffer2[serialTxLen2++] = ch; + return; + } + consumeSerialTxBudget(1); serialTxBuffer[serialTxLen++] = ch; } void serialWriteBuf(serialPort_t *instance, const uint8_t *data, int count) { - UNUSED(instance); + if (instance == &testSerialPort2) { + memcpy(&serialTxBuffer2[serialTxLen2], data, (size_t)count); + serialTxLen2 += (size_t)count; + return; + } + consumeSerialTxBudget(count); memcpy(&serialTxBuffer[serialTxLen], data, (size_t)count); serialTxLen += (size_t)count; } @@ -3607,6 +3891,9 @@ bool isSerialTransmitBufferEmpty(const serialPort_t *instance) void waitForSerialPortToFinishTransmitting(serialPort_t *serialPort) { + if (serialTxBudget >= 0) { + serialTxBudget = testUartTxBufferFree; + } waitForSerialPortToFinishTransmittingCalls++; lastPostProcessPort = serialPort; } From 8005a050db4f4ae25172b5892388c7b83d68b1f0 Mon Sep 17 00:00:00 2001 From: b14ckyy <33039058+b14ckyy@users.noreply.github.com> Date: Fri, 25 Sep 2026 22:53:31 +0200 Subject: [PATCH 2/2] MAVLink tunnel: do not transmit from the busy check The busy check flushed the pending reply from inside RX processing, right after lastRxFrameUs was updated, so on a half-duplex port a pipelined request could put a chunk on the line inside the backoff window. The per-cycle flush already runs before RX and respects the backoff, so the check now only reports whether a reply is pending. Co-Authored-By: Claude Fable 5.1 --- src/main/fc/fc_mavlink.c | 9 +------- src/test/unit/mavlink_unittest.cc | 35 +++++++++++++++++++++++++++++++ 2 files changed, 36 insertions(+), 8 deletions(-) diff --git a/src/main/fc/fc_mavlink.c b/src/main/fc/fc_mavlink.c index 99bf4f34e57..221d9ce8265 100644 --- a/src/main/fc/fc_mavlink.c +++ b/src/main/fc/fc_mavlink.c @@ -89,13 +89,6 @@ static bool mavlinkSendTunnelMspReply(uint8_t ingressPortIndex, uint8_t targetSy return true; } -static bool mavlinkTunnelMspReplyIsBusy(uint8_t ingressPortIndex) -{ - // Flushes only the ingress port; writing another port here would bypass its half-duplex backoff. - mavlinkFlushTunnelMspReply(ingressPortIndex); - return mavlinkTunnelMspReplyIsPending(); -} - static bool mavlinkTunnelMessageTargetsLocalFc(const mavlink_tunnel_t *msg) { return msg->payload_type == MAVLINK_TUNNEL_PAYLOAD_TYPE_INAV_MSP && @@ -134,7 +127,7 @@ static bool mavlinkProcessCompletedTunnelCommand(uint8_t ingressPortIndex) uint8_t *replyPayloadHead = reply.buf.ptr; // Clients wait for the full reply before the next request, so only pipelining clients or a second port hit this. - if (mavlinkTunnelMspReplyIsBusy(ingressPortIndex)) { + if (mavlinkTunnelMspReplyIsPending()) { mspPort->c_state = MSP_IDLE; return false; } diff --git a/src/test/unit/mavlink_unittest.cc b/src/test/unit/mavlink_unittest.cc index b1c76f0b938..4d383f1240e 100644 --- a/src/test/unit/mavlink_unittest.cc +++ b/src/test/unit/mavlink_unittest.cc @@ -834,6 +834,41 @@ TEST(MavlinkTelemetryTest, TunnelRequestOnOtherPortWhileReplyPendingIsDroppedWit EXPECT_FALSE(serialTxOverrun); } +TEST(MavlinkTelemetryTest, TunnelPipelinedRequestOnHalfDuplexPortDoesNotTransmitInsideBackoff) +{ + initMavlinkTestState(); + rxConfigMutable()->receiverType = RX_TYPE_SERIAL; + rxConfigMutable()->serialrx_provider = SERIALRX_MAVLINK; + rxConfigMutable()->halfDuplex = TRISTATE_ON; + testPortConfig.functionMask |= FUNCTION_RX_SERIAL; + serialTxBudget = testUartTxBufferFree; + + const std::vector largeRequest = makeMspV1Request(testLargeReplyMspCommand); + pushTunnelPayload((uint8_t)largeRequest.size(), largeRequest); + handleMAVLinkTelemetry(1000); + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(filterTunnelMessages(parseTxMessages()).size(), 1U); + + const size_t txLenBeforePipelinedRequest = serialTxLen; + const std::vector simpleRequest = makeMspV1Request(testSimpleMspCommand); + pushTunnelPayload((uint8_t)simpleRequest.size(), simpleRequest); + serialTxBudget = testUartTxBufferFree; + handleMAVLinkTelemetry(1000); + + EXPECT_EQ(serialRxPos, serialRxLen); + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ(serialTxLen, txLenBeforePipelinedRequest); + + serialTxBudget = testUartTxBufferFree; + handleMAVLinkTelemetry(1000 + TELEMETRY_MAVLINK_DELAY); + + EXPECT_EQ(mspCommandCallCount, 1); + EXPECT_EQ( + collectTunnelPayload(filterTunnelMessages(parseTxMessages())), + encodeMspV1Reply(testLargeReplyMspCommand, MSP_RESULT_ACK, makeTestReplyPayload(300))); + EXPECT_FALSE(serialTxOverrun); +} + TEST(MavlinkTelemetryTest, TunnelRebootReplyIsFlushedBeforeRebootWhenTxRingIsFull) { initMavlinkTestState();