Repository navigation
Expand file tree
/
Copy pathgroup_sync.cpp
More file actions
577 lines (514 loc) · 18.3 KB
/
Copy pathgroup_sync.cpp
File metadata and controls
577 lines (514 loc) · 18.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
#include "group_sync.h"
#if FEATURE_ESPNOW
#include <Arduino.h>
#include <WiFi.h>
#include <esp_now.h>
#include <esp_wifi.h>
#include <string.h>
#include "settings.h"
#include "display.h"
#include "signal_light.h"
namespace group_sync {
namespace {
// Broadcast rather than unicast: no discovery, no peer list to keep in sync
// across reboots, and no six-peer encryption cap. The group name in every
// message is what separates one set of lights from another.
const uint8_t kBroadcast[6] = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF};
enum MessageType : uint8_t {
MSG_LEVEL = 1,
MSG_SETTINGS = 2,
MSG_MODE = 3,
MSG_SELF_TEST = 4,
MSG_CONFIG = 5, // addressed at a single unit
};
// Packed and fixed-width so the format does not depend on how a particular
// compiler lays out a struct. Well inside ESP-NOW's 250-byte limit.
struct __attribute__((packed)) Message {
uint8_t magic[2]; // 'd', 'L'
uint8_t version;
uint8_t type;
uint32_t group; // hash of the group name
// All zeroes means everyone in the group. Anything else is one unit's MAC,
// which is how per-unit settings are configured from elsewhere.
uint8_t target[6];
// A union rather than variable-length payloads, so every message is the
// same size and a wrong length is enough to reject a foreign packet.
union {
struct {
int16_t centiDb;
uint8_t zones; // what the sender lights for
uint8_t follows; // whether the sender tracks the group at all
char name[GROUP_NAME_MAX + 1];
uint8_t inactiveLevel;
uint8_t displayBrightness;
uint8_t ip[4]; // where the sender's web interface lives
} level;
struct {
uint8_t dbMin;
uint8_t dbMax;
uint8_t brightness;
// Group-wide, unlike zones: it decides how the group's readings become
// one number, and two units disagreeing about it would quietly show
// different colours.
uint8_t combine;
} settings;
struct {
uint8_t mode; // signal_light::Mode
uint8_t r, g, b;
} mode;
struct {
uint8_t zones;
uint8_t follows;
uint8_t inactiveLevel;
uint8_t displayBrightness;
char name[GROUP_NAME_MAX + 1];
} config;
} body;
};
struct Peer {
uint8_t mac[6];
float levelDb;
uint8_t zones;
bool followsGroup;
uint8_t inactiveLevel;
uint8_t displayBrightness;
uint8_t ip[4];
char name[GROUP_NAME_MAX + 1];
uint32_t heardMs;
bool used;
};
bool running = false;
uint32_t groupId = 0;
uint32_t lastSentMs = 0;
// What actually prevents an applied change echoing around the group is the
// shape of the code: publishing happens only where a change originates -
// remote_control and web_control - and the apply path below calls settings
// and signal_light directly, never a publish. This flag is a second line of
// defence for the day someone wires publishing into one of those setters. It
// is deliberately not load-bearing today, and a test cannot reach it.
bool applyingRemote = false;
// This unit's label on the air, filled in at begin().
char localNameBuf[GROUP_NAME_MAX + 1] = {0};
uint8_t localMac[6] = {0};
uint32_t localIp = 0;
// Commands are repeated a few times because broadcast is unacknowledged. The
// repeats are paced by tick() rather than a delay, so a held remote key
// cannot stall the loop.
Message repeatMsg;
uint8_t repeatsLeft = 0;
uint32_t nextRepeatMs = 0;
// Commands land here from the radio callback and are applied by tick() on the
// main thread. Only the most recent of each kind is kept: they are absolute
// states, not increments, so an older one has nothing to contribute.
struct Inbox {
bool haveSettings;
uint8_t dbMin, dbMax, brightness, combine;
bool haveMode;
uint8_t mode;
uint32_t color;
bool haveSelfTest;
bool haveConfig;
uint8_t zones, inactiveLevel, displayBrightness;
bool follows;
char name[GROUP_NAME_MAX + 1];
};
Inbox inbox = {};
// Written by the receive callback, which runs from the WiFi task. Only ever
// touched under this lock; the table is small enough that holding it for a
// linear scan costs nothing.
portMUX_TYPE peerLock = portMUX_INITIALIZER_UNLOCKED;
Peer peers[GROUP_MAX_PEERS] = {};
// FNV-1a. Any stable hash would do; this one is four lines and has no
// dependencies. Collisions would merge two groups, which is a cosmetic
// problem for a classroom light rather than a correctness one.
uint32_t hashGroup(const char* name) {
uint32_t h = 2166136261u;
for (const char* c = name; *c != '\0'; c++) {
h ^= uint8_t(*c);
h *= 16777619u;
}
return h;
}
bool sameMac(const uint8_t* a, const uint8_t* b) { return memcmp(a, b, 6) == 0; }
void macToHex(const uint8_t* mac, char* out) {
static const char kHex[] = "0123456789ABCDEF";
for (int i = 0; i < 6; i++) {
out[i * 2] = kHex[mac[i] >> 4];
out[i * 2 + 1] = kHex[mac[i] & 0xF];
}
out[12] = '\0';
}
// Returns false on anything that is not twelve hex digits, so a malformed id
// addresses nobody rather than everybody.
bool hexToMac(const char* hex, uint8_t* mac) {
if (hex == nullptr || strlen(hex) != 12) return false;
for (int i = 0; i < 6; i++) {
uint8_t byte = 0;
for (int n = 0; n < 2; n++) {
const char c = hex[i * 2 + n];
const int v = (c >= '0' && c <= '9') ? c - '0'
: (c >= 'A' && c <= 'F') ? c - 'A' + 10
: (c >= 'a' && c <= 'f') ? c - 'a' + 10
: -1;
if (v < 0) return false;
byte = uint8_t((byte << 4) | v);
}
mac[i] = byte;
}
return true;
}
// Common header, plus the two conditions under which nothing may be sent.
bool fill(Message& msg, uint8_t type) {
if (!running || applyingRemote) return false;
msg.magic[0] = 'd';
msg.magic[1] = 'L';
msg.version = GROUP_PROTOCOL_VERSION;
msg.type = type;
msg.group = groupId;
memset(msg.target, 0, sizeof(msg.target)); // everyone unless addressed
memset(&msg.body, 0, sizeof(msg.body));
return true;
}
void transmit(const Message& msg) {
esp_now_send(kBroadcast, reinterpret_cast<const uint8_t*>(&msg), sizeof(msg));
}
// One copy now; the rest are paced out by tick(). A newer command replaces
// any still-pending repeats, since both carry absolute state and the older
// one has nothing left to say.
void send(const Message& msg, uint8_t repeats = 1) {
transmit(msg);
if (repeats > 1) {
repeatMsg = msg;
repeatsLeft = repeats - 1;
nextRepeatMs = millis() + GROUP_COMMAND_GAP_MS;
}
}
void onReceive(const uint8_t* mac, const uint8_t* data, int len) {
if (len != int(sizeof(Message))) return;
Message msg;
memcpy(&msg, data, sizeof(msg));
if (msg.magic[0] != 'd' || msg.magic[1] != 'L') return;
if (msg.version != GROUP_PROTOCOL_VERSION) return;
if (msg.group != groupId) return;
// An addressed message is for exactly one unit; everything else is for the
// whole group.
static const uint8_t kEveryone[6] = {0, 0, 0, 0, 0, 0};
if (!sameMac(msg.target, kEveryone) && !sameMac(msg.target, localMac)) return;
const uint32_t now = millis();
if (msg.type != MSG_LEVEL) {
// Everything else is a command. Record it and get out; applying it here
// would mean touching settings and the LEDs from the WiFi task.
portENTER_CRITICAL(&peerLock);
switch (msg.type) {
case MSG_SETTINGS:
inbox.dbMin = msg.body.settings.dbMin;
inbox.dbMax = msg.body.settings.dbMax;
inbox.brightness = msg.body.settings.brightness;
inbox.combine = msg.body.settings.combine;
inbox.haveSettings = true;
break;
case MSG_MODE:
inbox.mode = msg.body.mode.mode;
inbox.color = (uint32_t(msg.body.mode.r) << 16) | (uint32_t(msg.body.mode.g) << 8) |
msg.body.mode.b;
inbox.haveMode = true;
break;
case MSG_SELF_TEST:
inbox.haveSelfTest = true;
break;
case MSG_CONFIG:
inbox.zones = msg.body.config.zones;
inbox.follows = msg.body.config.follows != 0;
inbox.inactiveLevel = msg.body.config.inactiveLevel;
inbox.displayBrightness = msg.body.config.displayBrightness;
memcpy(inbox.name, msg.body.config.name, sizeof(inbox.name));
inbox.name[GROUP_NAME_MAX] = '\0';
inbox.haveConfig = true;
break;
default:
break;
}
portEXIT_CRITICAL(&peerLock);
return;
}
const float level = float(msg.body.level.centiDb) / 100.0f;
portENTER_CRITICAL(&peerLock);
int slot = -1;
int oldest = -1;
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (peers[i].used && sameMac(peers[i].mac, mac)) {
slot = i;
break;
}
if (!peers[i].used && slot < 0) slot = i;
if (peers[i].used && (oldest < 0 || peers[i].heardMs < peers[oldest].heardMs)) oldest = i;
}
// A full table evicts whoever has been quiet longest, so a group larger
// than the cap degrades to the most recently active members rather than
// ignoring newcomers entirely.
if (slot < 0) slot = oldest;
if (slot >= 0) {
memcpy(peers[slot].mac, mac, 6);
peers[slot].levelDb = level;
peers[slot].zones = msg.body.level.zones;
peers[slot].followsGroup = msg.body.level.follows != 0;
memcpy(peers[slot].name, msg.body.level.name, GROUP_NAME_MAX);
peers[slot].name[GROUP_NAME_MAX] = '\0';
peers[slot].inactiveLevel = msg.body.level.inactiveLevel;
peers[slot].displayBrightness = msg.body.level.displayBrightness;
memcpy(peers[slot].ip, msg.body.level.ip, 4);
peers[slot].heardMs = now;
peers[slot].used = true;
}
portEXIT_CRITICAL(&peerLock);
}
} // namespace
bool begin() {
// Re-entrant on purpose: a station reconnect can land on a different
// channel, which invalidates the broadcast peer, so the network layer
// calls this again rather than restarting the unit.
if (running) esp_now_deinit();
running = false;
lastSentMs = 0;
groupId = hashGroup(settings::groupName());
memset(peers, 0, sizeof(peers));
repeatsLeft = 0;
WiFi.macAddress(localMac);
strncpy(localNameBuf, settings::unitName(), GROUP_NAME_MAX);
localNameBuf[GROUP_NAME_MAX] = '\0';
if (localNameBuf[0] == '\0') {
snprintf(localNameBuf, sizeof(localNameBuf), "%02X%02X", localMac[4], localMac[5]);
}
if (settings::groupName()[0] == '\0') {
Serial.println(F("group: no group name set, working alone"));
return false;
}
if (esp_now_init() != ESP_OK) {
Serial.println(F("group: esp_now_init failed"));
return false;
}
esp_now_register_recv_cb(onReceive);
// The peer has to sit on whichever interface WiFi actually brought up, and
// channel 0 means "whatever channel we are on", which keeps working if the
// station reconnects to an access point on a different one.
esp_now_peer_info_t peer = {};
memcpy(peer.peer_addr, kBroadcast, 6);
peer.channel = 0;
peer.encrypt = false;
peer.ifidx = (WiFi.getMode() & WIFI_MODE_STA) ? WIFI_IF_STA : WIFI_IF_AP;
if (esp_now_add_peer(&peer) != ESP_OK) {
Serial.println(F("group: could not add the broadcast peer"));
esp_now_deinit();
return false;
}
running = true;
Serial.printf("group: '%s' on channel %u\n", settings::groupName(), channel());
return true;
}
bool active() { return running; }
uint8_t channel() {
uint8_t primary = 0;
wifi_second_chan_t second;
esp_wifi_get_channel(&primary, &second);
return primary;
}
uint8_t peerCount() {
uint8_t n = 0;
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (peers[i].used) n++;
}
portEXIT_CRITICAL(&peerLock);
return n;
}
const char* localName() { return localNameBuf; }
void setAddress(uint32_t ipv4) { localIp = ipv4; }
bool publishPeerConfig(const char* id, const PeerConfig& config) {
uint8_t mac[6];
if (!hexToMac(id, mac)) return false;
Message msg;
if (!fill(msg, MSG_CONFIG)) return false;
memcpy(msg.target, mac, 6);
msg.body.config.zones = config.zones;
msg.body.config.follows = config.followsGroup ? 1 : 0;
msg.body.config.inactiveLevel = config.inactiveLevel;
msg.body.config.displayBrightness = config.displayBrightness;
strncpy(msg.body.config.name, config.name, GROUP_NAME_MAX);
msg.body.config.name[GROUP_NAME_MAX] = '\0';
send(msg, GROUP_COMMAND_REPEATS);
return true;
}
uint8_t peerList(PeerInfo* out, uint8_t max) {
uint8_t n = 0;
portENTER_CRITICAL(&peerLock);
const uint32_t now = millis();
for (int i = 0; i < GROUP_MAX_PEERS && n < max; i++) {
if (!peers[i].used) continue;
macToHex(peers[i].mac, out[n].id);
memcpy(out[n].name, peers[i].name, sizeof(out[n].name));
snprintf(out[n].ip, sizeof(out[n].ip), "%u.%u.%u.%u", peers[i].ip[0], peers[i].ip[1],
peers[i].ip[2], peers[i].ip[3]);
out[n].inactiveLevel = peers[i].inactiveLevel;
out[n].displayBrightness = peers[i].displayBrightness;
out[n].levelDb = peers[i].levelDb;
out[n].zones = peers[i].zones;
out[n].followsGroup = peers[i].followsGroup;
out[n].ageMs = now - peers[i].heardMs;
n++;
}
portEXIT_CRITICAL(&peerLock);
return n;
}
// Only units that actually follow the group count towards coverage: one
// deliberately running independently is not part of the arrangement, and
// counting it would report overlaps that do not matter.
uint8_t zoneCoverage() {
uint8_t covered = settings::get().groupLevel ? settings::get().zones : 0;
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (peers[i].used && peers[i].followsGroup) covered |= peers[i].zones;
}
portEXIT_CRITICAL(&peerLock);
return covered;
}
uint8_t zoneOverlap() {
uint8_t seen = settings::get().groupLevel ? settings::get().zones : 0;
uint8_t twice = 0;
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (!peers[i].used || !peers[i].followsGroup) continue;
twice |= seen & peers[i].zones;
seen |= peers[i].zones;
}
portEXIT_CRITICAL(&peerLock);
return twice;
}
uint32_t lastHeardMs() {
uint32_t newest = 0;
bool any = false;
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (peers[i].used && (!any || peers[i].heardMs > newest)) {
newest = peers[i].heardMs;
any = true;
}
}
portEXIT_CRITICAL(&peerLock);
return any ? millis() - newest : 0;
}
void publishLevel(float leqDb) {
Message msg;
if (!fill(msg, MSG_LEVEL)) return;
// Only the periodic level message is throttled; commands are rare and
// should go out the moment they happen.
if (millis() - lastSentMs < GROUP_BROADCAST_MS) return;
lastSentMs = millis();
msg.group = groupId;
msg.body.level.centiDb = int16_t(leqDb * 100.0f);
msg.body.level.zones = settings::get().zones;
msg.body.level.follows = settings::get().groupLevel ? 1 : 0;
strncpy(msg.body.level.name, localNameBuf, GROUP_NAME_MAX);
msg.body.level.name[GROUP_NAME_MAX] = '\0';
msg.body.level.inactiveLevel = settings::get().inactiveLevel;
msg.body.level.displayBrightness = settings::get().displayBrightness;
msg.body.level.ip[0] = localIp & 0xFF;
msg.body.level.ip[1] = (localIp >> 8) & 0xFF;
msg.body.level.ip[2] = (localIp >> 16) & 0xFF;
msg.body.level.ip[3] = (localIp >> 24) & 0xFF;
send(msg);
}
void publishSettings() {
Message msg;
if (!fill(msg, MSG_SETTINGS)) return;
const Settings& s = settings::get();
msg.body.settings.dbMin = s.dbMin;
msg.body.settings.dbMax = s.dbMax;
msg.body.settings.brightness = s.brightness;
msg.body.settings.combine = s.combine;
send(msg, GROUP_COMMAND_REPEATS);
}
void publishMode() {
Message msg;
if (!fill(msg, MSG_MODE)) return;
const uint32_t rgb = signal_light::manualColor();
msg.body.mode.mode = uint8_t(signal_light::mode());
msg.body.mode.r = (rgb >> 16) & 0xFF;
msg.body.mode.g = (rgb >> 8) & 0xFF;
msg.body.mode.b = rgb & 0xFF;
send(msg, GROUP_COMMAND_REPEATS);
}
void publishSelfTest() {
Message msg;
if (!fill(msg, MSG_SELF_TEST)) return;
send(msg, GROUP_COMMAND_REPEATS);
}
float groupLevel(float ownDb, uint8_t combine) {
if (!running) return ownDb;
float loudest = ownDb;
float total = ownDb;
uint8_t counted = 1;
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
if (!peers[i].used) continue;
if (peers[i].levelDb > loudest) loudest = peers[i].levelDb;
total += peers[i].levelDb;
counted++;
}
portEXIT_CRITICAL(&peerLock);
// With nobody else talking both rules collapse to our own reading, so a
// lone unit needs no special case at the call site.
return combine == COMBINE_AVERAGE ? total / float(counted) : loudest;
}
void tick() {
if (!running) return;
const uint32_t now = millis();
// Pace out any pending command repeats.
if (repeatsLeft > 0 && int32_t(now - nextRepeatMs) >= 0) {
transmit(repeatMsg);
repeatsLeft--;
nextRepeatMs = now + GROUP_COMMAND_GAP_MS;
}
portENTER_CRITICAL(&peerLock);
for (int i = 0; i < GROUP_MAX_PEERS; i++) {
// Unsigned arithmetic, so this survives the millis() rollover.
if (peers[i].used && now - peers[i].heardMs >= GROUP_PEER_TIMEOUT_MS) peers[i].used = false;
}
const Inbox pending = inbox;
inbox = Inbox{};
portEXIT_CRITICAL(&peerLock);
if (!pending.haveSettings && !pending.haveMode && !pending.haveSelfTest &&
!pending.haveConfig) {
return;
}
applyingRemote = true;
if (pending.haveSettings) {
settings::setDbMin(pending.dbMin);
settings::setDbMax(pending.dbMax);
settings::setBrightness(pending.brightness);
settings::setCombine(pending.combine);
signal_light::setBrightness(settings::get().brightness);
}
if (pending.haveMode) {
switch (signal_light::Mode(pending.mode)) {
case signal_light::Mode::Manual: signal_light::setManualColor(pending.color); break;
case signal_light::Mode::Off: signal_light::setMode(signal_light::Mode::Off); break;
default: signal_light::setMode(signal_light::Mode::Auto); break;
}
}
// Guarded because commands are repeated: restarting a running test would
// show as the sequence stuttering.
if (pending.haveConfig) {
settings::setZones(pending.zones);
settings::setGroupLevel(pending.follows);
settings::setInactiveLevel(pending.inactiveLevel);
settings::setDisplayBrightness(pending.displayBrightness);
if (pending.name[0] != '\0') settings::setUnitName(pending.name);
const Settings& s = settings::get();
signal_light::setZones(s.zones, s.inactiveLevel);
display::setBrightness(s.displayBrightness);
}
if (pending.haveSelfTest && !signal_light::selfTestRunning()) signal_light::startSelfTest();
applyingRemote = false;
}
} // namespace group_sync
#endif // FEATURE_ESPNOW