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(); +}