This commit is contained in:
Jacob Dahl
2025-07-24 14:09:40 -08:00
parent dce58637af
commit bbb84e25e6
7 changed files with 469 additions and 150 deletions
@@ -0,0 +1,178 @@
/****************************************************************************
*
* Copyright (c) 2025 PX4 Development Team. All rights reserved.
*
* Redistribution and use in source and binary forms, with or without
* modification, are permitted provided that the following conditions
* are met:
*
* 1. Redistributions of source code must retain the above copyright
* notice, this list of conditions and the following disclaimer.
* 2. Redistributions in binary form must reproduce the above copyright
* notice, this list of conditions and the following disclaimer in
* the documentation and/or other materials provided with the
* distribution.
* 3. Neither the name PX4 nor the names of its contributors may be
* used to endorse or promote products derived from this software
* without specific prior written permission.
*
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
* "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
* LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS
* FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE
* COPYRIGHT OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT,
* INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING,
* BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS
* OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED
* AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT
* LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN
* ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE
* POSSIBILITY OF SUCH DAMAGE.
*
****************************************************************************/
class TokenBucketRateLimiter
{
private:
uint64_t _last_token_update{0};
float _available_tokens{0.0f}; // Currently available tokens (in bytes)
float _max_bucket_tokens{0.0f}; // Maximum tokens bucket can hold
float _tokens_per_microsecond{0.0f}; // Token generation rate
uint32_t _configured_rate_bytes_per_sec{0}; // For debugging/telemetry
// Statistics
uint64_t _total_tokens_requested{0};
uint64_t _total_tokens_granted{0};
uint64_t _total_tokens_denied{0};
public:
/**
* Configure the rate limiter
* @param bytes_per_sec Maximum data rate in bytes per second
*/
void configure_rate(uint32_t bytes_per_sec)
{
_configured_rate_bytes_per_sec = bytes_per_sec;
_tokens_per_microsecond = static_cast<float>(bytes_per_sec) / 1e6f;
// TODO: how large of a burst should we allow? Probably MAVLINK_MAX_PACKET_LEN
// Bucket size: allow burst of 100ms worth of data
// This permits parameter sending bursts while maintaining average rate
// _max_bucket_tokens = static_cast<float>(bytes_per_sec) * 0.1f;
// Fill bucket with max bandwidth worth of tokens, duhhh. All our traffic is bursty -_-
_max_bucket_tokens = static_cast<float>(bytes_per_sec);
// Don't let available tokens exceed new limit
if (_available_tokens > _max_bucket_tokens) {
_available_tokens = _max_bucket_tokens;
}
}
/**
* Request tokens for transmission
* @param bytes Number of bytes to transmit
* @param priority Allow overdraft for high priority messages
* @return true if tokens were granted
*/
bool request_tokens(size_t bytes, bool priority = false)
{
replenish_tokens();
_total_tokens_requested += bytes;
const float tokens_needed = static_cast<float>(bytes);
if (priority && _available_tokens < tokens_needed) {
// Priority messages can overdraft up to 20% of bucket size
const float max_overdraft = _max_bucket_tokens * 0.2f;
if (_available_tokens >= -max_overdraft) {
_available_tokens -= tokens_needed;
_total_tokens_granted += bytes;
return true;
}
}
if (_available_tokens >= tokens_needed) {
_available_tokens -= tokens_needed;
_total_tokens_granted += bytes;
return true;
}
_total_tokens_denied += bytes;
return false;
}
bool check_tokens_available(size_t bytes) const
{
return _available_tokens >= static_cast<float>(bytes);
}
float get_available_tokens()
{
replenish_tokens();
return _available_tokens;
}
/**
* Get percentage of bucket filled
* @return 0.0 to 1.0 representing bucket fill level
*/
float get_bucket_fill_ratio()
{
replenish_tokens();
return _available_tokens / _max_bucket_tokens;
}
struct Statistics {
uint64_t tokens_requested;
uint64_t tokens_granted;
uint64_t tokens_denied;
float current_tokens;
float max_tokens;
float fill_ratio;
uint32_t configured_rate;
};
Statistics get_statistics()
{
replenish_tokens();
return {
_total_tokens_requested,
_total_tokens_granted,
_total_tokens_denied,
_available_tokens,
_max_bucket_tokens,
get_bucket_fill_ratio(),
_configured_rate_bytes_per_sec
};
}
private:
void replenish_tokens()
{
hrt_abstime timestamp = hrt_absolute_time();
// Handle first call
if (_last_token_update == 0) {
_last_token_update = timestamp;
_available_tokens = _max_bucket_tokens;
return;
}
const float dt_us = static_cast<float>(timestamp - _last_token_update);
_last_token_update = timestamp;
// Generate new tokens based on elapsed time
const float new_tokens = _tokens_per_microsecond * dt_us;
_available_tokens += new_tokens;
// Cap at bucket maximum
if (_available_tokens > _max_bucket_tokens) {
_available_tokens = _max_bucket_tokens;
}
}
};
+36 -2
View File
@@ -2169,7 +2169,8 @@ Mavlink::task_main(int argc, char *argv[])
/* USB serial is indicated by /dev/ttyACMx */
if (strncmp(_device_name, "/dev/ttyACM", 11) == 0) {
if (_datarate == 0) {
_datarate = 100000;
// _datarate = 100000;
_datarate = 30000;
}
/* USB has no baudrate, but use a magic number for 'fast' */
@@ -2209,6 +2210,7 @@ Mavlink::task_main(int argc, char *argv[])
fflush(stdout);
}
#if defined(MAVLINK_UDP)
else if (get_protocol() == Protocol::UDP) {
@@ -2223,6 +2225,9 @@ Mavlink::task_main(int argc, char *argv[])
#endif // MAVLINK_UDP
PX4_INFO("JAKE JAKE JAKE: DATARATE: %d", _datarate);
_tx_rate_limiter.configure_rate(_datarate);
if (set_instance_id()) {
if (!set_channel()) {
PX4_ERR("set channel failed");
@@ -2297,6 +2302,8 @@ Mavlink::task_main(int argc, char *argv[])
_main_loop_delay = MAVLINK_MAX_INTERVAL;
}
PX4_INFO("_main_loop_delay %u", _main_loop_delay);
/* open the UART device after setting the instance, as it might block */
if (get_protocol() == Protocol::SERIAL) {
@@ -2344,6 +2351,31 @@ Mavlink::task_main(int argc, char *argv[])
_task_running.store(true);
// Calculate required bandwidth for all our streams.
// - If exceeding configured _datarate (MAV_x_RATE), reduce all streams via update_rate_mult()
// - Always leave 20% headroom for other services
float total_bytes_per_s = 0;
for (auto stream : _streams) {
uint32_t bytes = stream->get_size();
uint32_t interval = stream->get_interval();
if (bytes == 0) {
continue;
}
PX4_INFO("Stream %s --> %lu bytes at %luus interval", stream->get_name(), bytes, interval);
float bytes_per_s = float(bytes) / (float(interval) / 1e6f);
PX4_INFO("bytes_per_s %f", (double)bytes_per_s);
total_bytes_per_s += bytes_per_s;
}
PX4_INFO("Total stream B/s --> %f", (double)total_bytes_per_s);
while (!should_exit()) {
/* main loop */
px4_usleep(_main_loop_delay);
@@ -2382,10 +2414,12 @@ Mavlink::task_main(int argc, char *argv[])
check_requested_subscriptions();
/* update streams */
// TODO: rate limit based on available bandwidth
// update streams
for (const auto &stream : _streams) {
stream->update(t);
// TODO: the below logic feels out of place
if (!_first_heartbeat_sent) {
if (_mode == MAVLINK_MODE_IRIDIUM) {
if (stream->get_id() == MAVLINK_MSG_ID_HIGH_LATENCY2) {
+41
View File
@@ -86,6 +86,8 @@
#include "mavlink_shell.h"
#include "mavlink_ulog.h"
#include "TokenBucketRateLimiter.h"
#define DEFAULT_BAUD_RATE 57600
#define DEFAULT_DEVICE_NAME "/dev/ttyS1"
@@ -525,6 +527,45 @@ public:
bool radio_status_critical() const { return _radio_status_critical; }
////////////////////////////////////////
// Token Bucket Rate Limiter
////////////////////////////////////////
bool request_data_tokens(size_t bytes, bool priority = false)
{
return _tx_rate_limiter.request_tokens(bytes, priority);
}
/**
* Check if tokens available without consuming
* @param bytes Number of bytes to check
* @return true if tokens available
*/
bool check_data_tokens(size_t bytes) const
{
return _tx_rate_limiter.check_tokens_available(bytes);
}
/**
* Get current token availability
* @return Available tokens in bytes
*/
float get_available_data_tokens()
{
return _tx_rate_limiter.get_available_tokens();
}
/**
* Update rate limiter when configuration changes
*/
void update_data_rate_limits()
{
// Use the configured data rate
_tx_rate_limiter.configure_rate(_datarate);
}
TokenBucketRateLimiter _tx_rate_limiter;
private:
MavlinkReceiver _receiver;
+188 -129
View File
@@ -66,14 +66,16 @@ MavlinkParametersManager::handle_message(const mavlink_message_t *msg)
mavlink_param_request_list_t req_list;
mavlink_msg_param_request_list_decode(msg, &req_list);
PX4_INFO("PARAM_REQUEST_LIST");
if (req_list.target_system == mavlink_system.sysid &&
(req_list.target_component == mavlink_system.compid || req_list.target_component == MAV_COMP_ID_ALL)) {
if (_send_all_index < 0) {
_send_all_index = PARAM_HASH;
if (_next_param_index < 0) {
_next_param_index = PARAM_HASH;
} else {
/* a restart should skip the hash check on the ground */
_send_all_index = 0;
_next_param_index = 0;
}
}
@@ -112,7 +114,7 @@ MavlinkParametersManager::handle_message(const mavlink_message_t *msg)
if (strncmp(name, "_HASH_CHECK", sizeof(name)) == 0) {
if (_mavlink.hash_check_enabled()) {
_send_all_index = -1;
_next_param_index = -1;
}
/* No other action taken, return */
@@ -284,10 +286,59 @@ MavlinkParametersManager::handle_message(const mavlink_message_t *msg)
}
}
// void
// MavlinkParametersManager::send()
// {
// if (_first_send) {
// // parameters QGC can't tolerate not finding (2020-11-11)
// param_find("BAT_CRIT_THR");
// param_find("BAT_EMERGEN_THR");
// param_find("BAT_LOW_THR");
// param_find("CAL_ACC0_ID");
// param_find("CAL_GYRO0_ID");
// param_find("CAL_MAG0_ID");
// param_find("CAL_MAG0_ROT");
// param_find("CAL_MAG1_ID");
// param_find("CAL_MAG1_ROT");
// param_find("CAL_MAG2_ID");
// param_find("CAL_MAG2_ROT");
// param_find("CAL_MAG3_ID");
// param_find("CAL_MAG3_ROT");
// param_find("SENS_BOARD_ROT");
// param_find("SENS_BOARD_X_OFF");
// param_find("SENS_BOARD_Y_OFF");
// param_find("SENS_BOARD_Z_OFF");
// param_find("SENS_DPRES_OFF");
// param_find("TRIG_MODE");
// param_find("UAVCAN_ENABLE");
// // parameter only used in startup script but should show on ground station
// param_find("SYS_PARAM_VER");
// _first_send = false;
// }
// int max_num_to_send;
// if (_mavlink.get_protocol() == Protocol::SERIAL && !_mavlink.is_usb_uart()) {
// max_num_to_send = 3;
// } else {
// // speed up parameter loading via UDP or USB: try to send 20 at once
// max_num_to_send = 20;
// }
// int i = 0;
// // Send while burst is not exceeded, we still have buffer space and still something to send
// while ((i++ < max_num_to_send) && (_mavlink.get_free_tx_buf() >= get_size()) && !_mavlink.radio_status_critical()
// && send_params()) {}
// }
void
MavlinkParametersManager::send()
{
if (!_first_send) {
if (_first_send) {
// parameters QGC can't tolerate not finding (2020-11-11)
param_find("BAT_CRIT_THR");
param_find("BAT_EMERGEN_THR");
@@ -313,179 +364,185 @@ MavlinkParametersManager::send()
// parameter only used in startup script but should show on ground station
param_find("SYS_PARAM_VER");
_first_send = true;
_first_send = false;
}
int max_num_to_send;
// Calculate how many parameters we could send based on bandwidth
const size_t bytes_per_param = get_size();
const float available_tokens = _mavlink.get_available_data_tokens();
const int max_by_bandwidth = static_cast<int>(available_tokens / bytes_per_param);
if (_mavlink.get_protocol() == Protocol::SERIAL && !_mavlink.is_usb_uart()) {
max_num_to_send = 3;
} else {
// speed up parameter loading via UDP or USB: try to send 20 at once
max_num_to_send = 20;
if (max_by_bandwidth <= 0) {
// PX4_INFO("No tokens available: %f", (double)available_tokens);
return;
}
int i = 0;
// TODO: determine how much of the available bandwidth we want to consume
int send_count = math::max(max_by_bandwidth / 2, 1);
// Send while burst is not exceeded, we still have buffer space and still something to send
while ((i++ < max_num_to_send) && (_mavlink.get_free_tx_buf() >= get_size()) && !_mavlink.radio_status_critical()
&& send_params()) {}
}
while (send_count) {
bool
MavlinkParametersManager::send_params()
{
#if defined(CONFIG_MAVLINK_UAVCAN_PARAMETERS)
if (send_uavcan()) {
return true;
}
bool success = _mavlink.request_data_tokens(bytes_per_param);
#endif // CONFIG_MAVLINK_UAVCAN_PARAMETERS
if (!success) {
PX4_INFO("Tokens unavailable");
break;
}
if (send_one()) {
return true;
if (send_uavcan()) {
send_count--;
}
#endif
} else if (send_untransmitted()) {
return true;
}
return false;
}
bool
MavlinkParametersManager::send_untransmitted()
{
bool sent_one = false;
if (_parameter_update_sub.updated()) {
// clear the update
parameter_update_s pupdate;
_parameter_update_sub.copy(&pupdate);
// Schedule an update if not already the case
if (_param_update_time == 0) {
_param_update_time = pupdate.timestamp;
_param_update_index = 0;
if (send_one()) {
send_count--;
} else {
// Finished
break;
}
}
if ((_param_update_time != 0) && ((_param_update_time + 5 * 1000) < hrt_absolute_time())) {
param_t param = 0;
// send out all changed values
do {
// skip over all parameters which are not invalid and not used
do {
param = param_for_index(_param_update_index);
++_param_update_index;
} while (param != PARAM_INVALID && !param_used(param));
// send parameters which are untransmitted while there is
// space in the TX buffer
if ((param != PARAM_INVALID) && param_value_unsaved(param)) {
int ret = send_param(param);
sent_one = true;
if (ret != PX4_OK) {
break;
}
}
} while ((_mavlink.get_free_tx_buf() >= get_size()) && !_mavlink.radio_status_critical()
&& (_param_update_index < (int) param_count()));
// Flag work as done once all params have been sent
if (_param_update_index >= (int) param_count()) {
_param_update_time = 0;
}
}
return sent_one;
}
// bool
// MavlinkParametersManager::send_untransmitted()
// {
// bool sent_one = false;
// if (_parameter_update_sub.updated()) {
// // clear the update
// parameter_update_s pupdate;
// _parameter_update_sub.copy(&pupdate);
// // Schedule an update if not already the case
// if (_param_update_time == 0) {
// _param_update_time = pupdate.timestamp;
// _param_update_index = 0;
// }
// }
// if ((_param_update_time != 0) && ((_param_update_time + 5 * 1000) < hrt_absolute_time())) {
// param_t param = 0;
// // send out all changed values
// do {
// // skip over all parameters which are not invalid and not used
// do {
// param = param_for_index(_param_update_index);
// ++_param_update_index;
// } while (param != PARAM_INVALID && !param_used(param));
// // send parameters which are untransmitted while there is
// // space in the TX buffer
// if ((param != PARAM_INVALID) && param_value_unsaved(param)) {
// int ret = send_param(param);
// sent_one = true;
// if (ret != PX4_OK) {
// break;
// }
// }
// } while ((_mavlink.get_free_tx_buf() >= get_size()) && !_mavlink.radio_status_critical()
// && (_param_update_index < (int) param_count()));
// // Flag work as done once all params have been sent
// if (_param_update_index >= (int) param_count()) {
// _param_update_time = 0;
// }
// }
// return sent_one;
// }
bool
MavlinkParametersManager::send_one()
{
const hrt_abstime now = hrt_absolute_time();
// If in low-bandwidth mode, throttle parameter transmission to 8 Hz
if (_mavlink.get_mode() == Mavlink::MAVLINK_MODE_LOW_BANDWIDTH
&& now < _last_param_sent_timestamp + 125_ms) {
// Finished
if (_next_param_index < 0) {
return false;
}
if (_send_all_index >= 0) {
/* send all parameters if requested, but only after the system has booted */
// The first thing we send is a hash of all values for the ground
// station to try and quickly load a cached copy of our params
if (_next_param_index == PARAM_HASH) {
send_param_hash();
return true;
}
/* The first thing we send is a hash of all values for the ground
* station to try and quickly load a cached copy of our params
*/
if (_send_all_index == PARAM_HASH) {
/* return hash check for cached params */
uint32_t hash = param_hash_check();
// Iterate over all parameter indices
param_t p = PARAM_INVALID;
/* build the one-off response message */
mavlink_param_value_t msg;
msg.param_count = param_count_used();
msg.param_index = -1;
strncpy(msg.param_id, HASH_PARAM, MAVLINK_MSG_PARAM_VALUE_FIELD_PARAM_ID_LEN);
msg.param_type = MAV_PARAM_TYPE_UINT32;
memcpy(&msg.param_value, &hash, sizeof(hash));
mavlink_msg_param_value_send_struct(_mavlink.get_channel(), &msg);
while (1) {
p = param_for_index(_next_param_index);
/* after this we should start sending all params */
_send_all_index = 0;
/* No further action, return now */
return true;
if (p == PARAM_INVALID || !param_used(p)) {
// There can be a lot of invalid or unused parameters, we skip those
_next_param_index++;
continue;
}
/* look for the first parameter which is used */
param_t p;
// Param is valid and used, send it
// TODO: use result to decide retries
auto result = send_param(p);
do {
/* walk through all parameters, including unused ones */
p = param_for_index(_send_all_index);
_send_all_index++;
} while (p != PARAM_INVALID && !param_used(p));
if (p != PARAM_INVALID) {
send_param(p);
if (result == 0) {
_next_param_index++;
_last_param_sent_timestamp = now;
}
if ((p == PARAM_INVALID) || (_send_all_index >= (int) param_count())) {
_send_all_index = -1;
return false;
} else {
return true;
}
// break from the loop after a send attempt
break;
}
return false;
if (_next_param_index >= (int) param_count()) {
// Finished sending all params
_next_param_index = -1;
return false;
}
return true;
}
int
MavlinkParametersManager::send_param_hash()
{
uint32_t hash = param_hash_check();
mavlink_param_value_t msg;
msg.param_count = param_count_used();
msg.param_index = -1;
strncpy(msg.param_id, HASH_PARAM, MAVLINK_MSG_PARAM_VALUE_FIELD_PARAM_ID_LEN);
msg.param_type = MAV_PARAM_TYPE_UINT32;
memcpy(&msg.param_value, &hash, sizeof(hash));
mavlink_msg_param_value_send_struct(_mavlink.get_channel(), &msg);
// after this we should start sending all params
_next_param_index = 0;
return true;
}
int
MavlinkParametersManager::send_param(param_t param, int component_id)
{
if (param == PARAM_INVALID) {
PX4_INFO("wtf redundant");
return 1;
}
/* no free TX buf to send this param */
// no free TX buf to send this param
if (_mavlink.get_free_tx_buf() < MAVLINK_MSG_ID_PARAM_VALUE_LEN) {
PX4_INFO("no free tx");
return 1;
}
mavlink_param_value_t msg;
mavlink_param_value_t msg = {};
/*
* get param value, since MAVLink encodes float and int params in the same
* space during transmission, copy param onto float val_buf
*/
if (param_type(param) == PARAM_TYPE_INT32) {
int32_t param_value;
@@ -505,6 +562,8 @@ MavlinkParametersManager::send_param(param_t param, int component_id)
msg.param_value = param_value;
}
// TODO: both operations below iterate over the entire parameter list. This is both
// expensive and redundant during a PARAM_REQUEST_LIST
msg.param_count = param_count_used();
msg.param_index = param_get_used_index(param);
+6 -12
View File
@@ -78,7 +78,7 @@ public:
void handle_message(const mavlink_message_t *msg);
private:
int _send_all_index{-1};
int _next_param_index{-1};
/* do not allow top copying this class */
MavlinkParametersManager(MavlinkParametersManager &);
@@ -89,17 +89,11 @@ protected:
/// @return true if a parameter was sent
bool send_one();
/**
* Handle any open param send transfer
*/
bool send_params();
/**
* Send untransmitted params
*/
bool send_untransmitted();
// TODO: what is the return value?
int send_param(param_t param, int component_id = -1);
// TODO: what is the return value?
int send_param_hash();
#if defined(CONFIG_MAVLINK_UAVCAN_PARAMETERS)
/**
@@ -161,6 +155,6 @@ protected:
Mavlink &_mavlink;
bool _first_send{false};
bool _first_send{true};
hrt_abstime _last_param_sent_timestamp{0}; // time at which the last parameter was sent
};
+14 -6
View File
@@ -3278,11 +3278,19 @@ MavlinkReceiver::run()
usleep(10000);
}
const hrt_abstime t = hrt_absolute_time();
const hrt_abstime now = hrt_absolute_time();
CheckHeartbeats(t);
CheckHeartbeats(now);
if (t - last_send_update > timeout * 1000) {
// FIXME:
// Mavlink receiver is emitting messages here. That's not what a "receiver" does.
// - mission items
// - params
// - ftp
// TODO: why limit at all?
// NOTE: limited to 10ms update interval
if (now - last_send_update > timeout * 1000) {
_mission_manager.check_active_mission();
_mission_manager.send();
@@ -3295,13 +3303,13 @@ MavlinkReceiver::run()
}
_mavlink_log_handler.send();
last_send_update = t;
last_send_update = now;
}
if (_tune_publisher != nullptr) {
_tune_publisher->publish_next_tune(t);
_tune_publisher->publish_next_tune(now);
}
}
} // end while-loop --> !_mavlink.should_exit()
}
bool MavlinkReceiver::component_was_seen(int system_id, int component_id)
+6 -1
View File
@@ -104,10 +104,15 @@ MavlinkStream::update(const hrt_abstime &t)
// needs to be accounted for as well.
// This method is not theoretically optimal but a suitable
// stopgap as it hits its deadlines well (0.5 Hz, 50 Hz and 250 Hz)
if (unlimited_rate || (dt > (interval - (_mavlink->get_main_loop_delay() / 10) * 3))) {
// interval expired, send message
// Only send if tokens are available
if (!_mavlink->request_data_tokens(get_size())) {
// TODO: per stream tokens_denied counter
return -1;
}
// If the interval is non-zero and dt is smaller than 1.5 times the interval
// do not use the actual time but increment at a fixed rate, so that processing delays do not
// distort the average rate. The check of the maximum interval is done to ensure that after a