From bbb84e25e6d548812e1865263b04912d803306cb Mon Sep 17 00:00:00 2001 From: Jacob Dahl Date: Wed, 23 Jul 2025 17:18:46 -0800 Subject: [PATCH] staging --- src/modules/mavlink/TokenBucketRateLimiter.h | 178 +++++++++++ src/modules/mavlink/mavlink_main.cpp | 38 ++- src/modules/mavlink/mavlink_main.h | 41 +++ src/modules/mavlink/mavlink_parameters.cpp | 317 +++++++++++-------- src/modules/mavlink/mavlink_parameters.h | 18 +- src/modules/mavlink/mavlink_receiver.cpp | 20 +- src/modules/mavlink/mavlink_stream.cpp | 7 +- 7 files changed, 469 insertions(+), 150 deletions(-) create mode 100644 src/modules/mavlink/TokenBucketRateLimiter.h diff --git a/src/modules/mavlink/TokenBucketRateLimiter.h b/src/modules/mavlink/TokenBucketRateLimiter.h new file mode 100644 index 0000000000..f21c556471 --- /dev/null +++ b/src/modules/mavlink/TokenBucketRateLimiter.h @@ -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(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(bytes_per_sec) * 0.1f; + + + // Fill bucket with max bandwidth worth of tokens, duhhh. All our traffic is bursty -_- + _max_bucket_tokens = static_cast(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(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(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(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; + } + } +}; diff --git a/src/modules/mavlink/mavlink_main.cpp b/src/modules/mavlink/mavlink_main.cpp index b86c9a6eca..1fe0ba684b 100644 --- a/src/modules/mavlink/mavlink_main.cpp +++ b/src/modules/mavlink/mavlink_main.cpp @@ -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) { diff --git a/src/modules/mavlink/mavlink_main.h b/src/modules/mavlink/mavlink_main.h index 304bd06d58..9e8de3438d 100644 --- a/src/modules/mavlink/mavlink_main.h +++ b/src/modules/mavlink/mavlink_main.h @@ -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; diff --git a/src/modules/mavlink/mavlink_parameters.cpp b/src/modules/mavlink/mavlink_parameters.cpp index ab341d0186..7fe5a16a67 100644 --- a/src/modules/mavlink/mavlink_parameters.cpp +++ b/src/modules/mavlink/mavlink_parameters.cpp @@ -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(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); diff --git a/src/modules/mavlink/mavlink_parameters.h b/src/modules/mavlink/mavlink_parameters.h index 6ad09889c2..395bd2245d 100644 --- a/src/modules/mavlink/mavlink_parameters.h +++ b/src/modules/mavlink/mavlink_parameters.h @@ -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 }; diff --git a/src/modules/mavlink/mavlink_receiver.cpp b/src/modules/mavlink/mavlink_receiver.cpp index 9c8e0099b0..62d2ffc4d9 100644 --- a/src/modules/mavlink/mavlink_receiver.cpp +++ b/src/modules/mavlink/mavlink_receiver.cpp @@ -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) diff --git a/src/modules/mavlink/mavlink_stream.cpp b/src/modules/mavlink/mavlink_stream.cpp index 6b00fd4e49..e4d697d9d6 100644 --- a/src/modules/mavlink/mavlink_stream.cpp +++ b/src/modules/mavlink/mavlink_stream.cpp @@ -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