Files
HaloKeymind/src/Dispatcher.cpp
T

1125 lines
44 KiB
C++

#include "Dispatcher.h"
#include "helpers/ota/OtaTiming.h"
#if MESH_PACKET_LOGGING
#include <Arduino.h>
#include <helpers/SerialPacketLog.h>
#endif
#include <math.h>
namespace mesh {
#define MAX_RX_DELAY_MILLIS 32000 // 32 seconds
#define MIN_TX_BUDGET_RESERVE_MS 100 // min budget (ms) required before allowing next TX
#define MIN_TX_BUDGET_AIRTIME_DIV 2 // require at least 1/N of estimated airtime as budget before TX
#ifndef NOISE_FLOOR_CALIB_INTERVAL
#define NOISE_FLOOR_CALIB_INTERVAL 2000 // refresh every 2 seconds
#endif
#ifndef RADIO_LIVENESS_SOFT_MS
#define RADIO_LIVENESS_SOFT_MS (30UL * 60UL * 1000UL)
#endif
#ifndef RADIO_LIVENESS_HARD_MS
#define RADIO_LIVENESS_HARD_MS (12UL * 60UL * 60UL * 1000UL)
#endif
#define MIN_CAD_FAIL_RETRY_DELAY_MS 50UL
#define MIN_CAD_FAIL_MAX_DURATION_MS 500UL
#define RADIO_PREPARE_BUSY_GRACE_MS 8000UL
static uint32_t validPrepareBusyAirtime(uint32_t airtime) {
// RadioLib estimates originate as 32-bit microseconds. An impossible or
// error-sentinel estimate must not retain a failed packet for days.
return airtime <= UINT32_MAX / 1000 ? airtime : 0;
}
static uint32_t scaleCADDelayForQueue(uint32_t normal_delay, int ready_count,
uint32_t minimum_delay) {
if (ready_count <= 1) return normal_delay;
uint32_t scaled = normal_delay / (uint32_t)ready_count;
uint32_t floor = normal_delay < minimum_delay ? normal_delay : minimum_delay;
return scaled < floor ? floor : scaled;
}
void Dispatcher::setRadioAvailable(bool available) {
if (radio_available == available) return;
radio_available = available;
if (!dispatcher_started || !radio_available) return;
const unsigned long now = _ms->getMillis();
radio_nonrx_start = now;
next_floor_calib_time = now;
armed_agc_reset_interval = 0;
agc_reset_armed = false;
_radio->begin();
prev_isrecv_mode = _radio->isInRecvMode();
#ifdef RADIO_LIVENESS_SOFT_ONLY
rx_watchdog_window_start = now;
#else
radio_liveness.begin(now);
nonrx_soft_recovery_attempted = false;
#endif
MESH_DEBUG_PRINTLN("Dispatcher: radio transport activated");
}
void Dispatcher::begin() {
n_sent_flood = n_sent_direct = 0;
n_recv_flood = n_recv_direct = 0;
last_meshcore_recv_millis = 0;
_err_flags = 0;
outbound_radio_retry_at = 0;
outbound_radio_prepare_deadline = 0;
outbound_radio_retry_pending = false;
outbound_radio_retry_used = false;
outbound_cancellation = OutboundCancellation::None;
ota_tx_airtime = 0;
radio_nonrx_start = _ms->getMillis();
duty_cycle_window_ms = getDutyCycleWindowMs();
float duty_cycle = 1.0f / (1.0f + getAirtimeBudgetFactor());
tx_budget_ms = (unsigned long)(duty_cycle_window_ms * duty_cycle);
last_budget_update = _ms->getMillis();
const unsigned long now = _ms->getMillis();
armed_agc_reset_interval = 0;
agc_reset_armed = false;
dispatcher_started = true;
if (radio_available) {
_radio->begin();
prev_isrecv_mode = _radio->isInRecvMode();
} else {
prev_isrecv_mode = false;
MESH_DEBUG_PRINTLN("Dispatcher: starting without a radio transport");
}
#ifdef RADIO_LIVENESS_SOFT_ONLY
rx_watchdog_window_start = now;
#else
radio_liveness.begin(now);
nonrx_soft_recovery_attempted = false;
#endif
}
float Dispatcher::getAirtimeBudgetFactor() const {
return 1.0;
}
void Dispatcher::updateTxBudget() {
unsigned long now = _ms->getMillis();
unsigned long elapsed = now - last_budget_update;
float duty_cycle = 1.0f / (1.0f + getAirtimeBudgetFactor());
unsigned long max_budget = (unsigned long)(getDutyCycleWindowMs() * duty_cycle);
unsigned long refill = (unsigned long)(elapsed * duty_cycle);
if (refill > 0) {
tx_budget_ms += refill;
if (tx_budget_ms > max_budget) {
tx_budget_ms = max_budget;
}
last_budget_update = now;
}
}
bool Dispatcher::restoreOutboundTxOverrides() {
if (!outbound_restore_cr) return true;
const uint32_t now = _ms->getMillis();
if (outbound_restore_retry_at && (int32_t)(now - outbound_restore_retry_at) < 0) return false;
// A scanner or an intervening settings change may have selected another
// profile since the failed restore. Restore its configured CR, not a stale
// default from the preceding packet.
const uint8_t cr = _radio->profiles()
? _radio->profiles()->params(_radio->receiveProfile()).cr : getDefaultTxCodingRate();
const auto result = _radio->tryRestoreCodingRate(cr >= 4 && cr <= 8 ? cr : outbound_restore_cr);
if (result == RadioParamApplyResult::APPLIED) {
outbound_restore_cr = 0;
outbound_restore_retry_at = outbound_restore_recovery_at = 0;
return true;
}
outbound_restore_retry_at = now + 250;
if (result == RadioParamApplyResult::BUSY) {
// A legitimate RX/preamble/duty-cycle sleep is not a failed modulation
// write. Let receive service drain it without forcing a radio reset.
outbound_restore_recovery_at = 0;
return false;
}
if (!outbound_restore_recovery_at) outbound_restore_recovery_at = now + RADIO_PREPARE_BUSY_GRACE_MS;
else if ((int32_t)(now - outbound_restore_recovery_at) >= 0) {
_err_flags |= ERR_EVENT_RADIO_WATCHDOG;
_radio->recoverRadio(true);
outbound_restore_recovery_at = now + RADIO_PREPARE_BUSY_GRACE_MS;
}
return false;
}
bool Dispatcher::startOutboundTransmit() {
if (outbound == NULL) return false;
if (!restoreOutboundTxOverrides()) return false;
if (!isPacketRadioCurrent(outbound)
|| _radio->prepareTransmitProfile(outbound->radio_profile,
outbound->radio_reply && outbound->radio_reply_force) != RadioParamApplyResult::APPLIED) return false;
int len = 0;
uint8_t raw[MAX_TRANS_UNIT];
raw[len++] = outbound->header;
if (outbound->hasTransportCodes()) {
memcpy(&raw[len], &outbound->transport_codes[0], 2); len += 2;
memcpy(&raw[len], &outbound->transport_codes[1], 2); len += 2;
}
raw[len++] = outbound->path_len;
len += Packet::writePath(&raw[len], outbound->path, outbound->path_len);
if (len + outbound->payload_len > MAX_TRANS_UNIT) {
MESH_DEBUG_PRINTLN("%s Dispatcher::startOutboundTransmit(): FATAL: Invalid packet queued... too long, len=%d",
getLogDateTime(), len + outbound->payload_len);
return false;
}
memcpy(&raw[len], outbound->payload, outbound->payload_len);
len += outbound->payload_len;
uint32_t max_airtime = _radio->getEstAirtimeFor(len) * 3 / 2;
outbound_restore_cr = 0;
uint8_t default_cr = getDefaultTxCodingRate();
if (_radio->profiles()) default_cr = _radio->profiles()->params(outbound->radio_profile).cr;
if (outbound->tx_cr >= 4 && outbound->tx_cr <= 8
&& default_cr >= 4 && default_cr <= 8
&& outbound->tx_cr != default_cr) {
if (_radio->setCodingRate(outbound->tx_cr)) {
outbound_restore_cr = default_cr;
max_airtime = _radio->getEstAirtimeFor(len) * 3 / 2;
} else {
MESH_DEBUG_PRINTLN("%s Dispatcher::startOutboundTransmit(): WARN: failed to set packet CR%d",
getLogDateTime(), (uint32_t)outbound->tx_cr);
}
}
outbound_start = _ms->getMillis();
if (!_radio->startSendRaw(raw, len)) {
MESH_DEBUG_PRINTLN("%s Dispatcher::startOutboundTransmit(): ERROR: send start failed!",
getLogDateTime());
restoreOutboundTxOverrides();
return false;
}
outbound_expiry = futureMillis(max_airtime);
#if MESH_PACKET_LOGGING
if (isUsbLoggingEnabled()) {
#if MESH_PACKET_LOGGING_COMPACT
SerialLogLine<> line;
line.printf("T");
line.hex(raw, len);
line.flush(usbLoggingPort(), false);
#else
logPacketLine("TX", outbound, len, false, 0.0f, 0);
#endif
}
#endif
return true;
}
bool Dispatcher::scheduleOutboundRadioRetry() {
if (outbound == NULL || outbound_radio_retry_used || outbound_cancellation != OutboundCancellation::None) return false;
outbound_radio_retry_used = true;
outbound_radio_retry_pending = true;
outbound_radio_prepare_deadline = 0;
outbound_radio_retry_at = futureMillis(getCADFailRetryDelay());
MESH_DEBUG_PRINTLN("%s Dispatcher: retrying packet once after radio fault",
getLogDateTime());
return true;
}
void Dispatcher::failOutboundTransmit() {
if (outbound == NULL) return;
restoreOutboundTxOverrides();
if (outbound_cancellation == OutboundCancellation::Delivered) {
// A downstream echo confirms delivery even if the driver's completion
// failed. Pending command replies must complete their application barrier.
onSendComplete(outbound);
} else {
if (outbound_cancellation == OutboundCancellation::None) logTxFail(outbound, outbound->getRawLength());
// Explicit cancellation still releases application references, including
// pending replies and battery alerts, through the ordinary failure hook.
onSendFail(outbound);
}
releasePacket(outbound);
outbound = NULL;
outbound_radio_retry_pending = false;
outbound_radio_retry_used = false;
outbound_radio_prepare_deadline = 0;
outbound_cancellation = OutboundCancellation::None;
}
int Dispatcher::calcRxDelay(float score, uint32_t air_time) const {
return (int) ((powf(10.0f, 0.85f - score) - 1.0f) * air_time);
}
uint32_t Dispatcher::getCADFailRetryDelay() const {
return 200;
}
uint32_t Dispatcher::getCADFailMaxDuration() const {
return 4000; // 4 seconds
}
#if MESH_PACKET_LOGGING
void Dispatcher::logPacketLine(const char* direction, const Packet* packet,
int len, bool include_rx_metrics, float score,
uint32_t air_time) {
SerialLogLine<256> line;
line.printf("%s: %s, len=%d (type=%d, route=%s, payload_len=%d)",
getLogDateTime(), direction, len, packet->getPayloadType(),
packet->isRouteDirect() ? "D" : "F", packet->payload_len);
if (include_rx_metrics) {
line.printf(" SNR=%d RSSI=%d score=%d time=%u", (int)packet->getSNR(),
(int)packet->getRSSI(), (int)(score * 1000),
(unsigned)air_time);
uint8_t packet_hash[MAX_HASH_SIZE];
packet->calculatePacketHash(packet_hash);
line.printf(" hash=");
line.hex(packet_hash, MAX_HASH_SIZE);
}
const uint8_t type = packet->getPayloadType();
if (packet->payload_len >= 2
&& (type == PAYLOAD_TYPE_PATH || type == PAYLOAD_TYPE_REQ
|| type == PAYLOAD_TYPE_RESPONSE || type == PAYLOAD_TYPE_TXT_MSG)) {
line.printf(" [%02X -> %02X]", (uint32_t)packet->payload[1],
(uint32_t)packet->payload[0]);
}
line.flush(usbLoggingPort());
}
#endif
bool Dispatcher::getNextQueueWakeDelay(uint32_t& delay_millis) const {
if (isDualRadioActive()) { delay_millis = 0; return true; }
const uint32_t now = _ms->getMillis();
bool found = false;
uint32_t shortest_delay = 0;
if (outbound != NULL) {
if (outbound_radio_retry_pending) {
int32_t signed_retry_delay = (int32_t)(outbound_radio_retry_at - now);
delay_millis = signed_retry_delay > 0 ? (uint32_t)signed_retry_delay : 0;
} else {
// TX completion/timeout still needs the normal fast lifecycle path.
delay_millis = 0;
}
return true;
}
if (outbound_restore_cr) {
const int32_t restore_delay = (int32_t)(outbound_restore_retry_at - now);
shortest_delay = restore_delay > 0 ? (uint32_t)restore_delay : 0;
found = true;
}
uint32_t scheduled_for;
if (_mgr->getNextOutboundTime(now, scheduled_for)) {
int32_t signed_queue_delay = (int32_t)(scheduled_for - now);
uint32_t outbound_delay = signed_queue_delay > 0 ? (uint32_t)signed_queue_delay : 0;
int32_t signed_tx_delay = (int32_t)(next_tx_time - now);
if (signed_tx_delay > 0 && (uint32_t)signed_tx_delay > outbound_delay) {
outbound_delay = (uint32_t)signed_tx_delay;
}
if (!found || outbound_delay < shortest_delay) shortest_delay = outbound_delay;
found = true;
}
if (_mgr->getNextInboundTime(now, scheduled_for)) {
int32_t signed_inbound_delay = (int32_t)(scheduled_for - now);
uint32_t inbound_delay = signed_inbound_delay > 0 ? (uint32_t)signed_inbound_delay : 0;
if (!found || inbound_delay < shortest_delay) {
shortest_delay = inbound_delay;
found = true;
}
}
if (found) delay_millis = shortest_delay;
return found;
}
#ifdef WITH_MQTT_BRIDGE
uint32_t Dispatcher::getRadioWatchdogMillis() const {
return RADIO_WATCHDOG_MS;
}
#endif
void Dispatcher::loop() {
if (!radio_available) {
// Keep packet-manager cleanup alive, but never touch an uninitialized or
// physically unavailable radio. sendPacket() rejects new outbound work
// while this state is active, so no queue can silently fill up.
releaseDroppedOutbound();
return;
}
if (millisHasNowPassed(next_floor_calib_time)) {
_radio->triggerNoiseFloorCalibrate(getInterferenceThreshold());
_radio->setCADEnabled(getCADEnabled());
next_floor_calib_time = futureMillis(NOISE_FLOOR_CALIB_INTERVAL);
}
_radio->loop();
// A maintenance carrier has no TxDone packet. Keep its deadline serviced,
// but leave queued packets and RX/AGC watchdogs alone until normal RX returns.
if (_radio->isCarrierWaveActive()) return;
const unsigned long now = _ms->getMillis();
// check for radio 'stuck' in mode other than Rx
bool is_recv = _radio->isInRecvMode();
if (is_recv != prev_isrecv_mode) {
prev_isrecv_mode = is_recv;
if (!is_recv) {
radio_nonrx_start = now;
}
#ifndef RADIO_LIVENESS_SOFT_ONLY
else {
nonrx_soft_recovery_attempted = false;
}
#endif
}
bool recovered_this_loop = false;
if (!is_recv && outbound == NULL && now - radio_nonrx_start > 8000) {
_err_flags |= ERR_EVENT_STARTRX_TIMEOUT;
#ifdef RADIO_LIVENESS_SOFT_ONLY
MESH_DEBUG_PRINTLN("Radio watchdog: radio outside RX for %lu ms; soft recovery",
now - radio_nonrx_start);
if (_radio->recoverRadio(false)) {
rx_watchdog_window_start = now;
recovered_this_loop = true;
}
#else
const bool hard = nonrx_soft_recovery_attempted;
MESH_DEBUG_PRINTLN("Radio watchdog: radio outside RX for %lu ms; %s recovery",
now - radio_nonrx_start, hard ? "hard" : "soft");
if (_radio->recoverRadio(hard)) recovered_this_loop = true;
nonrx_soft_recovery_attempted = true;
#endif
radio_nonrx_start = now; // bounded retry cadence if recovery is unsupported
}
if (outbound) { // waiting for outbound send to complete, or for its one radio retry
if (outbound_radio_retry_pending) {
// The failed send has already returned the chip to RX. Drain a packet
// arriving during backoff before asking to retune; otherwise BUSY can
// keep the retry waiting forever on an unread RxDone interrupt.
checkRecv();
// A cancelled or expired packet must retire even if the chip never
// becomes available for another retune (including an armed OTA reboot).
if (outbound_cancellation != OutboundCancellation::None
|| !isPacketRadioCurrent(outbound) || !allowPacketTransmit(outbound)) {
failOutboundTransmit();
return;
}
if (!millisHasNowPassed(outbound_radio_retry_at)) return;
if (!restoreOutboundTxOverrides()) {
// Release a failed retained retry; restoration itself remains pending
// while normal RX/watchdog service and future queued work continue.
failOutboundTransmit();
return;
}
const auto prepared = _radio->prepareTransmitProfile(outbound->radio_profile,
outbound->radio_reply && outbound->radio_reply_force);
if (prepared == RadioParamApplyResult::BUSY) {
if (outbound_radio_prepare_deadline == 0) {
// A frame already being received may legitimately own either scan
// profile. Allow its longest full-packet airtime plus recovery grace.
// This timer covers only prepare BUSY, not ordinary CAD backoff.
uint32_t airtime = validPrepareBusyAirtime(_radio->getProfileAirtime(0, MAX_TRANS_UNIT));
if (isDualRadioActive()) {
const uint32_t other = validPrepareBusyAirtime(_radio->getProfileAirtime(1, MAX_TRANS_UNIT));
if (other > airtime) airtime = other;
}
const uint32_t timeout = airtime * 3 / 2 + RADIO_PREPARE_BUSY_GRACE_MS;
outbound_radio_prepare_deadline = _ms->getMillis() + timeout;
if (outbound_radio_prepare_deadline == 0) outbound_radio_prepare_deadline = 1;
} else if (millisHasNowPassed(outbound_radio_prepare_deadline)) {
// Retaining outbound disables the non-RX watchdog, and this backoff
// bypasses the liveness watchdog. Release ownership before recovery
// so even a failed reset cannot pin all later queue entries forever.
_err_flags |= ERR_EVENT_RADIO_WATCHDOG;
MESH_DEBUG_PRINTLN("%s Dispatcher: radio retry profile remained busy", getLogDateTime());
failOutboundTransmit();
_radio->recoverRadio(true);
radio_nonrx_start = _ms->getMillis();
return;
}
outbound_radio_retry_at = futureMillis(10);
return;
}
outbound_radio_prepare_deadline = 0;
if (prepared == RadioParamApplyResult::FAILED || !isPacketRadioCurrent(outbound)
|| !allowPacketTransmit(outbound)) {
MESH_DEBUG_PRINTLN("%s Dispatcher::loop(): radio retry packet no longer allowed, type=%u",
getLogDateTime(), (uint32_t)outbound->getPayloadType());
failOutboundTransmit();
} else {
uint32_t retry_delay;
if (!isTransmitChannelReady(outbound, retry_delay)) {
outbound_radio_retry_at = futureMillis(retry_delay);
return;
}
outbound_radio_retry_pending = false;
if (!allowPacketTransmit(outbound) || !startOutboundTransmit()) failOutboundTransmit();
else return;
}
} else if (_radio->isSendComplete()) {
long t = _ms->getMillis() - outbound_start;
total_air_time += t;
if (outbound->getPayloadType() == PAYLOAD_TYPE_OTA) {
ota_tx_finished_at = _ms->getMillis();
ota_tx_airtime = t > 0 ? (uint32_t)t : 0;
}
//Serial.print(" airtime="); Serial.println(t);
updateTxBudget();
if (t > tx_budget_ms) {
tx_budget_ms = 0;
} else {
tx_budget_ms -= t;
}
if (tx_budget_ms < MIN_TX_BUDGET_RESERVE_MS) {
float duty_cycle = 1.0f / (1.0f + getAirtimeBudgetFactor());
unsigned long needed = MIN_TX_BUDGET_RESERVE_MS - tx_budget_ms;
next_tx_time = futureMillis((unsigned long)(needed / duty_cycle));
} else {
next_tx_time = _ms->getMillis();
}
_radio->onSendFinished();
// RX can occur entirely between two loop observations during a TX
// burst. Start a fresh RX recovery allowance after each successful TX;
// otherwise a later CAD pause can trip the eight-second watchdog using
// the start time of the whole healthy burst.
radio_nonrx_start = _ms->getMillis();
restoreOutboundTxOverrides();
logTx(outbound, 2 + outbound->getPathByteLen() + outbound->payload_len);
onSendComplete(outbound);
if (auto* p = _radio->profiles()) ++p->tx_packets[outbound->radio_profile];
if (outbound->isRouteFlood()) {
n_sent_flood++;
} else {
n_sent_direct++;
}
releasePacket(outbound); // return to pool
outbound = NULL;
outbound_radio_retry_pending = false;
outbound_radio_retry_used = false;
outbound_cancellation = OutboundCancellation::None;
} else if (millisHasNowPassed(outbound_expiry)) {
MESH_DEBUG_PRINTLN("%s Dispatcher::loop(): WARNING: outbound packed send timed out!", getLogDateTime());
_radio->onSendFinished();
restoreOutboundTxOverrides();
_err_flags |= ERR_EVENT_RADIO_WATCHDOG;
const bool recovered = _radio->recoverRadio(true);
recovered_this_loop = true;
if (!recovered) {
MESH_DEBUG_PRINTLN("%s Dispatcher::loop(): WARNING: hard radio recovery after TX timeout failed!",
getLogDateTime());
}
if (recovered && scheduleOutboundRadioRetry()) return;
failOutboundTransmit();
} else {
return; // can't do any more radio activity until send is complete or timed out
}
// going back into receive mode now...
const int agc_interval = getAGCResetInterval();
if (agc_interval > 0) {
next_agc_reset_time = futureMillis(agc_interval);
armed_agc_reset_interval = agc_interval;
agc_reset_armed = true;
} else {
armed_agc_reset_interval = 0;
agc_reset_armed = false;
}
}
const int agc_interval = getAGCResetInterval();
if (agc_interval <= 0) {
armed_agc_reset_interval = 0;
agc_reset_armed = false;
} else if (!agc_reset_armed
|| agc_interval != armed_agc_reset_interval) {
// MyMesh loads persisted preferences after Dispatcher::begin(). Arm from
// the first loop so the real configured interval is used, and so a freshly
// initialized radio/frontend gets one full settle interval before the
// first warm-sleep AGC reset. A live interval change also gets a complete
// new interval instead of inheriting the old setting's deadline.
next_agc_reset_time = futureMillis(agc_interval);
armed_agc_reset_interval = agc_interval;
agc_reset_armed = true;
} else if (millisHasNowPassed(next_agc_reset_time)) {
_radio->resetAGC();
next_agc_reset_time = futureMillis(agc_interval);
armed_agc_reset_interval = agc_interval;
}
// check inbound (delayed) queue
{
const uint32_t now = _ms->getMillis();
uint32_t next_inbound;
// Packet managers with deadline support avoid a full priority scan while
// every delayed packet is still in the future. Legacy managers retain the
// original getNextInbound() behavior.
if (!_mgr->getNextInboundTime(now, next_inbound)
|| (int32_t)(next_inbound - now) <= 0) {
Packet* pkt = _mgr->getNextInbound(now);
if (pkt) {
processRecvPacket(pkt);
}
}
}
checkRecv();
// Count only successful recvRaw() delivery as RX activity. Check after
// draining RX and completing TX, before starting the next queued transmit:
// a packet arriving at the deadline prevents unnecessary recovery, and a
// continuous TX queue cannot hide a deaf receiver. The driver defers any
// recovery that would interrupt a pending packet or an active reception.
if (_radio->isInRecvMode() && outbound == NULL && !recovered_this_loop) {
const unsigned long now = _ms->getMillis();
uint32_t soft_liveness_ms = RADIO_LIVENESS_SOFT_MS;
#ifdef WITH_MQTT_BRIDGE
const uint32_t configured_watchdog_ms = getRadioWatchdogMillis();
if (configured_watchdog_ms > 0) soft_liveness_ms = configured_watchdog_ms;
#endif
#ifdef RADIO_LIVENESS_SOFT_ONLY
if (soft_liveness_ms > 0 && (uint32_t)(now - rx_watchdog_window_start) >= soft_liveness_ms) {
_err_flags |= ERR_EVENT_RADIO_WATCHDOG;
MESH_DEBUG_PRINTLN("Radio watchdog: no packet received in observation window for %lu ms, state=%d, soft recovery",
(unsigned long)(uint32_t)(now - rx_watchdog_window_start), _radio->getRadioState());
_radio->recoverRadio(false);
// This radio has no independently resettable RF peripheral. Allow a
// full observation interval between attempts on a quiet mesh.
rx_watchdog_window_start = now;
}
#else
uint32_t hard_liveness_ms = RADIO_LIVENESS_HARD_MS;
if (soft_liveness_ms > hard_liveness_ms / 2) {
hard_liveness_ms = soft_liveness_ms <= 0x7FFFFFFFUL
? soft_liveness_ms * 2 : 0xFFFFFFFFUL;
}
const RadioRecoveryAction action = radio_liveness.poll(
now, soft_liveness_ms, hard_liveness_ms);
if (action != RadioRecoveryAction::NONE) {
const bool hard = action == RadioRecoveryAction::HARD;
_err_flags |= ERR_EVENT_RADIO_WATCHDOG;
MESH_DEBUG_PRINTLN("Radio watchdog: no packet received in observation window for %lu ms, state=%d, %s recovery",
(unsigned long)(uint32_t)(now - radio_liveness.windowStart()), _radio->getRadioState(),
hard ? "hard" : "soft");
const bool recovered = _radio->recoverRadio(hard);
if (hard) radio_liveness.noteHardRecoveryResult(now, recovered);
}
#endif
}
checkSend();
releaseDroppedOutbound();
}
void Dispatcher::releaseDroppedOutbound() {
Packet* dropped;
while ((dropped = _mgr->getNextDroppedOutbound()) != NULL) {
onSendFail(dropped);
releasePacket(dropped);
}
}
bool Dispatcher::tryParsePacket(Packet* pkt, const uint8_t* raw, int len) {
if (pkt == NULL || raw == NULL || len < 2 || len > MAX_TRANS_UNIT) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): packet length invalid, len=%d", getLogDateTime(), len);
return false;
}
int i = 0;
pkt->tx_cr = 0;
pkt->flood_retry_policy = FLOOD_RETRY_POLICY_DEFAULT;
pkt->radio_profile = pkt->radio_origin = _radio->receiveProfile();
pkt->radio_bound = false;
pkt->radio_local = false;
pkt->tx_radio = RADIO_TX_AUTO;
pkt->radio_reply = pkt->radio_reply_force = false;
pkt->radio_generation = pkt->radio_origin_generation = _radio->receiveProfileGeneration();
pkt->header = raw[i++];
if (pkt->getPayloadVer() > PAYLOAD_VER_1) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): unsupported packet version", getLogDateTime());
return false;
}
if (pkt->hasTransportCodes()) {
if (i + (int)sizeof(pkt->transport_codes) + 1 > len) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): incomplete transport header, len=%d", getLogDateTime(), len);
return false;
}
memcpy(&pkt->transport_codes[0], &raw[i], 2); i += 2;
memcpy(&pkt->transport_codes[1], &raw[i], 2); i += 2;
} else {
pkt->transport_codes[0] = pkt->transport_codes[1] = 0;
}
if (i >= len) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): missing path header, len=%d", getLogDateTime(), len);
return false;
}
pkt->path_len = raw[i++];
uint8_t path_mode = pkt->path_len >> 6; // upper 2 bits (legacy firmware: 00)
if (path_mode == 3) { // Reserved for future
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): unsupported path mode: 3", getLogDateTime());
return false;
}
uint8_t path_byte_len = (pkt->path_len & 63) * pkt->getPathHashSize();
if (path_byte_len > MAX_PATH_SIZE || path_byte_len > len - i) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): partial or corrupt packet received, len=%d", getLogDateTime(), len);
return false;
}
memcpy(pkt->path, &raw[i], path_byte_len); i += path_byte_len;
pkt->payload_len = len - i; // payload is remainder
if (pkt->payload_len > sizeof(pkt->payload)) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): packet payload too big, payload_len=%d", getLogDateTime(), (uint32_t)pkt->payload_len);
return false;
}
memcpy(pkt->payload, &raw[i], pkt->payload_len);
return true; // success
}
void Dispatcher::checkRecv() {
if (_radio->isCarrierWaveActive()) return;
Packet* pkt;
float score;
float snr;
float rssi;
uint32_t air_time;
{
uint8_t raw[MAX_TRANS_UNIT+1];
int len = _radio->recvRaw(raw, MAX_TRANS_UNIT);
if (len > 0) {
// A successful driver read proves RX works even if application parsing
// rejects the packet or the packet pool has no space. IRQs alone do not.
#ifdef RADIO_LIVENESS_SOFT_ONLY
rx_watchdog_window_start = _ms->getMillis();
#else
radio_liveness.noteReceive(_ms->getMillis());
#endif
snr = _radio->getLastSNR();
rssi = _radio->getLastRSSI();
// TRACE and CONTROL are direct-only. Reject a flood framing before
// raw hooks can expose it to USB logs, MQTT, ESP-NOW, or another bridge.
// The header check is deliberate: it also holds when the packet pool is
// exhausted, so an invalid frame never slips out through a raw sink.
if (Packet::violatesRoutePolicy(raw[0])) {
pkt = NULL;
} else {
logRxRaw(snr, rssi, raw, len);
pkt = _mgr->allocNew();
if (pkt == NULL) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(): WARNING: received data, no unused packets available!", getLogDateTime());
} else {
if (tryParsePacket(pkt, raw, len)) {
last_meshcore_recv_millis = _ms->getMillis();
if (auto* p = _radio->profiles()) ++p->rx_packets[pkt->radio_profile];
pkt->_snr = snr * 4.0f;
pkt->_rssi = (int16_t)rssi;
score = _radio->packetScore(snr, len);
air_time = _radio->getEstAirtimeFor(len);
rx_air_time += air_time;
} else {
_mgr->free(pkt); // put back into pool
pkt = NULL;
}
}
}
} else {
pkt = NULL;
}
}
if (pkt) {
#if MESH_PACKET_LOGGING && !MESH_PACKET_LOGGING_COMPACT
if (isUsbLoggingEnabled()) {
logPacketLine("RX", pkt, pkt->getRawLength(), true, score, air_time);
}
#endif
logRx(pkt, pkt->getRawLength(), score); // hook for custom logging
if (pkt->isRouteFlood()) {
n_recv_flood++;
int _delay = calcRxDelayForPacket(pkt, score, air_time);
if (_delay < 50) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(), score delay below threshold (%d)", getLogDateTime(), _delay);
processRecvPacket(pkt); // is below the score delay threshold, so process immediately
} else {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkRecv(), score delay is: %d millis", getLogDateTime(), _delay);
if (_delay > MAX_RX_DELAY_MILLIS) {
_delay = MAX_RX_DELAY_MILLIS;
}
_mgr->queueInbound(pkt, futureMillis(_delay)); // add to delayed inbound queue
}
} else {
n_recv_direct++;
processRecvPacket(pkt);
}
}
_radio->onReceiveProcessed();
}
void Dispatcher::processRecvPacket(Packet* pkt) {
const auto* profiles = _radio->profiles();
if (profiles && (pkt->radio_profile > 1
|| (pkt->radio_profile == 1 && !profiles->enabled())
|| pkt->radio_generation != profiles->generation[pkt->radio_profile])) {
releasePacket(pkt);
return;
}
const uint8_t previous_profile = receive_context_profile;
const uint32_t previous_generation = receive_context_generation;
const bool previous_active = receive_context_active;
receive_context_profile = pkt->radio_profile;
receive_context_generation = pkt->radio_generation;
receive_context_active = true;
DispatcherAction action = onRecvPacket(pkt);
receive_context_profile = previous_profile;
receive_context_generation = previous_generation;
receive_context_active = previous_active;
if (action == ACTION_RELEASE) {
_mgr->free(pkt);
} else if (action == ACTION_MANUAL_HOLD) {
// sub-class is wanting to manually hold Packet instance, and call releasePacket() at appropriate time
} else { // ACTION_RETRANSMIT*
uint8_t priority = (action >> 24) - 1;
uint32_t _delay = action & 0xFFFFFF;
if (!queueOutboundPacket(pkt, priority, _delay)) {
onSendFail(pkt);
releasePacket(pkt);
} else if (pkt->isRouteDirect() && pkt->getPayloadType() == PAYLOAD_TYPE_TRACE) {
onTracePacketQueuedForSend(pkt);
}
}
}
bool Dispatcher::isTransmitChannelReady(const Packet* packet, uint32_t& retry_delay) {
const bool busy = packet != NULL && usePassiveChannelCheck(packet)
? _radio->isReceivingPassive(getRetryInterferenceMargin())
: _radio->isReceiving();
const uint8_t profile = packet ? packet->radio_profile : 0;
uint32_t& profile_busy = profile_cad_busy[profile];
if (busy) {
const uint32_t now = _ms->getMillis();
// A driver retry owns its packet outside the queue, but remains one of
// the packets waiting for the channel's bounded busy allowance.
const int ready_count = _mgr->getOutboundCount(now) + (packet && outbound == packet ? 1 : 0);
cad_busy_start = profile_busy;
if (cad_busy_start == 0) {
cad_busy_start = now;
profile_busy = now;
}
const uint32_t maximum = scaleCADDelayForQueue(
getCADFailMaxDuration(), ready_count, MIN_CAD_FAIL_MAX_DURATION_MS);
if (now - cad_busy_start <= maximum) {
retry_delay = scaleCADDelayForQueue(
getCADFailRetryDelay(), ready_count, MIN_CAD_FAIL_RETRY_DELAY_MS);
return false;
}
_err_flags |= ERR_EVENT_CAD_TIMEOUT;
MESH_DEBUG_PRINTLN("%s Dispatcher: CAD busy max duration reached!", getLogDateTime());
}
cad_busy_start = 0;
profile_busy = 0;
return true;
}
void Dispatcher::checkSend() {
if (_radio->isCarrierWaveActive()) return;
// Keep queued packets intact while recovering a failed per-packet CR
// 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, getOtaPacketSpeedFactor())) {
ota_tx_airtime = 0;
}
uint32_t next_outbound;
if (_mgr->getNextOutboundTime(now, next_outbound)) {
if ((int32_t)(next_outbound - now) > 0) return;
} else if (_mgr->getOutboundCount(now) == 0) {
// Compatibility fallback for custom PacketManager implementations that do
// not provide the optional O(1) deadline query.
return;
}
// Discard work bound to an expired/changed profile before it can be retuned
// onto a different channel. This also retires its retry ownership normally.
Packet* pending = _mgr->peekNextOutbound(now);
if (pending && !isPacketRadioCurrent(pending)) {
outbound = _mgr->getNextOutbound(now);
failOutboundTransmit();
return;
}
if (ota_tx_airtime) {
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) {
// Preserve this packet's priority and retry ownership. A short deferral
// lets ordinary traffic run and makes live speed changes take effect
// within 100 ms, even when the previous speed was much slower.
const uint32_t remaining = gap - elapsed;
int candidates = _mgr->getOutboundTotal();
do {
if (candidates-- <= 0 || !_mgr->deferOutboundForPacing(pending, now,
now + (remaining > 100 ? 100 : remaining))) return;
pending = _mgr->peekNextOutbound(now);
} while (pending && pending->getPayloadType() == PAYLOAD_TYPE_OTA);
if (!pending) return;
if (!isPacketRadioCurrent(pending)) {
outbound = _mgr->getNextOutbound(now);
failOutboundTransmit();
return;
}
}
}
updateTxBudget();
uint32_t est_airtime = pending && _radio->profiles()
? _radio->getProfileAirtime(pending->radio_profile, MAX_TRANS_UNIT)
: _radio->getEstAirtimeFor(MAX_TRANS_UNIT);
if (tx_budget_ms < est_airtime / MIN_TX_BUDGET_AIRTIME_DIV) {
float duty_cycle = 1.0f / (1.0f + getAirtimeBudgetFactor());
unsigned long needed = est_airtime / MIN_TX_BUDGET_AIRTIME_DIV - tx_budget_ms;
next_tx_time = futureMillis((unsigned long)(needed / duty_cycle));
return;
}
if (!millisHasNowPassed(next_tx_time)) return;
// Waiting for airtime credit must leave the receive scanner free to run.
// Only retune once this queue entry is actually eligible to transmit.
if (pending && _radio->profiles()) {
const auto prepared = _radio->prepareTransmitProfile(pending->radio_profile,
pending->radio_reply && pending->radio_reply_force);
if (prepared == RadioParamApplyResult::BUSY) return;
if (prepared == RadioParamApplyResult::FAILED) {
outbound = _mgr->getNextOutbound(now);
failOutboundTransmit();
return;
}
}
uint32_t retry_delay;
if (!isTransmitChannelReady(pending, retry_delay)) {
if (!isDualRadioActive() || !pending || !_mgr->deferOutbound(pending, futureMillis(retry_delay))) {
next_tx_time = futureMillis(retry_delay);
}
return;
}
// Retuning and CAD can advance the clock. Keep the selection timestamp so
// a newly due higher-priority packet cannot skip its own profile, airtime,
// and OTA pacing checks. A manager may also expire/reorder work on peek.
if (pending && _mgr->peekNextOutbound(now) != pending) return;
outbound = _mgr->getNextOutbound(now);
if (outbound) {
outbound_radio_retry_pending = false;
outbound_radio_retry_used = false;
outbound_cancellation = OutboundCancellation::None;
if (!allowPacketTransmit(outbound)) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkSend(): packet no longer allowed, type=%u", getLogDateTime(),
(uint32_t)outbound->getPayloadType());
failOutboundTransmit();
return;
}
if (outbound->getRawLength() > MAX_TRANS_UNIT) {
MESH_DEBUG_PRINTLN("%s Dispatcher::checkSend(): FATAL: Invalid packet queued... too long, len=%d",
getLogDateTime(), outbound->getRawLength());
failOutboundTransmit();
return;
}
if (!startOutboundTransmit()) {
// RadioLib performs its radio-specific cleanup (including LR1110 hard
// recovery for a stuck BUSY command) before returning false. Retain the
// packet and let the firmware make one fresh start without app help.
if (scheduleOutboundRadioRetry()) return;
failOutboundTransmit();
}
}
}
Packet* Dispatcher::obtainNewPacket() {
auto pkt = _mgr->allocNew(); // TODO: zero out all fields
if (pkt == NULL) {
_err_flags |= ERR_EVENT_FULL;
} else {
pkt->payload_len = pkt->path_len = 0;
pkt->_rssi = 0;
pkt->_snr = 0;
pkt->tx_cr = 0;
pkt->flood_retry_policy = FLOOD_RETRY_POLICY_DEFAULT;
pkt->radio_profile = pkt->radio_origin = receive_context_active ? receive_context_profile : 0;
pkt->radio_generation = pkt->radio_origin_generation = receive_context_active ? receive_context_generation : 0;
pkt->radio_bound = false;
pkt->radio_local = !receive_context_active;
pkt->tx_radio = RADIO_TX_AUTO;
pkt->radio_reply = receive_context_active;
pkt->radio_reply_force = false;
}
return pkt;
}
void Dispatcher::releasePacket(Packet* packet) {
_mgr->free(packet);
}
bool Dispatcher::queueOutboundPacket(Packet* packet, uint8_t priority, uint32_t delay_millis) {
if (packet->violatesRoutePolicy()) {
return false;
}
if (!radio_available) {
MESH_DEBUG_PRINTLN("%s Dispatcher::sendPacket(): radio unavailable", getLogDateTime());
return false;
}
if (!Packet::isValidPathLen(packet->path_len) || packet->payload_len > MAX_PACKET_PAYLOAD) {
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);
return false;
}
auto* profiles = _radio->profiles();
if (packet->radio_profile > 1) return false;
if (profiles == nullptr) {
if (packet->tx_radio != RADIO_TX_AUTO
&& !(explicitRadioTxMask(packet->tx_radio, false) & 1)) return false;
return _mgr->queueOutbound(packet, priority, futureMillis(delay_millis));
}
if (packet->radio_bound) {
if (!isPacketRadioCurrent(packet)) return false;
if (profiles && packet->radio_profile != packet->radio_origin
&& !allowRadioProfileCross(packet)) return false;
return _mgr->queueOutbound(packet, priority, futureMillis(delay_millis));
}
if (!packet->radio_local && packet->radio_generation
&& packet->radio_generation != profiles->generation[packet->radio_profile]) return false;
if (packet->radio_reply && packet->tx_radio == RADIO_TX_AUTO) {
packet->tx_radio = profiles->reply_tx;
packet->radio_reply_force = profiles->reply_force;
}
// Locally generated OTA traffic uses the temporary update profile. Replies
// created while processing RX inherit that request's profile instead.
// Keep that origin in RX-only mode too: isolation must drop a disallowed
// transmission instead of silently sending it on the normal channel.
if (packet->tx_radio == RADIO_TX_AUTO && packet->radio_local && packet->getPayloadType() == PAYLOAD_TYPE_OTA
&& profiles->secondary_temporary && !profiles->primary_temporary) {
packet->radio_profile = 1;
}
const uint8_t origin = packet->radio_profile;
uint8_t mask = getTransmitProfileMask(packet);
const uint8_t cross_mask = (uint8_t)(1U << (origin ^ 1));
if ((mask & cross_mask) && !allowRadioProfileCross(packet)) {
mask &= (uint8_t)~cross_mask;
}
if (!mask) return false;
packet->radio_origin = origin;
packet->radio_origin_generation = profiles->generation[origin];
if (!(mask & (1U << origin))) {
onSendFail(packet); // retire any retry reserved on the RX-only profile
packet->radio_profile ^= 1;
}
packet->radio_generation = profiles->generation[packet->radio_profile];
packet->radio_bound = true;
if (!_mgr->queueOutbound(packet, priority, futureMillis(delay_millis))) return false;
if (packet->radio_profile != origin) onRadioProfileCopyQueued(packet, nullptr, priority);
const uint8_t other = packet->radio_profile ^ 1;
if (mask & (1U << other)) {
Packet* copy = obtainNewPacket();
if (copy) {
*copy = *packet;
copy->radio_profile = other;
copy->radio_generation = profiles->generation[other];
if (_mgr->queueOutbound(copy, priority, futureMillis(delay_millis))) {
onRadioProfileCopyQueued(copy, packet, priority);
onTracePacketQueuedForSend(copy);
} else {
onSendFail(copy);
releasePacket(copy);
}
}
}
return true;
}
uint8_t Dispatcher::getTransmitProfileMask(const Packet* packet) const {
const auto* profiles = _radio->profiles();
if (!profiles) return packet->tx_radio == RADIO_TX_AUTO ? 1 : explicitRadioTxMask(packet->tx_radio, false);
if (packet->radio_profile > 1) return 0;
const bool reply = packet->radio_reply && !packet->radio_bound && packet->tx_radio == RADIO_TX_AUTO;
const uint8_t policy = reply ? profiles->reply_tx : packet->tx_radio;
const bool force = packet->radio_reply && (reply ? profiles->reply_force : packet->radio_reply_force);
uint8_t origin = packet->radio_profile;
if (policy == RADIO_TX_AUTO && packet->radio_local && packet->getPayloadType() == PAYLOAD_TYPE_OTA
&& profiles->secondary_temporary && !profiles->primary_temporary) origin = 1;
return policy == RADIO_TX_AUTO ? profiles->transmitMask(origin)
: explicitRadioTxMask(policy, profiles->canTransmit(1, force));
}
uint32_t Dispatcher::getTransmitAirtime(const Packet* packet) const {
const uint8_t mask = getTransmitProfileMask(packet);
uint32_t total = 0;
for (uint8_t profile = 0; profile < 2; ++profile) {
if (mask & (1U << profile))
total += _radio->getProfileAirtime(profile, packet->getRawLength(), packet->tx_cr);
}
return total;
}
bool Dispatcher::isPacketRadioCurrent(const Packet* packet) const {
const auto* p = _radio->profiles();
if (!p || !packet->radio_bound) return true;
const uint8_t target = packet->radio_profile;
const uint8_t origin = packet->radio_origin;
const bool force = packet->radio_reply && packet->radio_reply_force;
return target < 2 && origin < 2 && p->canTransmit(target, force)
&& packet->radio_generation == p->generation[target]
&& packet->radio_origin_generation == p->generation[origin]
&& (packet->tx_radio == RADIO_TX_AUTO ? (origin == target || p->canCross())
: (explicitRadioTxMask(packet->tx_radio, p->canTransmit(1, force)) & (1U << target)) != 0);
}
bool Dispatcher::sendPacket(Packet* packet, uint8_t priority, uint32_t delay_millis) {
if (!queueOutboundPacket(packet, priority, delay_millis)) {
onSendFail(packet);
releasePacket(packet);
return false;
}
if (packet->isRouteDirect() && packet->getPayloadType() == PAYLOAD_TYPE_TRACE) {
onTracePacketQueuedForSend(packet);
}
return true;
}
// Utility function -- handles the case where millis() wraps around back to zero
// 2's complement arithmetic will handle any unsigned subtraction up to HALF the word size (32-bits in this case)
bool Dispatcher::millisHasNowPassed(unsigned long timestamp) const {
// millis is a 32-bit counter even on native hosts where long is 64 bits.
return (int32_t)((uint32_t)_ms->getMillis() - (uint32_t)timestamp) > 0;
}
unsigned long Dispatcher::futureMillis(int millis_from_now) const {
return _ms->getMillis() + millis_from_now;
}
}