Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/Mavlink.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
93 changes: 69 additions & 24 deletions src/main/fc/fc_mavlink.c
Original file line number Diff line number Diff line change
Expand Up @@ -20,40 +20,72 @@ 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;
}

if (pending->portIndex != portIndex) {
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));
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;
}

mavlink_msg_tunnel_pack(
mavSystemId,
mavComponentId,
&mavSendMsg,
pending->targetSystem,
pending->targetComponent,
MAVLINK_TUNNEL_PAYLOAD_TYPE_INAV_MSP,
chunkLength,
chunk);
if (!mavlinkSendMessageToPortIfRoom(portIndex)) {
return;
}
mavlinkSendTunnelReply(targetSystem, targetComponent, chunk, chunkLength);

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;
}

Expand Down Expand Up @@ -94,11 +126,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 (mavlinkTunnelMspReplyIsPending()) {
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,
Expand All @@ -120,6 +159,7 @@ static bool mavlinkProcessCompletedTunnelCommand(uint8_t ingressPortIndex)

if (status != MSP_RESULT_NO_REPLY) {
mavlinkSendTunnelMspReply(
ingressPortIndex,
mavlinkContext.recvMsg.sysid,
mavlinkContext.recvMsg.compid,
&reply,
Expand All @@ -132,6 +172,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;
}
Expand Down
3 changes: 3 additions & 0 deletions src/main/fc/fc_mavlink.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
11 changes: 11 additions & 0 deletions src/main/mavlink/mavlink_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand Down Expand Up @@ -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)
Expand Down
3 changes: 3 additions & 0 deletions src/main/mavlink/mavlink_ports.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
80 changes: 56 additions & 24 deletions src/main/mavlink/mavlink_runtime.c
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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;
}
Expand All @@ -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];
Expand Down Expand Up @@ -266,15 +292,16 @@ 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];

return state->portConfig &&
(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)
Expand All @@ -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;
}

Expand Down
1 change: 1 addition & 0 deletions src/main/mavlink/mavlink_runtime.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Loading
Loading