From 9bc6060831a0ad6064911836518e143d8a7c5ab9 Mon Sep 17 00:00:00 2001 From: Scott Powell Date: Fri, 4 Sep 2026 15:13:17 +1000 Subject: [PATCH] * new: REQ_TYPE_SUBSCRIBE and REQ_TYPE_UNSUBSCRIBE --- examples/simple_sensor/SensorMesh.cpp | 124 ++++++++++++++++++++++---- examples/simple_sensor/SensorMesh.h | 9 +- src/helpers/ClientACL.h | 5 ++ 3 files changed, 116 insertions(+), 22 deletions(-) diff --git a/examples/simple_sensor/SensorMesh.cpp b/examples/simple_sensor/SensorMesh.cpp index 287a3c1bf..0a78f5daf 100644 --- a/examples/simple_sensor/SensorMesh.cpp +++ b/examples/simple_sensor/SensorMesh.cpp @@ -55,6 +55,8 @@ #define REQ_TYPE_GET_TELEMETRY_DATA 0x03 #define REQ_TYPE_GET_AVG_MIN_MAX 0x04 #define REQ_TYPE_GET_ACCESS_LIST 0x05 +#define REQ_TYPE_SUBSCRIBE 0x10 +#define REQ_TYPE_UNSUBSCRIBE 0x11 #define RESP_SERVER_LOGIN_OK 0 // response to ANON_REQ @@ -74,6 +76,8 @@ static File openAppend(FILESYSTEM* _fs, const char* fname) { #endif } +/* --------------------- Cayenne LPP helpers ----------------------------*/ + static uint8_t getDataSize(uint8_t type) { switch (type) { case LPP_GPS: @@ -171,7 +175,53 @@ static uint8_t putFloat(uint8_t * dest, float value, uint8_t size, uint32_t mult return size; } -uint8_t SensorMesh::handleRequest(uint8_t perms, uint32_t sender_timestamp, uint8_t req_type, uint8_t* payload, size_t payload_len) { +static float findTelemValue(const uint8_t* buf, uint8_t size, uint8_t channel, uint8_t type) { + uint8_t i = 0; + + while (i + 2 < size) { + // Get channel # + uint8_t ch = buf[i++]; + // Get data type + uint8_t t = buf[i++]; + uint8_t sz = getDataSize(t); + + if (ch == channel && t == type) { + return getFloat(&buf[i], sz, getMultiplier(t), isSigned(t)); + } + i += sz; // skip + } + return 0.0f; // not found +} + +/* ------------------ end Cayenne LPP helpers ----------------------*/ + +bool SensorMesh::telemHasChanged(const uint8_t* min_deltas, uint8_t min_deltas_len) { + if (telemetry.getSize() != prev_telem_size) return true; + + auto buf = telemetry.getBuffer(); + uint8_t size = telemetry.getSize(); + uint8_t i = 0; + + while (i + 2 < size) { + // Get channel # + uint8_t ch = buf[i++]; + // Get data type + uint8_t t = buf[i++]; + uint8_t sz = getDataSize(t); + + float v = getFloat(&buf[i], sz, getMultiplier(t), isSigned(t)); + float pv = getFloat(&prev_telem[i], sz, getMultiplier(t), isSigned(t)); + float min_delta = findTelemValue(min_deltas, min_deltas_len, ch, t); + if (abs(v - pv) > min_delta) return true; // Yes, has changed + + i += sz; // skip + } + return false; // no changes +} + +uint8_t SensorMesh::handleRequest(ClientInfo* from, uint32_t sender_timestamp, uint8_t req_type, uint8_t* payload, size_t payload_len) { + uint8_t perms = from->isAdmin() ? 0xFF : from->permissions; + memcpy(reply_data, &sender_timestamp, 4); // reflect sender_timestamp back in response packet (kind of like a 'tag') if (req_type == REQ_TYPE_GET_TELEMETRY_DATA) { // allow all @@ -235,6 +285,29 @@ uint8_t SensorMesh::handleRequest(uint8_t perms, uint32_t sender_timestamp, uint return ofs; } } + if (req_type == REQ_TYPE_SUBSCRIBE && payload_len >= 2 && (perms & PERM_ACL_ROLE_MASK) >= PERM_ACL_READ_ONLY) { + uint8_t reserved = payload[0]; + RegionEntry* r; + if (recv_pkt_region && !recv_pkt_region->isWildcard()) { // use request scope + r = recv_pkt_region; + } else { // use default scope + r = region_map.getDefaultRegion(); + } + from->extra.sensor.scope_region_id = r ? r->id : 0; + from->extra.sensor.min_deltas_len = payload[1]; + // NOTE: curr impl truncates LPP min_diffs spec (re-do if better impl is needed) + memcpy(from->extra.sensor.min_deltas, &payload[2], min(sizeof(from->extra.sensor.min_deltas), (size_t)payload[1])); + + getRNG()->random(&reply_data[4], 2); // just some entropy for better packet-hash uniqueness + strcpy((char *)&reply_data[6], r ? r->name : ""); // reply with name of scope that will be used + return 6 + strlen((char *)&reply_data[6]); + } + if (req_type == REQ_TYPE_UNSUBSCRIBE && (perms & PERM_ACL_ROLE_MASK) >= PERM_ACL_READ_ONLY) { + from->extra.sensor.scope_region_id = 0; + reply_data[4] = 0; // success + getRNG()->random(&reply_data[5], 3); // just some entropy for better packet-hash uniqueness + return 8; + } return 0; // unknown command } @@ -548,7 +621,7 @@ void SensorMesh::onPeerDataRecv(mesh::Packet* packet, uint8_t type, int sender_i memcpy(×tamp, data, 4); if (timestamp > from->last_timestamp) { // prevent replay attacks - uint8_t reply_len = handleRequest(from->isAdmin() ? 0xFF : from->permissions, timestamp, data[4], &data[5], len - 5); + uint8_t reply_len = handleRequest(from, timestamp, data[4], &data[5], len - 5); if (reply_len == 0) return; // invalid command from->last_timestamp = timestamp; @@ -726,6 +799,7 @@ SensorMesh::SensorMesh(mesh::MainBoard& board, mesh::Radio& radio, mesh::Millise num_alert_tasks = 0; set_radio_at = revert_radio_at = 0; recv_pkt_region = NULL; + prev_telem_size = 0; // defaults _prefs.airtime_factor = 1.0; @@ -916,23 +990,7 @@ void SensorMesh::formatPacketStatsReply(char *reply) { } float SensorMesh::getTelemValue(uint8_t channel, uint8_t type) { - auto buf = telemetry.getBuffer(); - uint8_t size = telemetry.getSize(); - uint8_t i = 0; - - while (i + 2 < size) { - // Get channel # - uint8_t ch = buf[i++]; - // Get data type - uint8_t t = buf[i++]; - uint8_t sz = getDataSize(t); - - if (ch == channel && t == type) { - return getFloat(&buf[i], sz, getMultiplier(t), isSigned(t)); - } - i += sz; // skip - } - return 0.0f; // not found + return findTelemValue(telemetry.getBuffer(), telemetry.getSize(), channel, type); } bool SensorMesh::getGPS(uint8_t channel, float& lat, float& lon, float& alt) { @@ -982,6 +1040,34 @@ void SensorMesh::loop() { // query other sensors -- target specific sensors.querySensors(0xFF, telemetry); // allow all telemetry permissions + // compare with previous telemetry, check if any deltas are greater than subscriber minimums + for (int i = 0; i < acl.getNumClients(); i++) { + auto c = acl.getClientByIdx(i); + if (c->permissions == 0 || c->extra.sensor.scope_region_id == 0) continue; // skip deleted entries, or Not subscribed to deltas + RegionEntry* r = region_map.findById(c->extra.sensor.scope_region_id); + if (r == NULL) continue; // unknown region scope + if (telemHasChanged(c->extra.sensor.min_deltas, c->extra.sensor.min_deltas_len)) { + TransportKey scope; + if (region_map.getTransportKeysFor(*r, &scope, 1) > 0) { + uint8_t tlen = telemetry.getSize(); + uint32_t timestamp = getRTCClock()->getCurrentTimeUnique(); // this will be an unknown 'tag' to the client + memcpy(reply_data, ×tamp, 4); + memcpy(&reply_data[4], telemetry.getBuffer(), tlen); + + mesh::Packet* reply = createDatagram(PAYLOAD_TYPE_RESPONSE, c->id, c->shared_secret, reply_data, 4 + tlen); + if (reply) { + if (c->out_path_len != OUT_PATH_UNKNOWN) { // we have an out_path, so send DIRECT + sendDirect(reply, c->out_path, c->out_path_len, 0); + } else { + sendFloodScoped(scope, reply, 0, _prefs.path_hash_mode + 1); + } + } + } + } + } + memcpy(prev_telem, telemetry.getBuffer(), telemetry.getSize()); // save snapshot for next compare cycle + prev_telem_size = telemetry.getSize(); + onSensorDataRead(); last_read_time = curr; diff --git a/examples/simple_sensor/SensorMesh.h b/examples/simple_sensor/SensorMesh.h index 845402bae..7ea87f206 100644 --- a/examples/simple_sensor/SensorMesh.h +++ b/examples/simple_sensor/SensorMesh.h @@ -96,12 +96,12 @@ protected: bool getGPS(uint8_t channel, float& lat, float& lon, float& alt); // alerts - enum AlertPriority { LOW_PRI_ALERT, HIGH_PRI_ALERT }; + enum AlertPriority : uint8_t { LOW_PRI_ALERT, HIGH_PRI_ALERT }; struct Trigger { uint32_t timestamp; - AlertPriority pri; uint32_t expected_acks[4]; + AlertPriority pri; int8_t curr_contact_idx; uint8_t attempt; unsigned long send_expiry; @@ -145,6 +145,8 @@ private: uint8_t reply_data[MAX_PACKET_PAYLOAD]; unsigned long dirty_contacts_expiry; CayenneLPP telemetry; + uint8_t prev_telem_size; + uint8_t prev_telem[MAX_PACKET_PAYLOAD - 4]; TransportKeyStore key_store; RegionMap region_map; RegionEntry* recv_pkt_region; @@ -159,8 +161,9 @@ private: uint8_t pending_sf; uint8_t pending_cr; + bool telemHasChanged(const uint8_t* min_deltas, uint8_t min_deltas_len); uint8_t handleLoginReq(const mesh::Identity& sender, const uint8_t* secret, uint32_t sender_timestamp, const uint8_t* data, bool is_flood); - uint8_t handleRequest(uint8_t perms, uint32_t sender_timestamp, uint8_t req_type, uint8_t* payload, size_t payload_len); + uint8_t handleRequest(ClientInfo* from, uint32_t sender_timestamp, uint8_t req_type, uint8_t* payload, size_t payload_len); mesh::Packet* createSelfAdvert(); void sendAlert(const ClientInfo* c, Trigger* t); diff --git a/src/helpers/ClientACL.h b/src/helpers/ClientACL.h index e06544647..9fde35ecb 100644 --- a/src/helpers/ClientACL.h +++ b/src/helpers/ClientACL.h @@ -28,6 +28,11 @@ struct ClientInfo { unsigned long ack_timeout; uint8_t push_failures; } room; + struct { + uint16_t scope_region_id; // scope to use when sending telemetry to this client/subscriber + uint8_t min_deltas_len; + uint8_t min_deltas[14]; // LPP encoded + } sensor; } extra; bool isAdmin() const { return (permissions & PERM_ACL_ROLE_MASK) == PERM_ACL_ADMIN; }