Complete ESP32 partition expander and adaptive LoRa OTA pacing

This commit is contained in:
mikecarper
2026-09-24 17:13:22 -07:00
parent e370db7f28
commit a74a09da56
33 changed files with 986 additions and 66 deletions
+2 -2
View File
@@ -849,7 +849,7 @@ void Dispatcher::checkSend() {
// override. RX processing and ordinary watchdogs continue in loop().
if (!restoreOutboundTxOverrides()) return;
const uint32_t now = _ms->getMillis();
if (ota_tx_airtime && now - ota_tx_finished_at >= ota::packetQuietTime(ota_tx_airtime, getOtaSpeedFactor())) {
if (ota_tx_airtime && now - ota_tx_finished_at >= ota::packetQuietTime(ota_tx_airtime, getOtaPacketSpeedFactor())) {
ota_tx_airtime = 0;
}
uint32_t next_outbound;
@@ -870,7 +870,7 @@ void Dispatcher::checkSend() {
return;
}
if (ota_tx_airtime) {
const uint32_t gap = ota::packetQuietTime(ota_tx_airtime, getOtaSpeedFactor());
const uint32_t gap = ota::packetQuietTime(ota_tx_airtime, getOtaPacketSpeedFactor());
const uint32_t elapsed = now - ota_tx_finished_at;
if (elapsed >= gap) ota_tx_airtime = 0;
else if (pending && pending->getPayloadType() == PAYLOAD_TYPE_OTA) {
+1
View File
@@ -423,6 +423,7 @@ protected:
virtual float getAirtimeBudgetFactor() const;
virtual float getOtaSpeedFactor() const { return 1.0f; }
virtual float getOtaPacketSpeedFactor() const { return getOtaSpeedFactor(); }
virtual int calcRxDelay(float score, uint32_t air_time) const;
virtual bool shouldBypassRxDelay(const Packet* packet) {
(void)packet;
+10
View File
@@ -663,6 +663,16 @@ float Mesh::getOtaSpeedFactor() const {
#endif
}
float Mesh::getOtaPacketSpeedFactor() const {
const float configured = getOtaSpeedFactor();
#if defined(ENABLE_OTA)
const float adaptive = ota::ota_ctx().manager.adaptivePacketSpeed();
return configured < adaptive ? configured : adaptive;
#else
return configured;
#endif
}
uint32_t Mesh::getOtaPacketAirtime() const {
const auto* profiles = _radio->profiles();
if (!profiles || !profiles->enabled()) return _radio->getProfileAirtime(0, MAX_TRANS_UNIT);
+1
View File
@@ -542,6 +542,7 @@ protected:
bool _ota_temp_was_active = false; // detects entry into a temporary-radio window
#endif
float getOtaSpeedFactor() const override;
float getOtaPacketSpeedFactor() const override;
uint32_t getOtaPacketAirtime() const;
/**
+26 -1
View File
@@ -792,6 +792,28 @@ void CommonCLI::loadPrefs(FILESYSTEM* fs) {
}
_com_prefs_needs_upgrade = false;
#endif
} else if (fs->exists("/prefs.json")) {
// Stock 1.17.1 stored NodePrefs as JSON. A partition migration preserves
// that file byte-for-byte; import it before writing the current binary
// image so node name, passwords and primary radio survive the Full handoff.
#if defined(NRF52_PLATFORM)
File legacy(*fs);
legacy.open("/prefs.json", FILE_O_READ);
#elif defined(STM32_PLATFORM)
File legacy = fs->open("/prefs.json", FILE_O_READ);
#else
File legacy = fs->open("/prefs.json", "r");
#endif
if (legacy && _prefs->loadSerial(legacy)) {
loaded = true;
is_upgrade = true;
#ifdef WITH_MQTT_BRIDGE
node_prefs_needs_migration = true;
#else
savePrefs(fs);
#endif
}
if (legacy) legacy.close();
} else {
// File doesn't exist - set defaults for a fresh install. Dual R1/R2
// scanning keeps the node awake, so only that configuration starts with
@@ -2579,7 +2601,10 @@ void CommonCLI::handleCommand(uint32_t sender_timestamp, char* command, char* re
_callbacks->wirelessCommandSource(sender_timestamp))) return;
if (_radio_profiles.handle(command, reply, 160, sender_timestamp != 0)) return;
#if defined(ENABLE_OTA)
if (mesh::ota::handleSpeedCommand(command, reply, 160)) return;
const auto* ota_context = mesh::ota::ota_context_if_active();
const float adaptive_ota_pace = ota_context
? ota_context->manager.adaptivePacketSpeed() : mesh::ota::OTA_SPEED_DEFAULT;
if (mesh::ota::handleSpeedCommand(command, reply, 160, adaptive_ota_pace)) return;
#endif
if (strncmp(command, "set tempradio ", 14) == 0) {
handleCommand(sender_timestamp, command + 4, reply);
+9 -5
View File
@@ -147,7 +147,9 @@ static bool is_cmd(const char* a, const char* names, const char** rest) {
static bool handle_dev(const char* d, char* reply, OtaContext& c);
bool handle_ota_command(const char* command, char* reply, mesh::MainBoard& board) {
if (handleSpeedCommand(command, reply, 160)) return true;
const auto* active = ota_context_if_active();
const float adaptive_pace = active ? active->manager.adaptivePacketSpeed() : OTA_SPEED_DEFAULT;
if (handleSpeedCommand(command, reply, 160, adaptive_pace)) return true;
const char* a = command + 3;
if (*a != 0 && *a != ' ') return false;
while (*a == ' ') a++;
@@ -1002,10 +1004,12 @@ bool handle_ota_command(const char* command, char* reply, mesh::MainBoard& board
strcpy(reply, "ERR unknown OTA config setting");
} else { // show current policy
uint8_t af = c.manager.autofetch();
char speed[16]; formatSpeed(speed, sizeof(speed));
char speed[16], packet[16];
formatSpeed(speed, sizeof(speed));
formatSpeedFactor(effectivePacketPace(c.manager.adaptivePacketSpeed()), packet, sizeof(packet));
#if defined(NRF52_PLATFORM) && defined(OTA_SD_STORE)
bool cache_ready = c.ensureSdCache();
snprintf(reply, 160, "ota config: speed=%sx cache=%s/%u autofetch=%s autoinstall=%s checkpoint=%u advert=%umin hops=%u keys=%u", speed,
snprintf(reply, 160, "ota config: speed=%sx packet=%sx cache=%s/%u autofetch=%s autoinstall=%s checkpoint=%u advert=%umin hops=%u keys=%u", speed, packet,
cache_ready ? (c.sd_cache.autoCaptureEnabled() ? "on" : "off") : "unavailable",
cache_ready ? (unsigned)c.sd_cache.capturedCount() : 0,
af == OtaManager::AUTOFETCH_ANY ? "any" : af == OtaManager::AUTOFETCH_SIGNED ? "signed" : "off",
@@ -1014,11 +1018,11 @@ bool handle_ota_command(const char* command, char* reply, mesh::MainBoard& board
(unsigned)c.manager.max_hops(), (unsigned)c.allow.count());
#else
#if defined(OTA_SEEDER_ONLY)
snprintf(reply, 160, "ota config: speed=%sx mode=seeder-only autofetch=off autoinstall=off checkpoint=%u advert=%umin hops=%u", speed,
snprintf(reply, 160, "ota config: speed=%sx packet=%sx mode=seeder-only autofetch=off autoinstall=off checkpoint=%u advert=%umin hops=%u", speed, packet,
(unsigned)c.manager.checkpoint_blocks(), (unsigned)c.manager.advert_mins(),
(unsigned)c.manager.max_hops());
#else
snprintf(reply, 160, "ota config: speed=%sx autofetch=%s autoinstall=%s checkpoint=%u advert=%umin hops=%u keys=%u (persisted)", speed,
snprintf(reply, 160, "ota config: speed=%sx packet=%sx autofetch=%s autoinstall=%s checkpoint=%u advert=%umin hops=%u keys=%u (persisted)", speed, packet,
af == OtaManager::AUTOFETCH_ANY ? "any" : af == OtaManager::AUTOFETCH_SIGNED ? "signed" : "off",
c.autoinstall == OtaContext::AUTOINSTALL_TRUSTED ? "trusted" : "off",
(unsigned)c.manager.checkpoint_blocks(), (unsigned)c.manager.advert_mins(),
+64 -3
View File
@@ -94,6 +94,10 @@ void OtaManager::begin(uint32_t my_target_id, OtaSend send, void* ctx) {
memset(_src_advertised, 0, sizeof(_src_advertised));
_n_src = 0; _n_cat = 0;
clearPendingEgress();
_adaptive_packet_speed = 1.0f;
_pace_clean_blocks = 0;
_pace_has_mid = false;
_pace_has_loss = false;
}
// ---------------- serve (multi-mota registry) ----------------
@@ -682,6 +686,45 @@ bool OtaManager::loadActiveServeBlock() {
return true;
}
void OtaManager::noteServedRequestPacing(const uint8_t* mid, uint16_t block,
uint16_t want, uint16_t full_mask) {
if (want != full_mask) {
// A sparse mask is the receiver's explicit evidence that our previous
// DATA burst left holes. One extra packet-airtime gap is much cheaper than
// repeated multi-second request turns on a fast link.
if (!_pace_has_loss || memcmp(_pace_loss_mid, mid, sizeof(_pace_loss_mid)) != 0
|| _pace_last_loss_block != block) {
memcpy(_pace_loss_mid, mid, sizeof(_pace_loss_mid));
_pace_last_loss_block = block;
_pace_has_loss = true;
if (_adaptive_packet_speed > 0.5f) _adaptive_packet_speed = 0.5f;
else if (_adaptive_packet_speed > 0.25f) {
_adaptive_packet_speed -= 0.1f;
if (_adaptive_packet_speed < 0.25f) _adaptive_packet_speed = 0.25f;
}
}
_pace_clean_blocks = 0;
return;
}
if (!_pace_has_mid || memcmp(_pace_mid, mid, sizeof(_pace_mid)) != 0) {
memcpy(_pace_mid, mid, sizeof(_pace_mid));
_pace_has_mid = true;
_pace_last_full_block = block;
_pace_clean_blocks = 0;
return;
}
// Count only forward progress; duplicates of a full request are not proof
// that the receiver tolerated the current packet rate.
if (block <= _pace_last_full_block) return;
_pace_last_full_block = block;
if (_adaptive_packet_speed >= 1.0f) return;
if (++_pace_clean_blocks >= 16) {
_adaptive_packet_speed += 0.1f;
if (_adaptive_packet_speed > 1.0f) _adaptive_packet_speed = 1.0f;
_pace_clean_blocks = 0;
}
}
bool OtaManager::handleReq(const uint8_t* m, uint16_t n) {
ReqWindowMsg rq;
if (!decode_req_window(m, n, rq)) return false;
@@ -708,8 +751,11 @@ bool OtaManager::handleReq(const uint8_t* m, uint16_t n) {
const uint16_t requested = wire_v2
? ota_req_v2_fragments(rq.items[i].want_mask) : rq.items[i].want_mask;
const uint16_t want = (uint16_t)(requested & valid_mask);
if (want != 0) accepted |= queueServeJob(v->m.merkle_root, (uint16_t)idx, want,
wire_v2, allow_deflate, extended_length);
if (want != 0 && queueServeJob(v->m.merkle_root, (uint16_t)idx, want,
wire_v2, allow_deflate, extended_length)) {
accepted = true;
noteServedRequestPacing(v->m.merkle_root, (uint16_t)idx, want, valid_mask);
}
}
return accepted;
}
@@ -1223,7 +1269,22 @@ uint32_t OtaManager::fetchRetryTimeoutMs() const {
service = (service * _tx_spacing_permille + 999u) / 1000u;
service = (service * 5u + 3u) / 4u; // 25% CAD/relay contention allowance
service += OTA_FETCH_RETRY_GUARD_MS;
if (service < OTA_FETCH_RETRY_MIN_MS) service = OTA_FETCH_RETRY_MIN_MS;
uint64_t floor = OTA_FETCH_RETRY_MIN_MS;
if (_observed_path_transmissions != 0 && anyWireDataReceived()) {
// A valid fragment proves this source has already answered the current flight.
// Missing-fragment requests need not inherit the five-second allowance for
// an unanswered host/relay request. Four full-packet service intervals cover
// queued DATA/proof/turnaround; measured airtime and the observed path raise
// the floor automatically on slower or relayed links.
uint64_t partial_floor = (uint64_t)_radio_packet_airtime_ms
* _observed_path_transmissions * _tx_spacing_permille * 4u;
partial_floor = (partial_floor + 999u) / 1000u;
partial_floor += OTA_FETCH_RETRY_GUARD_MS;
if (partial_floor < OTA_FETCH_RETRY_PARTIAL_MIN_MS)
partial_floor = OTA_FETCH_RETRY_PARTIAL_MIN_MS;
if (partial_floor < floor) floor = partial_floor;
}
if (service < floor) service = floor;
if (service > OTA_FETCH_RETRY_MAX_MS) service = OTA_FETCH_RETRY_MAX_MS;
return retryDelay((uint32_t)service);
}
+17 -1
View File
@@ -215,7 +215,10 @@ static constexpr uint16_t ota_max_block_capability() { return (uint16_t)OTA_MAX_
#error "OTA_FETCH_PIPELINE_INITIAL must be between 1 and OTA_FETCH_PIPELINE"
#endif
#ifndef OTA_FETCH_RETRY_MIN_MS
#define OTA_FETCH_RETRY_MIN_MS 5000 // quiet-time floor: covers host/queue/turnaround latency on fast links
#define OTA_FETCH_RETRY_MIN_MS 5000 // unanswered/empty flight: allow host and relay turnaround
#endif
#ifndef OTA_FETCH_RETRY_PARTIAL_MIN_MS
#define OTA_FETCH_RETRY_PARTIAL_MIN_MS 1500 // after valid DATA, recover missing holes promptly on fast links
#endif
#ifndef OTA_FETCH_RETRY_MAX_MS
#define OTA_FETCH_RETRY_MAX_MS 60000 // bound loss recovery when configured radio settings are impractical
@@ -409,6 +412,9 @@ public:
void set_clock(uint32_t ms) { _now_ms = ms; }
bool set_speed(float speed);
float speed() const { return _speed; }
// Sender-only, loss-responsive packet pacing. This does not alter discovery,
// proof, or retry timers; the user's persisted speed remains an upper bound.
float adaptivePacketSpeed() const { return _adaptive_packet_speed; }
uint32_t pacedDelay(uint32_t ms) const { return scaleDelay(ms, _speed); }
// Faster pacing cannot shorten a loss-recovery deadline below the existing
// physical packet-flight allowance. Slower pacing expands that allowance.
@@ -620,6 +626,8 @@ private:
bool queueServeJob(const uint8_t* mid, uint16_t block, uint16_t want_mask,
bool wire_v2 = false, bool allow_deflate = false,
bool extended_length = false);
void noteServedRequestPacing(const uint8_t* mid, uint16_t block,
uint16_t want, uint16_t full_mask);
bool queueManifestJob(const uint8_t* mid, uint16_t want_mask);
uint32_t manifestEgressGapMs() const;
uint32_t proofEgressGapMs() const;
@@ -699,6 +707,14 @@ private:
};
ServeJob _serve_jobs[OTA_SERVE_QUEUE];
uint8_t _n_serve_jobs = 0;
float _adaptive_packet_speed = 1.0f;
uint8_t _pace_clean_blocks = 0;
uint8_t _pace_mid[4] = {0};
uint16_t _pace_last_full_block = 0;
bool _pace_has_mid = false;
uint8_t _pace_loss_mid[4] = {0};
uint16_t _pace_last_loss_block = 0;
bool _pace_has_loss = false;
struct ManifestServeJob {
OtaReplyRoute route;
uint8_t mid[4];
+18 -5
View File
@@ -156,8 +156,8 @@ void beginSpeedConfig(FILESYSTEM* fs) {
float speedFactor() { return factor; }
void formatSpeed(char* text, size_t capacity) {
const uint32_t scaled = (uint32_t)(factor * 1000000.0f + 0.5f);
void formatSpeedFactor(float speed, char* text, size_t capacity) {
const uint32_t scaled = (uint32_t)(speed * 1000000.0f + 0.5f);
snprintf(text, capacity, "%u.%06u", (unsigned)(scaled / 1000000), (unsigned)(scaled % 1000000));
if (!capacity) return;
size_t n = strlen(text);
@@ -165,7 +165,17 @@ void formatSpeed(char* text, size_t capacity) {
if (n && text[n - 1] == '.') text[n - 1] = 0;
}
bool handleSpeedCommand(const char* command, char* reply, size_t capacity) {
void formatSpeed(char* text, size_t capacity) {
formatSpeedFactor(factor, text, capacity);
}
float effectivePacketPace(float adaptive_speed) {
if (!validSpeed(adaptive_speed)) adaptive_speed = OTA_SPEED_DEFAULT;
return factor < adaptive_speed ? factor : adaptive_speed;
}
bool handleSpeedCommand(const char* command, char* reply, size_t capacity,
float adaptive_speed) {
const char* value = nullptr;
bool require_value = false, read_only = false;
const char* const forms[] = {"set ota.speed", "get ota.speed", "ota config speed", "ota cfg speed", "ota set speed", "ota speed"};
@@ -194,8 +204,11 @@ bool handleSpeedCommand(const char* command, char* reply, size_t capacity) {
return true;
}
}
char number[16]; formatSpeed(number, sizeof(number));
snprintf(reply, capacity, *value ? "OK ota.speed=%sx (saved)" : "> ota.speed=%sx", number);
char number[16], packet[16];
formatSpeed(number, sizeof(number));
formatSpeedFactor(effectivePacketPace(adaptive_speed), packet, sizeof(packet));
if (*value) snprintf(reply, capacity, "OK ota.speed=%sx (saved) packet=%sx", number, packet);
else snprintf(reply, capacity, "> ota.speed=%sx packet=%sx", number, packet);
return true;
}
+6 -1
View File
@@ -10,6 +10,11 @@ namespace mesh { namespace ota {
void beginSpeedConfig(FILESYSTEM* fs);
float speedFactor();
void formatSpeed(char* text, size_t capacity);
bool handleSpeedCommand(const char* command, char* reply, size_t capacity);
void formatSpeedFactor(float speed, char* text, size_t capacity);
// Automatic sender pacing only adds packet quiet time. The saved setting is
// still the upper bound and continues to control the other OTA timers.
float effectivePacketPace(float adaptive_speed);
bool handleSpeedCommand(const char* command, char* reply, size_t capacity,
float adaptive_speed = OTA_SPEED_DEFAULT);
} }