From 2d22d95cac9d1cbcf523ac8aabfd9301e52bd95b Mon Sep 17 00:00:00 2001 From: Dhillon Kannabhiran Date: Sat, 19 Sep 2026 12:46:51 +0800 Subject: [PATCH] Add a rolling-window relay airtime budget with forward accounting Nodes that repeat traffic for other nodes have no visibility into how much airtime the relay role consumes, and no way to bound it separately from their own traffic. The existing token bucket (set dutycycle) bounds total airtime but is shared by both, accumulates credit that can be drawn down in a burst, and reports only one total. Add mesh::AirtimeBudget, a rolling-window airtime ledger: airtime is accumulated into slots and a slot is only forgotten once it lies entirely outside the window, so the reported usage never understates the true trailing-window total. The dispatcher marks the packets it re-transmits, admits a forward only when the estimated airtime fits the relay budget, counts refused forwards, and records the measured airtime of completed forwards. The relay duty cycle is configurable per node (set fwd.dutycycle / pref fwd_af); with the default factor of 0 the ledger reports usage and refuses nothing, so existing behaviour is unchanged. stats-radio now reports fwd_air_secs, fwd_window_secs/fwd_limit_secs and fwd_dropped. New unit tests cover the ledger (accumulation, limits, conservative expiry, millis() wrap) and the dispatcher's forwarding admission and accounting. --- docs/cli_commands.md | 24 ++ examples/companion_radio/MyMesh.cpp | 5 + examples/companion_radio/MyMesh.h | 1 + examples/companion_radio/NodePrefs.h | 4 + examples/simple_repeater/MyMesh.cpp | 5 +- examples/simple_repeater/MyMesh.h | 4 + examples/simple_room_server/MyMesh.cpp | 3 +- examples/simple_sensor/SensorMesh.cpp | 3 +- platformio.ini | 2 + src/AirtimeBudget.cpp | 68 ++++ src/AirtimeBudget.h | 75 +++++ src/Dispatcher.cpp | 27 +- src/Dispatcher.h | 13 + src/Packet.cpp | 3 +- src/Packet.h | 1 + src/helpers/CommonCLI.h | 4 + src/helpers/CommonRadioPrefs.cpp | 21 ++ src/helpers/CommonRadioPrefs.h | 3 + src/helpers/StatsFormatHelper.h | 14 +- .../test_airtime_budget.cpp | 293 ++++++++++++++++++ 20 files changed, 564 insertions(+), 9 deletions(-) create mode 100644 src/AirtimeBudget.cpp create mode 100644 src/AirtimeBudget.h create mode 100644 test/test_airtime_budget/test_airtime_budget.cpp diff --git a/docs/cli_commands.md b/docs/cli_commands.md index f4a43fad3c..acd1252330 100644 --- a/docs/cli_commands.md +++ b/docs/cli_commands.md @@ -577,6 +577,30 @@ This document provides an overview of CLI commands that can be sent to MeshCore --- +#### View or change the relay (forward) duty cycle limit +**Usage:** +- `get fwd.dutycycle` +- `set fwd.dutycycle ` + +**Parameters:** +- `value`: Duty cycle percentage (1-100) that this node may spend re-transmitting other + nodes' packets. Airtime spent forwarding is accounted over a rolling window, and once + the budget is spent further forwards are dropped (and counted) rather than delayed. + Traffic this node originates, including its own ACKs, is not charged to this budget. + +**Default:** `100%` (no separate limit; forwards are still bounded by the overall duty cycle) + +**Examples:** +- `set fwd.dutycycle 100` — no separate relay limit +- `set fwd.dutycycle 10` — relay at most 10% of each window +- `set fwd.dutycycle 1` — relay at most 1% of each window (strictest EU requirement) + +> **Note:** Takes effect after reboot. `stats-radio` reports relay usage as +> `fwd_air_secs` (total), `fwd_window_secs` / `fwd_limit_secs` (current window and cap) +> and `fwd_dropped` (forwards refused). + +--- + #### View or change the airtime factor (duty cycle limit) > **Deprecated** as of firmware v1.15.0. Use [`get/set dutycycle`](#view-or-change-the-duty-cycle-limit) instead. diff --git a/examples/companion_radio/MyMesh.cpp b/examples/companion_radio/MyMesh.cpp index c38fa70ae8..25fa8b0f16 100644 --- a/examples/companion_radio/MyMesh.cpp +++ b/examples/companion_radio/MyMesh.cpp @@ -267,6 +267,10 @@ float MyMesh::getAirtimeBudgetFactor() const { return _prefs.airtime_factor; } +float MyMesh::getForwardAirtimeBudgetFactor() const { + return _prefs.fwd_airtime_factor; +} + int MyMesh::getInterferenceThreshold() const { return _prefs.interference_threshold; } @@ -1019,6 +1023,7 @@ void MyMesh::begin(bool has_display) { _prefs.tx_delay_factor = constrain(_prefs.tx_delay_factor, 0, 2.0f); _prefs.direct_tx_delay_factor = constrain(_prefs.direct_tx_delay_factor, 0, 2.0f); _prefs.airtime_factor = constrain(_prefs.airtime_factor, 0, 9.0f); + _prefs.fwd_airtime_factor = constrain(_prefs.fwd_airtime_factor, 0, 9.0f); _prefs.freq = constrain(_prefs.freq, 150.0f, 2500.0f); _prefs.bw = constrain(_prefs.bw, 7.8f, 500.0f); _prefs.sf = constrain(_prefs.sf, 5, 12); diff --git a/examples/companion_radio/MyMesh.h b/examples/companion_radio/MyMesh.h index 3b98a4f674..83aeba9467 100644 --- a/examples/companion_radio/MyMesh.h +++ b/examples/companion_radio/MyMesh.h @@ -119,6 +119,7 @@ class MyMesh : public BaseChatMesh, public DataStoreHost { protected: float getAirtimeBudgetFactor() const override; + float getForwardAirtimeBudgetFactor() const override; int getInterferenceThreshold() const override; bool getCADEnabled() const override; int getAGCResetInterval() const override { diff --git a/examples/companion_radio/NodePrefs.h b/examples/companion_radio/NodePrefs.h index 4b169d73fa..aa422c82e2 100644 --- a/examples/companion_radio/NodePrefs.h +++ b/examples/companion_radio/NodePrefs.h @@ -14,6 +14,7 @@ class NodePrefs : public ConfigSerializer { // persisted to file public: float airtime_factor = 0; + float fwd_airtime_factor = 0; // separate duty cycle limit for relayed traffic (0 = no separate limit) char node_name[32]; double node_lat = 0, node_lon = 0; float freq = 0; @@ -75,6 +76,7 @@ class NodePrefs : public ConfigSerializer { // persisted to file def("fem_txgain", _parent->radio_fem_txgain); def("tx", _parent->tx_power_dbm); def("af", _parent->airtime_factor); + def("fwd_af", _parent->fwd_airtime_factor); def("rxdelay", _parent->rx_delay_base); def("f_txdelay", _parent->tx_delay_factor); def("d_txdelay", _parent->direct_tx_delay_factor); @@ -96,6 +98,8 @@ class NodePrefs : public ConfigSerializer { // persisted to file void setCodingRate(uint8_t cr) override { _parent->cr = cr; markDirty(); } float getAirtimeFactor() const override { return _parent->airtime_factor; } void setAirtimeFactor(float af) override { _parent->airtime_factor = af; markDirty(); } + float getForwardAirtimeFactor() const override { return _parent->fwd_airtime_factor; } + void setForwardAirtimeFactor(float af) override { _parent->fwd_airtime_factor = af; markDirty(); } bool isCadEnabled() const override { return _parent->cad_enabled; } void setCadEnabled(bool en) override { _parent->cad_enabled = en; markDirty(); } uint8_t getIntThresh() const override { return _parent->interference_threshold; } diff --git a/examples/simple_repeater/MyMesh.cpp b/examples/simple_repeater/MyMesh.cpp index ca6a3e607e..4172d0190e 100644 --- a/examples/simple_repeater/MyMesh.cpp +++ b/examples/simple_repeater/MyMesh.cpp @@ -202,7 +202,7 @@ uint8_t MyMesh::handleAnonClockReq(const mesh::Identity& sender, uint32_t sender if (_prefs.disable_fwd) { // is this repeater currently disabled reply_data[8] |= 0x80; // is disabled } - // TODO: add some kind of moving-window utilisation metric, so can query 'how busy' is this repeater + // relay utilisation (moving-window airtime spent forwarding) is reported by the 'stats-radio' CLI command return 9; // reply length } return 0; @@ -1167,7 +1167,8 @@ void MyMesh::formatStatsReply(char *reply) { } void MyMesh::formatRadioStatsReply(char *reply) { - StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime()); + StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime(), + getForwardAirTime(), getForwardBudgetUsed(), getForwardBudgetLimit(), getNumForwardDropped()); } void MyMesh::formatPacketStatsReply(char *reply) { diff --git a/examples/simple_repeater/MyMesh.h b/examples/simple_repeater/MyMesh.h index cac6c4a281..e97648839f 100644 --- a/examples/simple_repeater/MyMesh.h +++ b/examples/simple_repeater/MyMesh.h @@ -132,6 +132,10 @@ class MyMesh : public mesh::Mesh, public CommonCLICallbacks { return _prefs.airtime_factor; } + float getForwardAirtimeBudgetFactor() const override { + return _prefs.fwd_airtime_factor; + } + bool allowPacketForward(const mesh::Packet* packet) override; const char* getLogDateTime() override; void logRxRaw(float snr, float rssi, const uint8_t raw[], int len) override; diff --git a/examples/simple_room_server/MyMesh.cpp b/examples/simple_room_server/MyMesh.cpp index 71b8d32a3a..21c874343a 100644 --- a/examples/simple_room_server/MyMesh.cpp +++ b/examples/simple_room_server/MyMesh.cpp @@ -888,7 +888,8 @@ void MyMesh::formatStatsReply(char *reply) { } void MyMesh::formatRadioStatsReply(char *reply) { - StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime()); + StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime(), + getForwardAirTime(), getForwardBudgetUsed(), getForwardBudgetLimit(), getNumForwardDropped()); } void MyMesh::formatPacketStatsReply(char *reply) { diff --git a/examples/simple_sensor/SensorMesh.cpp b/examples/simple_sensor/SensorMesh.cpp index 23d0cdc353..41aadde8bb 100644 --- a/examples/simple_sensor/SensorMesh.cpp +++ b/examples/simple_sensor/SensorMesh.cpp @@ -857,7 +857,8 @@ void SensorMesh::formatStatsReply(char *reply) { } void SensorMesh::formatRadioStatsReply(char *reply) { - StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime()); + StatsFormatHelper::formatRadioStats(reply, _radio, radio_driver, getTotalAirTime(), getReceiveAirTime(), + getForwardAirTime(), getForwardBudgetUsed(), getForwardBudgetLimit(), getNumForwardDropped()); } void SensorMesh::formatPacketStatsReply(char *reply) { diff --git a/platformio.ini b/platformio.ini index de4d6c29c3..888e0a624b 100644 --- a/platformio.ini +++ b/platformio.ini @@ -168,6 +168,8 @@ test_build_src = yes test_ignore = test_kiss_modem build_src_filter = -<*> + +<../src/AirtimeBudget.cpp> + +<../src/Dispatcher.cpp> +<../src/Utils.cpp> +<../src/Packet.cpp> +<../src/helpers/ConfigSerializer.cpp> diff --git a/src/AirtimeBudget.cpp b/src/AirtimeBudget.cpp new file mode 100644 index 0000000000..4164aafc66 --- /dev/null +++ b/src/AirtimeBudget.cpp @@ -0,0 +1,68 @@ +#include "AirtimeBudget.h" + +namespace mesh { + +void AirtimeBudget::clear() { + _window_ms = 0; + _slot_ms = 0; + _limit_ms = 0; + _used_ms = 0; + _slot_start = 0; + _slot_count = 0; + _current = 0; + for (uint8_t i = 0; i < MAX_SLOTS; i++) { + _slots[i] = 0; + } +} + +void AirtimeBudget::begin(uint32_t window_ms, uint32_t limit_ms, uint32_t now, uint8_t slot_count) { + clear(); + if (window_ms == 0 || slot_count == 0) return; // ledger disabled + if (slot_count > MAX_SLOTS) { + slot_count = MAX_SLOTS; + } + + _window_ms = window_ms; + _limit_ms = limit_ms; + _slot_count = slot_count; + _slot_ms = (window_ms + slot_count - 1) / slot_count; // round up, so the slots cover the whole window + _slot_start = now; +} + +void AirtimeBudget::update(uint32_t now) { + if (_slot_count == 0) return; + + uint32_t elapsed = now - _slot_start; // wraps correctly with millis() + if (elapsed >= _window_ms) { // everything is stale + _used_ms = 0; + _current = 0; + for (uint8_t i = 0; i < _slot_count; i++) { + _slots[i] = 0; + } + _slot_start = now; + return; + } + + uint32_t steps = elapsed / _slot_ms; + for (uint32_t s = 0; s < steps; s++) { + _current = (_current + 1) % _slot_count; + _used_ms -= _slots[_current]; + _slots[_current] = 0; + } + _slot_start += steps * _slot_ms; +} + +bool AirtimeBudget::canSpend(uint32_t airtime_ms, uint32_t now) { + if (!isEnabled()) return true; + update(now); + return _used_ms + airtime_ms <= _limit_ms; +} + +void AirtimeBudget::record(uint32_t airtime_ms, uint32_t now) { + if (_slot_count == 0 || airtime_ms == 0) return; + update(now); + _slots[_current] += airtime_ms; + _used_ms += airtime_ms; +} + +} diff --git a/src/AirtimeBudget.h b/src/AirtimeBudget.h new file mode 100644 index 0000000000..772a9070ec --- /dev/null +++ b/src/AirtimeBudget.h @@ -0,0 +1,75 @@ +#pragma once + +#include +#include + +namespace mesh { + +/** + * \brief Rolling-window airtime ledger. + * + * Accounts for airtime (milliseconds) spent transmitting, and reports how much + * of the trailing window has been used, so a caller can bound how much of the + * channel it occupies over a moving window (eg. a 1% duty-cycle band). + * + * Airtime is accumulated into fixed-size slots. A slot is only forgotten once + * it lies entirely outside the window, so the reported usage is never lower + * than the true trailing-window total: the ledger errs on the side of caution. + * More slots give a finer (less conservative) figure at the cost of RAM. + */ +class AirtimeBudget { +public: + static const uint8_t MAX_SLOTS = 12; + +private: + uint32_t _window_ms; + uint32_t _slot_ms; + uint32_t _limit_ms; + uint32_t _used_ms; + uint32_t _slot_start; + uint32_t _slots[MAX_SLOTS]; + uint8_t _slot_count; + uint8_t _current; + + void clear(); + +public: + AirtimeBudget() { clear(); } + + /** + * \brief Drop any usage that has fallen outside the window. + * \param now current clock, in milliseconds + */ + void update(uint32_t now); + + /** + * \brief (Re)initialise the ledger, eg. when prefs change. + * \param window_ms length of the rolling window + * \param limit_ms maximum airtime allowed per window (0 = no limit) + * \param now current clock, in milliseconds + * \param slot_count slots to divide the window into (1..MAX_SLOTS) + */ + void begin(uint32_t window_ms, uint32_t limit_ms, uint32_t now, uint8_t slot_count); + + /** + * \returns true if airtime_ms of transmit airtime still fits in the budget. + */ + bool canSpend(uint32_t airtime_ms, uint32_t now); + + /** + * \brief Record airtime that was actually spent transmitting. + */ + void record(uint32_t airtime_ms, uint32_t now); + + void setLimit(uint32_t limit_ms) { _limit_ms = limit_ms; } + + uint32_t getWindow() const { return _window_ms; } + uint32_t getLimit() const { return _limit_ms; } + uint32_t getUsed() const { return _used_ms; } + uint32_t getRemaining() const { return (!isEnabled() || _used_ms >= _limit_ms) ? 0 : _limit_ms - _used_ms; } + + /// A budget only binds when it is set below a full window of airtime. + bool isEnabled() const { return _slot_count > 0 && _limit_ms > 0 && _limit_ms < _window_ms; } +}; + +} diff --git a/src/Dispatcher.cpp b/src/Dispatcher.cpp index c0610b7f8a..9403836a2f 100644 --- a/src/Dispatcher.cpp +++ b/src/Dispatcher.cpp @@ -27,6 +27,13 @@ void Dispatcher::begin() { tx_budget_ms = (unsigned long)(duty_cycle_window_ms * duty_cycle); last_budget_update = _ms->getMillis(); + float fwd_factor = getForwardAirtimeBudgetFactor(); + if (fwd_factor < 0.0f) { + fwd_factor = 0.0f; // treat as no limit + } + float fwd_duty_cycle = 1.0f / (1.0f + fwd_factor); + _fwd_budget.begin(getDutyCycleWindowMs(), (uint32_t)(duty_cycle_window_ms * fwd_duty_cycle), _ms->getMillis(), AirtimeBudget::MAX_SLOTS); + _radio->begin(); prev_isrecv_mode = _radio->isInRecvMode(); } @@ -35,6 +42,10 @@ float Dispatcher::getAirtimeBudgetFactor() const { return 1.0; } +float Dispatcher::getForwardAirtimeBudgetFactor() const { + return 0; // by default, no separate budget for forwarded traffic +} + void Dispatcher::updateTxBudget() { unsigned long now = _ms->getMillis(); unsigned long elapsed = now - last_budget_update; @@ -89,6 +100,11 @@ void Dispatcher::loop() { total_air_time += t; //Serial.print(" airtime="); Serial.println(t); + if (outbound->_forwarded) { // airtime we spent relaying another node's traffic + fwd_air_time += t; + _fwd_budget.record(t, _ms->getMillis()); + } + updateTxBudget(); if (t > tx_budget_ms) { @@ -268,6 +284,14 @@ void Dispatcher::processRecvPacket(Packet* pkt) { uint8_t priority = (action >> 24) - 1; uint32_t _delay = action & 0xFFFFFF; + uint32_t air_time = _radio->getEstAirtimeFor(pkt->getRawLength()); + if (!_fwd_budget.canSpend(air_time, _ms->getMillis())) { // relay budget spent: don't queue it at all + n_fwd_dropped++; + MESH_DEBUG_PRINTLN("%s Dispatcher::processRecvPacket(): forward refused by relay airtime budget", getLogDateTime()); + _mgr->free(pkt); + return; + } + pkt->_forwarded = true; _mgr->queueOutbound(pkt, priority, futureMillis(_delay)); } } @@ -373,6 +397,7 @@ void Dispatcher::sendPacket(Packet* packet, uint8_t priority, uint32_t delay_mil MESH_DEBUG_PRINTLN("%s Dispatcher::sendPacket(): ERROR: invalid packet... path_len=%d, payload_len=%d", getLogDateTime(), (uint32_t) packet->path_len, (uint32_t) packet->payload_len); _mgr->free(packet); } else { + packet->_forwarded = false; // this node's own traffic, not a relay _mgr->queueOutbound(packet, priority, futureMillis(delay_millis)); } } @@ -387,4 +412,4 @@ unsigned long Dispatcher::futureMillis(int millis_from_now) const { return _ms->getMillis() + millis_from_now; } -} \ No newline at end of file +} diff --git a/src/Dispatcher.h b/src/Dispatcher.h index aad6cba3ec..bcd11e95ee 100644 --- a/src/Dispatcher.h +++ b/src/Dispatcher.h @@ -1,6 +1,7 @@ #pragma once #include +#include #include #include #include @@ -128,6 +129,9 @@ class Dispatcher { unsigned long tx_budget_ms; unsigned long last_budget_update; unsigned long duty_cycle_window_ms; + AirtimeBudget _fwd_budget; // rolling budget for airtime spent relaying + unsigned long fwd_air_time; // total airtime spent relaying (stat) + uint32_t n_fwd_dropped; // forwards refused by the relay budget (stat) void processRecvPacket(Packet* pkt); void updateTxBudget(); @@ -152,6 +156,8 @@ class Dispatcher { tx_budget_ms = 0; last_budget_update = 0; duty_cycle_window_ms = 3600000; + fwd_air_time = 0; + n_fwd_dropped = 0; } virtual DispatcherAction onRecvPacket(Packet* pkt) = 0; @@ -164,6 +170,7 @@ class Dispatcher { virtual const char* getLogDateTime() { return ""; } virtual float getAirtimeBudgetFactor() const; + virtual float getForwardAirtimeBudgetFactor() const; virtual int calcRxDelay(float score, uint32_t air_time) const; virtual uint32_t getCADFailRetryDelay() const; virtual uint32_t getCADFailMaxDuration() const; @@ -182,6 +189,11 @@ class Dispatcher { unsigned long getTotalAirTime() const { return total_air_time; } unsigned long getReceiveAirTime() const {return rx_air_time; } + unsigned long getForwardAirTime() const { return fwd_air_time; } + uint32_t getNumForwardDropped() const { return n_fwd_dropped; } + uint32_t getForwardBudgetUsed() { _fwd_budget.update(_ms->getMillis()); return _fwd_budget.getUsed(); } + uint32_t getForwardBudgetLimit() const { return _fwd_budget.getLimit(); } + uint32_t getForwardBudgetWindow() const { return _fwd_budget.getWindow(); } unsigned long getRemainingTxBudget() const { return tx_budget_ms; } uint32_t getNumSentFlood() const { return n_sent_flood; } uint32_t getNumSentDirect() const { return n_sent_direct; } @@ -189,6 +201,7 @@ class Dispatcher { uint32_t getNumRecvDirect() const { return n_recv_direct; } void resetStats() { n_sent_flood = n_sent_direct = n_recv_flood = n_recv_direct = 0; + n_fwd_dropped = 0; _err_flags = 0; } diff --git a/src/Packet.cpp b/src/Packet.cpp index aad3e2f48e..9cfe02ce7f 100644 --- a/src/Packet.cpp +++ b/src/Packet.cpp @@ -8,6 +8,7 @@ Packet::Packet() { header = 0; path_len = 0; payload_len = 0; + _forwarded = false; } bool Packet::isValidPathLen(uint8_t path_len) { @@ -84,4 +85,4 @@ bool Packet::readFrom(const uint8_t src[], uint8_t len) { return true; // success } -} \ No newline at end of file +} diff --git a/src/Packet.h b/src/Packet.h index c19d9e9d8f..cfd4661268 100644 --- a/src/Packet.h +++ b/src/Packet.h @@ -49,6 +49,7 @@ class Packet { uint8_t path[MAX_PATH_SIZE]; uint8_t payload[MAX_PACKET_PAYLOAD]; int8_t _snr; + bool _forwarded; // true when this packet is being re-transmitted for another node /** * \brief calculate the hash of payload + type diff --git a/src/helpers/CommonCLI.h b/src/helpers/CommonCLI.h index 8591cdc140..16c9f94aef 100644 --- a/src/helpers/CommonCLI.h +++ b/src/helpers/CommonCLI.h @@ -26,6 +26,7 @@ class NodePrefs : public ConfigSerializer { public: // in-memory backing data float airtime_factor = 0; + float fwd_airtime_factor = 0; // separate duty cycle limit for relayed traffic (0 = no separate limit) char node_name[32]; double node_lat = 0, node_lon = 0; char password[16]; @@ -89,6 +90,7 @@ class NodePrefs : public ConfigSerializer { def("fem_txgain", _parent->radio_fem_txgain); def("tx", _parent->tx_power_dbm); def("af", _parent->airtime_factor); + def("fwd_af", _parent->fwd_airtime_factor); def("rxdelay", _parent->rx_delay_base); def("f_txdelay", _parent->tx_delay_factor); def("d_txdelay", _parent->direct_tx_delay_factor); @@ -109,6 +111,8 @@ class NodePrefs : public ConfigSerializer { void setCodingRate(uint8_t cr) override { _parent->cr = cr; markDirty(); } float getAirtimeFactor() const override { return _parent->airtime_factor; } void setAirtimeFactor(float af) override { _parent->airtime_factor = af; markDirty(); } + float getForwardAirtimeFactor() const override { return _parent->fwd_airtime_factor; } + void setForwardAirtimeFactor(float af) override { _parent->fwd_airtime_factor = af; markDirty(); } bool isCadEnabled() const override { return _parent->cad_enabled; } void setCadEnabled(bool en) override { _parent->cad_enabled = en; markDirty(); } uint8_t getIntThresh() const override { return _parent->interference_threshold; } diff --git a/src/helpers/CommonRadioPrefs.cpp b/src/helpers/CommonRadioPrefs.cpp index 9d92cad9a8..4084825be7 100644 --- a/src/helpers/CommonRadioPrefs.cpp +++ b/src/helpers/CommonRadioPrefs.cpp @@ -99,6 +99,27 @@ bool CommonRadioPrefs::handleCommand(const char* command, uint32_t sender_timest return true; } + if (strcmp(command, "get fwd.dutycycle") == 0) { + float dc = 100.0f / (getForwardAirtimeFactor() + 1.0f); + int dc_int = (int)dc; + int dc_frac = (int)((dc - dc_int) * 10.0f + 0.5f); + sprintf(reply, "> %d.%d%%", dc_int, dc_frac); + return true; + } + if (memcmp(command, "set fwd.dutycycle ", 18) == 0) { + float dc = atof(&command[18]); + if (dc < 1 || dc > 100) { + strcpy(reply, "ERROR: dutycycle must be 1-100"); + } else { + setForwardAirtimeFactor((100.0f / dc) - 1.0f); + float actual = 100.0f / (getForwardAirtimeFactor() + 1.0f); + int a_int = (int)actual; + int a_frac = (int)((actual - a_int) * 10.0f + 0.5f); + sprintf(reply, "OK - %d.%d%%", a_int, a_frac); + } + return true; + } + if (strcmp(command, "get int.thresh") == 0) { sprintf(reply, "> %d", (uint32_t) getIntThresh()); return true; diff --git a/src/helpers/CommonRadioPrefs.h b/src/helpers/CommonRadioPrefs.h index 96895bb2eb..9e73563bec 100644 --- a/src/helpers/CommonRadioPrefs.h +++ b/src/helpers/CommonRadioPrefs.h @@ -26,6 +26,9 @@ class CommonRadioPrefs : public ConfigSerializer, public KeyValueStore { virtual float getAirtimeFactor() const = 0; virtual void setAirtimeFactor(float af) = 0; + virtual float getForwardAirtimeFactor() const = 0; + virtual void setForwardAirtimeFactor(float af) = 0; + virtual bool isCadEnabled() const = 0; virtual void setCadEnabled(bool en) = 0; diff --git a/src/helpers/StatsFormatHelper.h b/src/helpers/StatsFormatHelper.h index bf619133e9..d8a0ec0c2f 100644 --- a/src/helpers/StatsFormatHelper.h +++ b/src/helpers/StatsFormatHelper.h @@ -23,14 +23,22 @@ class StatsFormatHelper { mesh::Radio* radio, RadioDriverType& driver, uint32_t total_air_time_ms, - uint32_t total_rx_air_time_ms) { + uint32_t total_rx_air_time_ms, + uint32_t fwd_air_time_ms = 0, + uint32_t fwd_window_used_ms = 0, + uint32_t fwd_window_limit_ms = 0, + uint32_t n_fwd_dropped = 0) { sprintf(reply, - "{\"noise_floor\":%d,\"last_rssi\":%d,\"last_snr\":%.2f,\"tx_air_secs\":%u,\"rx_air_secs\":%u}", + "{\"noise_floor\":%d,\"last_rssi\":%d,\"last_snr\":%.2f,\"tx_air_secs\":%u,\"rx_air_secs\":%u,\"fwd_air_secs\":%u,\"fwd_window_secs\":%u,\"fwd_limit_secs\":%u,\"fwd_dropped\":%u}", (int16_t)radio->getNoiseFloor(), (int16_t)driver.getLastRSSI(), driver.getLastSNR(), total_air_time_ms / 1000, - total_rx_air_time_ms / 1000 + total_rx_air_time_ms / 1000, + fwd_air_time_ms / 1000, + fwd_window_used_ms / 1000, + fwd_window_limit_ms / 1000, + n_fwd_dropped ); } diff --git a/test/test_airtime_budget/test_airtime_budget.cpp b/test/test_airtime_budget/test_airtime_budget.cpp new file mode 100644 index 0000000000..f626dd70f0 --- /dev/null +++ b/test/test_airtime_budget/test_airtime_budget.cpp @@ -0,0 +1,293 @@ +#include +#include +#include +#include + +using namespace mesh; + +// --------------------------------------------------------------------------- +// AirtimeBudget: the rolling-window ledger itself +// --------------------------------------------------------------------------- + +TEST(AirtimeBudget, StartsEmptyAndAccumulates) { + AirtimeBudget budget; + budget.begin(3600000, 36000, 0, AirtimeBudget::MAX_SLOTS); + + EXPECT_TRUE(budget.isEnabled()); + EXPECT_EQ(0u, budget.getUsed()); + EXPECT_EQ(36000u, budget.getRemaining()); + + budget.record(1000, 0); + EXPECT_EQ(1000u, budget.getUsed()); + EXPECT_EQ(35000u, budget.getRemaining()); +} + +TEST(AirtimeBudget, DeniesSpendBeyondLimit) { + AirtimeBudget budget; + budget.begin(3600000, 1000, 0, AirtimeBudget::MAX_SLOTS); + + budget.record(900, 0); + EXPECT_FALSE(budget.canSpend(200, 0)); + EXPECT_TRUE(budget.canSpend(100, 0)); + + budget.record(100, 0); + EXPECT_EQ(1000u, budget.getUsed()); + EXPECT_EQ(0u, budget.getRemaining()); + EXPECT_FALSE(budget.canSpend(1, 0)); +} + +TEST(AirtimeBudget, IsNotEnabledWhenLimitIsFullWindow) { + AirtimeBudget budget; + budget.begin(3600000, 3600000, 0, AirtimeBudget::MAX_SLOTS); // 100% = no limit + + EXPECT_FALSE(budget.isEnabled()); + EXPECT_TRUE(budget.canSpend(3600000, 0)); +} + +// Usage is only forgotten once the slot holding it lies entirely outside the +// window, so the ledger over-holds rather than releasing airtime too early. +TEST(AirtimeBudget, ForgetsUsageOnlyAfterTheWholeWindow) { + AirtimeBudget budget; + budget.begin(1000, 500, 0, 10); // 10 x 100 ms slots + + budget.record(500, 0); + budget.record(500, 600); + EXPECT_EQ(1000u, budget.getUsed()); + + budget.update(1000); // the record at t=0 has now fallen out of the window + EXPECT_EQ(500u, budget.getUsed()); + + budget.update(1600); // and now the record at t=600 has too + EXPECT_EQ(0u, budget.getUsed()); + EXPECT_TRUE(budget.canSpend(500, 1600)); +} + +TEST(AirtimeBudget, MillisWrapIsHandled) { + AirtimeBudget budget; + uint32_t start = 0xFFFFFF00u; + budget.begin(1000, 500, start, 10); + + budget.record(300, start); + uint32_t now = start + 0x150; // wraps past zero + EXPECT_EQ(300u, budget.getUsed()); + EXPECT_TRUE(budget.canSpend(200, now)); + EXPECT_FALSE(budget.canSpend(300, now)); +} + +// --------------------------------------------------------------------------- +// Dispatcher: forwarding admission and relay accounting +// --------------------------------------------------------------------------- + +namespace { + +class FakeClock : public MillisecondClock { +public: + uint32_t now = 0; + unsigned long getMillis() override { return now; } +}; + +class FakeRadio : public Radio { +public: + uint32_t airtime_per_tx = 1000; + bool send_complete = false; + int sends_started = 0; + + int recvRaw(uint8_t* bytes, int sz) override { return 0; } + uint32_t getEstAirtimeFor(int len_bytes) override { return airtime_per_tx; } + float packetScore(float snr, int packet_len) override { return 0; } + bool startSendRaw(const uint8_t* bytes, int len) override { sends_started++; send_complete = false; return true; } + bool isSendComplete() override { return send_complete; } + void onSendFinished() override { } + bool isInRecvMode() const override { return true; } +}; + +class FakeManager : public PacketManager { + std::vector _pool; + Packet* _outbound = nullptr; + uint32_t _outbound_at = 0; + Packet* _inbound = nullptr; + uint32_t _inbound_at = 0; + + bool scheduled(uint32_t at, uint32_t now) const { return (int32_t)(now - at) >= 0; } + +public: + int n_freed = 0; + + ~FakeManager() { + for (auto p : _pool) delete p; + } + + Packet* allocNew() override { + auto p = new Packet(); + _pool.push_back(p); + return p; + } + void free(Packet* packet) override { + n_freed++; + if (packet == _outbound) _outbound = nullptr; + if (packet == _inbound) _inbound = nullptr; + } + + void queueOutbound(Packet* packet, uint8_t priority, uint32_t scheduled_for) override { + _outbound = packet; + _outbound_at = scheduled_for; + } + Packet* getNextOutbound(uint32_t now) override { + if (_outbound && scheduled(_outbound_at, now)) { + auto p = _outbound; + _outbound = nullptr; + return p; + } + return nullptr; + } + int getOutboundCount(uint32_t now) const override { return (_outbound && scheduled(_outbound_at, now)) ? 1 : 0; } + int getOutboundTotal() const override { return _outbound ? 1 : 0; } + int getFreeCount() const override { return 4; } + Packet* getOutboundByIdx(int i) override { return i == 0 ? _outbound : nullptr; } + Packet* removeOutboundByIdx(int i) override { + auto p = _outbound; + _outbound = nullptr; + return p; + } + void queueInbound(Packet* packet, uint32_t scheduled_for) override { + _inbound = packet; + _inbound_at = scheduled_for; + } + Packet* getNextInbound(uint32_t now) override { + if (_inbound && scheduled(_inbound_at, now)) { + auto p = _inbound; + _inbound = nullptr; + return p; + } + return nullptr; + } +}; + +class TestDispatcher : public Dispatcher { +public: + float fwd_factor = 0.0f; + unsigned long window_ms = 3600000; + + TestDispatcher(Radio& radio, MillisecondClock& ms, PacketManager& mgr) : Dispatcher(radio, ms, mgr) { } + + float getForwardAirtimeBudgetFactor() const override { return fwd_factor; } + +protected: + unsigned long getDutyCycleWindowMs() const override { return window_ms; } + DispatcherAction onRecvPacket(Packet* pkt) override { return ACTION_RETRANSMIT(3); } +}; + +void queueInbound(FakeManager& mgr, FakeClock& clock) { + auto pkt = mgr.allocNew(); + pkt->header = ROUTE_TYPE_FLOOD | (PAYLOAD_TYPE_TXT_MSG << PH_TYPE_SHIFT); + pkt->payload_len = 8; + mgr.queueInbound(pkt, clock.now); +} + +} // namespace + +TEST(ForwardAirtime, ForwardIsAdmittedAndMeasuredAirtimeIsRecorded) { + FakeClock clock; + clock.now = 1000; // firmware never starts the dispatcher at millis()==0 + FakeRadio radio; + FakeManager mgr; + TestDispatcher dispatcher(radio, clock, mgr); + radio.airtime_per_tx = 300; + dispatcher.begin(); + + queueInbound(mgr, clock); + clock.now += 1; // let the scheduler's 'next_tx_time' pass + dispatcher.loop(); // admit + queue + start the send + EXPECT_EQ(1, radio.sends_started); + EXPECT_EQ(0, mgr.n_freed); + + radio.send_complete = true; + clock.now += 400; // radio was busy 400 ms (measured airtime, not the estimate) + dispatcher.loop(); + + EXPECT_EQ(0u, dispatcher.getNumForwardDropped()); + EXPECT_EQ(400u, dispatcher.getForwardAirTime()); + EXPECT_EQ(400u, dispatcher.getForwardBudgetUsed()); +} + +TEST(ForwardAirtime, ForwardIsDroppedWhenBudgetIsSpent) { + FakeClock clock; + clock.now = 1000; + FakeRadio radio; + FakeManager mgr; + TestDispatcher dispatcher(radio, clock, mgr); + dispatcher.window_ms = 600000; // 10 minutes + dispatcher.fwd_factor = 9.0f; // 10% of the window may be relayed + radio.airtime_per_tx = 70000; // one forward alone exceeds that + dispatcher.begin(); + + queueInbound(mgr, clock); + dispatcher.loop(); + + EXPECT_EQ(0, radio.sends_started); + EXPECT_EQ(1, mgr.n_freed); + EXPECT_EQ(1u, dispatcher.getNumForwardDropped()); + EXPECT_EQ(0u, dispatcher.getForwardAirTime()); +} + +TEST(ForwardAirtime, OwnTrafficDoesNotConsumeTheRelayBudget) { + FakeClock clock; + clock.now = 1000; + FakeRadio radio; + FakeManager mgr; + TestDispatcher dispatcher(radio, clock, mgr); + dispatcher.window_ms = 600000; + dispatcher.fwd_factor = 9.0f; + radio.airtime_per_tx = 1000; + dispatcher.begin(); + + auto pkt = mgr.allocNew(); + pkt->header = ROUTE_TYPE_FLOOD | (PAYLOAD_TYPE_TXT_MSG << PH_TYPE_SHIFT); + dispatcher.sendPacket(pkt, 0); + EXPECT_FALSE(pkt->_forwarded); + + clock.now += 1; + dispatcher.loop(); + EXPECT_EQ(1, radio.sends_started); + radio.send_complete = true; + clock.now += 900; + dispatcher.loop(); + + EXPECT_EQ(0u, dispatcher.getForwardAirTime()); + EXPECT_EQ(0u, dispatcher.getForwardBudgetUsed()); +} + +TEST(ForwardAirtime, RelayBudgetFreesUpAfterTheWindow) { + FakeClock clock; + clock.now = 1000; + FakeRadio radio; + FakeManager mgr; + TestDispatcher dispatcher(radio, clock, mgr); + dispatcher.window_ms = 1000; + dispatcher.fwd_factor = 1.0f; // 50% of the window + radio.airtime_per_tx = 400; + dispatcher.begin(); + + queueInbound(mgr, clock); + clock.now += 1; + dispatcher.loop(); + EXPECT_EQ(1, radio.sends_started); + radio.send_complete = true; + clock.now += 400; + dispatcher.loop(); + EXPECT_EQ(400u, dispatcher.getForwardBudgetUsed()); + + queueInbound(mgr, clock); + dispatcher.loop(); // 400 + 400 > 500: refused + EXPECT_EQ(1u, dispatcher.getNumForwardDropped()); + + clock.now += 1000; // let the whole window pass + queueInbound(mgr, clock); + dispatcher.loop(); + EXPECT_EQ(2, radio.sends_started); // admitted again +} + +int main(int argc, char** argv) { + ::testing::InitGoogleTest(&argc, argv); + return RUN_ALL_TESTS(); +}