SensESP 3.5.1-alpha
Universal Signal K sensor toolkit ESP32
Loading...
Searching...
No Matches
signalk_ws_client.cpp
Go to the documentation of this file.
1#include "sensesp.h"
2
3#include "signalk_ws_client.h"
4
5#include <ArduinoJson.h>
6#include <ESPmDNS.h>
7#include <esp_http_client.h>
8
9#include <map>
10
12
13#ifdef SENSESP_SSL_SUPPORT
14#include <mbedtls/pem.h>
15#include <mbedtls/sha256.h>
16#include <mbedtls/ssl.h>
17#include <mbedtls/x509_crt.h>
18
20#endif
21
22#include <memory>
23#include <new>
24
25#include "Arduino.h"
26#include "elapsedMillis.h"
27#include "esp_arduino_version.h"
31#include "sensesp/system/uuid.h"
32#include "sensesp_app.h"
33
34namespace sensesp {
35
36constexpr int kWsClientTaskStackSize = 8192; // Stack for the connect worker
37constexpr int kWsTransportTaskStackSize = 6144; // Stack for esp_websocket_client internal task
38constexpr TickType_t kWsSendTimeoutTicks = pdMS_TO_TICKS(5000);
39// Periodic delta telemetry must never block the caller (it will be driven from
40// the event loop). 0 = enqueue if the ws-client lock and transport are free
41// right now, otherwise fail fast. Deltas are supersedable. See SignalK/SensESP#1033.
42constexpr TickType_t kWsDeltaSendTimeoutTicks = 0;
43// A device that keeps producing a delta larger than the buffer would drop one
44// every send cycle; rate-limit the warning so it does not flood the log.
45constexpr uint32_t kOversizeDropLogIntervalMs = 10000;
46
48
49static const char* kRequestPermission = "readwrite";
50
51#ifdef SENSESP_SSL_SUPPORT
52// Convert a SHA256 hash to hex string
53static void sha256_to_hex(const uint8_t* sha256, char* hex) {
54 for (int i = 0; i < 32; i++) {
55 sprintf(hex + (i * 2), "%02x", sha256[i]);
56 }
57 hex[64] = '\0';
58}
59
60// SHA256 of a certificate's raw DER, as a 64-char hex String.
61static String cert_fingerprint(const mbedtls_x509_crt* crt) {
62 uint8_t sha256[32];
63 mbedtls_sha256_context ctx;
64 mbedtls_sha256_init(&ctx);
65 mbedtls_sha256_starts(&ctx, 0); // 0 = SHA256 (not SHA224)
66 mbedtls_sha256_update(&ctx, crt->raw.p, crt->raw.len);
67 mbedtls_sha256_finish(&ctx, sha256);
68 mbedtls_sha256_free(&ctx);
69 char hex[65];
70 sha256_to_hex(sha256, hex);
71 return String(hex);
72}
73
74// Maximum stored/displayed length of a certificate CN.
75static constexpr size_t kMaxPinCnLen = 64;
76
77// Extract the CN from a certificate subject, for display. The CN is
78// attacker-controlled, so the result is length-bounded and stripped of
79// non-printable and quote/backslash characters before it is stored or shown.
80static String cert_common_name(const mbedtls_x509_crt* crt) {
81 char dn[256];
82 int len = mbedtls_x509_dn_gets(dn, sizeof(dn), &crt->subject);
83 if (len <= 0) {
84 return String("");
85 }
86 const char* cn = strstr(dn, "CN=");
87 if (cn == nullptr) {
88 return String("");
89 }
90 cn += 3; // skip "CN="
91 String out;
92 for (size_t i = 0; i < kMaxPinCnLen && cn[i] != '\0' && cn[i] != ','; i++) {
93 char c = cn[i];
94 if (c >= 0x20 && c < 0x7f && c != '"' && c != '\\') {
95 out += c;
96 }
97 }
98 return out;
99}
100
101// PEM-encode a certificate's DER for storage. The encode buffer is allocated on
102// demand and freed on return rather than held for the device's lifetime: TOFU CA
103// capture happens only during the TLS handshake, so a permanent .bss buffer
104// would waste ~4 KB for the whole uptime. It is heap- rather than stack-
105// allocated because the verify callback runs on a small TLS task stack.
106static String cert_to_pem(const mbedtls_x509_crt* crt) {
107 // Sized for a large CA cert (RSA-4096 + SANs/extensions). The size is fixed,
108 // not derived from crt->raw.len: PEM is base64 (~4/3 of the DER) plus the
109 // header/footer and line breaks, so the encoded form is always larger than the
110 // DER. On overflow mbedtls_pem_write_buffer returns an error and this returns
111 // "" -- callers must treat an empty PEM as "no usable CA" and fail safe, never
112 // store it.
113 constexpr size_t kPemBufSize = 4096;
114 std::unique_ptr<unsigned char[]> pem_buf(
115 new (std::nothrow) unsigned char[kPemBufSize]);
116 if (!pem_buf) {
117 ESP_LOGE("SKWSClient", "TOFU: PEM buffer allocation failed");
118 return String("");
119 }
120 size_t olen = 0;
121 int r = mbedtls_pem_write_buffer(
122 "-----BEGIN CERTIFICATE-----\n", "-----END CERTIFICATE-----\n",
123 crt->raw.p, crt->raw.len, pem_buf.get(), kPemBufSize, &olen);
124 if (r != 0) {
125 ESP_LOGE("SKWSClient", "TOFU: PEM encode failed (-0x%x)", -r);
126 return String("");
127 }
128 return String(reinterpret_cast<const char*>(pem_buf.get()));
129}
130
131// Normalized (lowercase, sorted, deduplicated, comma-joined) set of the
132// certificate's dNSName SANs, for TOFU identity binding. Empty if the cert has
133// no DNS SAN. Two certs with the same set of names produce the same string
134// regardless of order, so an exact-match comparison is stable across leaf
135// rotation but changes when the identity itself changes.
136static String cert_dns_sans(const mbedtls_x509_crt* crt) {
137 std::set<String> names;
138 for (const mbedtls_x509_sequence* cur = &crt->subject_alt_names;
139 cur != nullptr && cur->buf.p != nullptr; cur = cur->next) {
140 mbedtls_x509_subject_alternative_name san;
141 memset(&san, 0, sizeof(san));
142 if (mbedtls_x509_parse_subject_alt_name(&cur->buf, &san) != 0) {
143 continue;
144 }
145 if (san.type == MBEDTLS_X509_SAN_DNS_NAME &&
146 san.san.unstructured_name.p != nullptr &&
147 san.san.unstructured_name.len > 0) {
148 size_t n = san.san.unstructured_name.len;
149 if (n > 255) {
150 n = 255;
151 }
152 String name;
153 name.reserve(n);
154 for (size_t i = 0; i < n; i++) {
155 char c = static_cast<char>(san.san.unstructured_name.p[i]);
156 if (c >= 'A' && c <= 'Z') {
157 c = static_cast<char>(c + ('a' - 'A'));
158 }
159 name += c;
160 }
161 names.insert(name);
162 }
163 mbedtls_x509_free_subject_alt_name(&san);
164 }
165 String out;
166 for (const String& s : names) {
167 if (!out.isEmpty()) {
168 out += ",";
169 }
170 out += s;
171 }
172 return out;
173}
174
175// TOFU verification callback - called during the TLS handshake, once per
176// presented certificate, highest depth (CA) first down to depth 0 (leaf).
177// Returns 0 to allow, non-zero to reject.
178static int tofu_verify_callback(void* ctx, mbedtls_x509_crt* crt, int depth,
179 uint32_t* flags) {
180 SKWSClient* client = static_cast<SKWSClient*>(ctx);
181 if (client == nullptr) {
182 ESP_LOGW("SKWSClient", "TOFU: no client context, allowing connection");
183 *flags = 0;
184 return 0;
185 }
186
187 if (!client->is_tofu_enabled()) {
188 // Verification disabled: accept any certificate (insecure opt-out).
189 *flags = 0;
190 return 0;
191 }
192
193 // CA-anchor mode: the stored CA was installed as the trust anchor and esp-tls
194 // runs VERIFY_REQUIRED, so mbedTLS has already validated this certificate and
195 // set *flags. Honor that result rather than clearing it.
196 if (client->has_tofu_ca()) {
197 if (*flags != 0) {
198 ESP_LOGE("SKWSClient", "TOFU: certificate failed CA validation (0x%lx)",
199 (unsigned long)*flags);
200 client->flag_cert_error();
201 return MBEDTLS_ERR_X509_CERT_VERIFY_FAILED;
202 }
203 // Identity binding: at the leaf, require the same SAN identity captured when
204 // the CA was pinned. This is what keeps a public CA safe — a valid leaf for
205 // a different name signed by the same CA (e.g. any Let's Encrypt cert) is
206 // rejected here even though it chains to the pinned CA.
207 if (depth == 0 && !client->get_tofu_san().isEmpty()) {
208 if (cert_dns_sans(crt) != client->get_tofu_san()) {
209 ESP_LOGE("SKWSClient", "TOFU: leaf identity (SAN) mismatch, rejecting");
210 client->flag_cert_error();
211 return MBEDTLS_ERR_X509_CERT_VERIFY_FAILED;
212 }
213 }
214 return 0;
215 }
216
217 // Capture / leaf-fingerprint mode (VERIFY_OPTIONAL). Collect a CA candidate
218 // from the higher-depth certificates (which arrive first), then decide at the
219 // leaf. The first CA:TRUE certificate seen is the highest in the presented
220 // chain (closest to the root), so it is the preferred anchor.
221 // Certificates above the leaf (depth > 0): collect a CA candidate -- the
222 // highest CA:TRUE cert, which arrives first -- and never fail on chain-trust
223 // flags during capture. A depth-0 certificate is ALWAYS treated as the leaf
224 // by the fingerprint/role decision below, even if it is self-signed with
225 // CA:TRUE: a single presented certificate is pinned as a leaf, never adopted
226 // as a CA trust anchor. (Adopting a CA:TRUE leaf here would skip the
227 // fingerprint check and let a mismatched self-signed cert be accepted.)
228 if (depth > 0) {
229 // basicConstraints CA:TRUE has no public getter in mbedTLS 3.x, so the
230 // flag is read through the MBEDTLS_PRIVATE accessor macro.
231 bool is_ca = crt->MBEDTLS_PRIVATE(ca_istrue) != 0;
232 if (is_ca && !client->has_pending_ca()) {
233 String ca_pem = cert_to_pem(crt);
234 // Skip on encode failure: leaving no pending CA fails safe to leaf-
235 // fingerprint mode rather than committing an empty (pin-disabling) anchor.
236 if (!ca_pem.isEmpty()) {
237 client->stash_pending_ca(ca_pem, cert_common_name(crt));
238 }
239 }
240 *flags = 0;
241 return 0;
242 }
243
244 // depth == 0: the leaf. Always run the fingerprint/role decision.
245 String leaf_fp = cert_fingerprint(crt);
246 String leaf_san = cert_dns_sans(crt);
247 bool leaf_matches =
248 client->has_tofu_fingerprint() && client->get_tofu_fingerprint() == leaf_fp;
250 client->has_tofu_fingerprint(), leaf_matches, client->has_pending_ca(),
251 !leaf_san.isEmpty());
252
253 switch (decision) {
255 ESP_LOGE("SKWSClient", "TOFU: leaf fingerprint mismatch, rejecting");
256 client->flag_cert_error();
257 return MBEDTLS_ERR_X509_CERT_VERIFY_FAILED;
259 ESP_LOGI("SKWSClient", "TOFU: first use, pinning leaf %s", leaf_fp.c_str());
260 client->stash_pending_leaf(leaf_fp, cert_common_name(crt));
261 break;
263 ESP_LOGI("SKWSClient", "TOFU: first use, pinning issuing CA (identity %s)",
264 leaf_san.c_str());
265 client->set_pending_san(leaf_san); // bind the leaf identity to the CA
266 break;
268 // Leaf matches the stored fingerprint: keep the leaf pin and drop any
269 // stray pending state (e.g. a CA stashed at depth > 0 this handshake --
270 // mode is fixed at first use, so a later-presented CA is not adopted).
271 client->clear_pending_tofu();
272 break;
273 }
274 *flags = 0;
275 return 0;
276}
277
278// Attach function installed for the TLS connection. In CA-anchor mode it
279// installs the pinned CA as the mbedTLS trust anchor (esp-tls runs
280// VERIFY_REQUIRED, so mbedTLS validates the chain against it); otherwise it
281// selects VERIFY_OPTIONAL so the verify callback can capture / fingerprint.
282static esp_err_t tofu_crt_bundle_attach(void* conf) {
283 mbedtls_ssl_config* ssl_conf = static_cast<mbedtls_ssl_config*>(conf);
284 SKWSClient* client = ws_client;
285
286 if (client != nullptr && client->is_tofu_enabled() && client->has_tofu_ca()) {
287 // Re-parse the stored CA into a static cert that outlives the handshake.
288 static mbedtls_x509_crt pinned_ca;
289 static bool pinned_ca_inited = false;
290 if (pinned_ca_inited) {
291 mbedtls_x509_crt_free(&pinned_ca);
292 }
293 mbedtls_x509_crt_init(&pinned_ca);
294 pinned_ca_inited = true;
295 const String& pem = client->get_tofu_ca();
296 int r = mbedtls_x509_crt_parse(
297 &pinned_ca, reinterpret_cast<const unsigned char*>(pem.c_str()),
298 pem.length() + 1);
299 if (r == 0) {
300 mbedtls_ssl_conf_ca_chain(ssl_conf, &pinned_ca, nullptr);
301 // esp-tls already set VERIFY_REQUIRED before calling us; keep it.
302 ESP_LOGD("SKWSClient", "TOFU: pinned CA installed as trust anchor");
303 } else {
304 // Stored CA won't parse (corruption/bug). Fail closed: leave
305 // VERIFY_REQUIRED with no trust anchor so the handshake is rejected and
306 // surfaces as a certificate error requiring a manual reset -- never
307 // silently downgrade to accept-any.
308 ESP_LOGE("SKWSClient",
309 "TOFU: stored CA failed to parse (-0x%x); connections will fail "
310 "until reset",
311 -r);
312 }
313 } else {
314 // Capture / leaf-fingerprint mode (or TOFU disabled): the callback decides.
315 mbedtls_ssl_conf_authmode(ssl_conf, MBEDTLS_SSL_VERIFY_OPTIONAL);
316 }
317
318 mbedtls_ssl_conf_verify(ssl_conf, tofu_verify_callback, client);
319 return ESP_OK;
320}
321#endif // SENSESP_SSL_SUPPORT
322
323static void websocket_event_handler(void* handler_args,
324 esp_event_base_t base,
325 int32_t event_id, void* event_data) {
326 // Drop events from a client that has been handed off for destruction: the
327 // handler was registered with the generation current at build time; if that no
328 // longer matches the live generation, this callback is from a reaped client.
329 if (ws_client == nullptr ||
330 static_cast<uint32_t>(reinterpret_cast<uintptr_t>(handler_args)) !=
332 return;
333 }
334 esp_websocket_event_data_t* data = (esp_websocket_event_data_t*)event_data;
335 switch (event_id) {
336 case WEBSOCKET_EVENT_CONNECTED:
338 break;
339 case WEBSOCKET_EVENT_DISCONNECTED:
341 break;
342 case WEBSOCKET_EVENT_DATA:
343 // Only process text frames (opcode 0x1) and continuation frames (0x0).
344 // Control frames (ping/pong/close: 0x8-0xA) have no JSON payload.
345 if (data->op_code <= 0x2) {
346 ws_client->on_receive_delta((uint8_t*)data->data_ptr, data->data_len);
347 }
348 break;
349 case WEBSOCKET_EVENT_ERROR:
350 // The HTTP status of the failed upgrade (e.g. 401 for a rejected token)
351 // lets us distinguish a bad token from a transient transport error. The
352 // handshake status field only exists in the newer esp_websocket_client
353 // component pulled for SSL builds; the version bundled with the EOL
354 // espressif32 Arduino platform lacks it, so fall back to 0 (not
355 // applicable), matching how on_error() treats a transient failure.
356#ifdef SENSESP_SSL_SUPPORT
357 ws_client->on_error(data->error_handle.esp_ws_handshake_status_code);
358#else
360#endif
361 break;
362 }
363}
364
365SKWSClient::SKWSClient(const String& config_path,
366 std::shared_ptr<SKDeltaQueue> sk_delta_queue,
367 const String& server_address, uint16_t server_port,
368 bool use_mdns)
369 : FileSystemSaveable{config_path},
370 conf_server_address_{server_address},
371 conf_server_port_{server_port},
372 use_mdns_{use_mdns},
373 sk_delta_queue_{sk_delta_queue} {
374 // a SKWSClient object observes its own connection_state_ member
375 // and simply passes through any notification it emits. As a result,
376 // whenever the value of connection_state_ is updated, observers of the
377 // SKWSClient object get automatically notified.
379 [this]() { this->emit(this->connection_state_.get()); });
380
381 // process any received updates in the main task
382 event_loop()->onRepeat(1, [this]() { this->process_received_updates(); });
383
384 // set the singleton object pointer
385 ws_client = this;
386
387 load();
388
389 // Connect the counters
391
392 event_loop()->onDelay(0, [this]() {
393 ESP_LOGD(__FILENAME__, "Starting SKWSClient");
394 // Run the connection lifecycle on the event loop instead of a dedicated
395 // task: connect() is a non-blocking dispatcher (it spawns a worker for the
396 // blocking auth legs) and send_delta() is non-blocking, so neither stalls
397 // the loop. The ~100 ms cadence matches the former task's vTaskDelay(100ms).
398 event_loop()->onRepeat(100, [this]() {
399 // Retry a teardown whose reaper failed to spawn under OOM (no-op if none).
401 if (is_connect_due()) {
402 connect();
403 }
404 send_delta();
405 });
406 MDNS.addService("signalk-sensesp", "tcp", 80);
407 });
408}
409
419
420bool should_clear_token_on_status(int handshake_status) {
421 return handshake_status == kHttpUnauthorized;
422}
423
430void SKWSClient::on_error(int handshake_status) {
432 if (auth_token_ != NULL_AUTH_TOKEN &&
433 should_clear_token_on_status(handshake_status)) {
434 // The server rejected the token on the WebSocket upgrade (e.g. the server
435 // was reinstalled, or the device was moved to a different server). Clear it
436 // so the next reconnect requests fresh access, and reset the backoff so
437 // attempts after the already-scheduled one come at the short interval. A
438 // 401 on a tokenless upgrade cannot mean a stale token: it falls through
439 // to the retry branch, keeping the backoff growing and avoiding a config
440 // flash write per attempt. A non-401 error (transport, TLS, network)
441 // leaves the token intact and simply retries.
442 ESP_LOGW(__FILENAME__, "Token rejected (%d), requesting new access",
443 handshake_status);
444 auth_token_ = NULL_AUTH_TOKEN;
445 save();
447 } else {
448 ESP_LOGW(__FILENAME__, "Websocket client error.");
449 }
450}
451
458 // The connection is fully established and the server proved possession of its
459 // certificate's private key, so it is safe to persist any anchor stashed
460 // during this attempt's handshake (first-use capture).
461 this->commit_pending_tofu();
464 this->sk_delta_queue_->reset_meta_send();
465 ESP_LOGI(__FILENAME__, "Subscribing to Signal K listeners...");
466 this->subscribe_listeners();
467}
468
476 bool output_available = false;
477 JsonDocument subscription;
478 subscription["context"] = "vessels.self";
479
481 const std::vector<SKListener*>& listeners = SKListener::get_listeners();
482
483 if (listeners.size() > 0) {
484 output_available = true;
485 JsonArray subscribe = subscription["subscribe"].to<JsonArray>();
486
487 // Collapse listeners that share a path into one subscription entry.
488 // With sendMeta=all (a connection-level flag), a single subscription
489 // delivers both value and meta deltas, so a value listener and a
490 // metadata listener on the same path — the common gauge pattern —
491 // would otherwise emit two identical entries. Keep the smallest
492 // period so the fastest listener's cadence wins.
493 std::map<String, int> path_period;
494 for (size_t i = 0; i < listeners.size(); i++) {
495 auto* listener = listeners.at(i);
496 const String& sk_path = listener->get_sk_path();
497 int listen_delay = listener->get_listen_delay();
498 auto it = path_period.find(sk_path);
499 if (it == path_period.end() || listen_delay < it->second) {
500 path_period[sk_path] = listen_delay;
501 }
502 }
503 for (const auto& [sk_path, listen_delay] : path_period) {
504 JsonObject subscribe_path = subscribe.add<JsonObject>();
505 subscribe_path["path"] = sk_path;
506 subscribe_path["period"] = listen_delay;
507 ESP_LOGI(__FILENAME__, "Adding %s subscription with listen_delay %d\n",
508 sk_path.c_str(), listen_delay);
509 }
510 }
512
513 if (output_available &&
515 String json_message;
516
517 serializeJson(subscription, json_message);
518 ESP_LOGI(__FILENAME__, "Subscription JSON message:\n %s",
519 json_message.c_str());
520 int result = esp_websocket_client_send_text(
521 client_.load(), json_message.c_str(), json_message.length(),
523 if (result < 0) {
524 ESP_LOGE(__FILENAME__, "Subscription send failed (result=%d)", result);
525 }
526 }
527}
528
536void SKWSClient::on_receive_delta(uint8_t* payload, size_t length) {
537 // Need to work on null-terminated strings
538 constexpr size_t kMaxWsMessageSize = 4096;
539 if (length > kMaxWsMessageSize) {
540 ESP_LOGW(__FILENAME__, "WebSocket message too large (%u bytes), dropping",
541 (unsigned)length);
542 return;
543 }
544 std::unique_ptr<char[]> buf(new char[length + 1]);
545 memcpy(buf.get(), payload, length);
546 buf[length] = 0;
547
548#ifdef SIGNALK_PRINT_RCV_DELTA
549 ESP_LOGD(__FILENAME__, "Websocket payload received: %s", buf.get());
550#endif
551
552 JsonDocument message;
553 auto error = deserializeJson(message, buf.get());
554
555 if (!error) {
556 if (message["updates"].is<JsonVariant>()) {
557 on_receive_updates(message);
558 }
559
560 if (message["put"].is<JsonVariant>()) {
561 on_receive_put(message);
562 }
563
564 // Putrequest contains also requestId Key GA
565 if (message["requestId"].is<JsonVariant>() &&
566 !message["put"].is<JsonVariant>()) {
568 }
569 } else {
570 ESP_LOGE(__FILENAME__, "deserializeJson error: %s", error.c_str());
571 }
572}
573
581void SKWSClient::on_receive_updates(JsonDocument& message) {
582 // Process updates from subscriptions...
583 JsonArray updates = message["updates"];
584
585 // With sendMeta=all enabled by default, the server pushes meta deltas to
586 // every client. Skip copying them onto the queue unless something actually
587 // consumes them. Compute this before taking received_updates_semaphore_ to
588 // preserve the global lock order (SKListener before received_updates).
589 bool has_meta_listener = false;
591 for (SKListener* listener : SKListener::get_listeners()) {
592 if (listener->wants_meta()) {
593 has_meta_listener = true;
594 break;
595 }
596 }
598
600 for (size_t i = 0; i < updates.size(); i++) {
601 JsonObject update = updates[i];
602
603 JsonArray values = update["values"];
604
605 for (size_t vi = 0; vi < values.size(); vi++) {
606 // Copy each value into an owned document for processing in the main
607 // task (decoupled from `message`, which is freed when on_receive_delta
608 // returns).
610 ru.is_meta = false;
611 ru.doc.set(values[vi]);
612 enqueue_received_update(std::move(ru));
613 }
614
615 // Meta deltas only arrive when subscribed with sendMeta=all and are
616 // typically one-shot per path (at subscribe + on metadata change). Copy
617 // each meta entry into an owned document the same way; SKMetadataListener
618 // consumers receive it (path-routed) on the main task. No user code runs
619 // here, so the critical section stays free of arbitrary callbacks. Skip
620 // entirely when no listener consumes meta (see has_meta_listener above).
621 if (has_meta_listener) {
622 JsonArray meta_entries = update["meta"];
623 for (size_t mi = 0; mi < meta_entries.size(); mi++) {
624 JsonObject entry = meta_entries[mi];
625 if (entry["path"].isNull() || entry["value"].isNull()) continue;
627 ru.is_meta = true;
628 ru.doc.set(entry); // {path, value: {...meta...}}
629 enqueue_received_update(std::move(ru));
630 }
631 }
632 }
634}
635
646 // Per-kind caps, tunable via build_flags (see signalk_ws_client.h). Meta and
647 // value deltas are budgeted independently so a metadata burst cannot evict
648 // pending values, and vice versa.
649 constexpr size_t kMaxReceivedValues = SENSESP_MAX_RECEIVED_VALUE_UPDATES;
650 constexpr size_t kMaxReceivedMeta = SENSESP_MAX_RECEIVED_META_UPDATES;
651
652 const bool is_meta = update.is_meta;
653 const size_t cap = is_meta ? kMaxReceivedMeta : kMaxReceivedValues;
654
655 size_t kind_count = 0;
656 for (const auto& ru : received_updates_) {
657 if (ru.is_meta == is_meta) kind_count++;
658 }
659
660 while (kind_count >= cap) {
661 // Drop the oldest entry of the SAME kind, leaving the other kind intact.
662 for (auto it = received_updates_.begin(); it != received_updates_.end();
663 ++it) {
664 if (it->is_meta == is_meta) {
665 received_updates_.erase(it);
666 kind_count--;
667 break;
668 }
669 }
670 ESP_LOGW(__FILENAME__, "Dropping oldest received %s update (queue full)",
671 is_meta ? "meta" : "value");
672 }
673
674 received_updates_.push_back(std::move(update));
675}
676
685
686 const std::vector<SKListener*>& listeners = SKListener::get_listeners();
687 const std::vector<SKPutListener*>& put_listeners =
689
691 // Count only value/put deltas toward the rx metric; meta deltas are
692 // low-frequency one-shots and would otherwise inflate it.
693 int num_updates = 0;
694 while (!received_updates_.empty()) {
695 ReceivedUpdate& ru = received_updates_.front();
696
697 if (ru.is_meta) {
698 // Capture the path into an owned String before moving the document: the
699 // const char* would dangle once ru.doc is moved into the shared_ptr.
700 String path = ru.doc["path"].as<String>();
701 // Move the already-owned queue document into a refcounted, read-only
702 // shared_ptr (no extra deep copy) so it safely outlives the queue entry
703 // and any deferred consumer, then fan it out to matching listeners.
704 std::shared_ptr<const JsonDocument> meta_doc =
705 std::make_shared<const JsonDocument>(std::move(ru.doc));
706 for (size_t i = 0; i < listeners.size(); i++) {
707 SKListener* listener = listeners[i];
708 if (listener->wants_meta() && listener->matches(path)) {
709 listener->parse_meta(meta_doc);
710 }
711 }
712 } else {
713 num_updates++;
714 const char* path = ru.doc["path"];
715 JsonObject value = ru.doc.as<JsonObject>();
716
717 for (size_t i = 0; i < listeners.size(); i++) {
718 SKListener* listener = listeners[i];
719 if (!listener->wants_meta() && listener->matches(path)) {
720 listener->parse_value(value);
721 }
722 }
723 // to be able to parse values of Put Listeners GA
724 for (size_t i = 0; i < put_listeners.size(); i++) {
725 SKPutListener* listener = put_listeners[i];
726 if (listener->get_sk_path().equals(path)) {
727 listener->parse_value(value);
728 }
729 }
730 }
731 received_updates_.pop_front();
732 }
734 delta_rx_count_producer_.set(num_updates);
735
737}
738
746void SKWSClient::on_receive_put(JsonDocument& message) {
747 // Process PUT requests...
748 JsonArray puts = message["put"];
749 bool all_matched = true;
750 for (size_t i = 0; i < puts.size(); i++) {
751 JsonObject value = puts[i];
752 const char* path = value["path"];
753 bool matched = false;
754
756 const std::vector<SKPutListener*>& listeners =
758 for (size_t j = 0; j < listeners.size(); j++) {
759 SKPutListener* listener = listeners[j];
760 if (listener->get_sk_path().equals(path)) {
762 ru.is_meta = false;
763 ru.doc.set(value);
765 enqueue_received_update(std::move(ru));
767 matched = true;
768 }
769 }
771
772 if (!matched) {
773 all_matched = false;
774 }
775 }
776
777 // Send back a single request response if still connected
779 JsonDocument put_response;
780 put_response["requestId"] = message["requestId"];
781 if (all_matched) {
782 put_response["state"] = "COMPLETED";
783 put_response["statusCode"] = 200;
784 } else {
785 put_response["state"] = "FAILED";
786 put_response["statusCode"] = 405;
787 }
788 String response_text;
789 serializeJson(put_response, response_text);
790 int result = esp_websocket_client_send_text(
791 client_.load(), response_text.c_str(), response_text.length(),
793 if (result < 0) {
794 ESP_LOGE(__FILENAME__, "PUT response send failed (result=%d)", result);
795 }
796 }
797}
798
806void SKWSClient::sendTXT(String& payload) {
808 int result = esp_websocket_client_send_text(
809 client_.load(), payload.c_str(), payload.length(), kWsSendTimeoutTicks);
810 if (result < 0) {
811 ESP_LOGE(__FILENAME__, "sendTXT failed (result=%d)", result);
812 }
813 }
814}
815
816bool SKWSClient::get_mdns_service(String& server_address,
817 uint16_t& server_port) {
818 // get IP address using an mDNS query
819 // Try SSL service first, then fall back to non-SSL
820 int num = MDNS.queryService("signalk-wss", "tcp");
821 if (num > 0) {
822 // Found SSL-enabled server
823 ssl_enabled_ = true;
824 ESP_LOGI(__FILENAME__, "Found Signal K server via mDNS (signalk-wss)");
825 } else {
826 // Try non-SSL service
827 num = MDNS.queryService("signalk-ws", "tcp");
828 if (num == 0) {
829 // no service found
830 return false;
831 }
832 // Found non-SSL server, disable SSL
833 ssl_enabled_ = false;
834 ESP_LOGI(__FILENAME__, "Found Signal K server via mDNS (signalk-ws)");
835 }
836
837#if ESP_ARDUINO_VERSION_MAJOR < 3
838 server_address = MDNS.IP(0).toString();
839#else
840 server_address = MDNS.address(0).toString();
841#endif
842 server_port = MDNS.port(0);
843 ESP_LOGI(__FILENAME__, "Found server %s (port %d)", server_address.c_str(),
844 server_port);
845 return true;
846}
847
848// Event handler for detect_ssl() to capture the Location response header
849static esp_err_t detect_ssl_event_handler(esp_http_client_event_t* evt) {
850 if (evt->event_id == HTTP_EVENT_ON_HEADER) {
851 // Check for Location header (case-insensitive)
852 if (strcasecmp(evt->header_key, "Location") == 0) {
853 String* location = static_cast<String*>(evt->user_data);
854 *location = evt->header_value;
855 }
856 }
857 return ESP_OK;
858}
859
861 // Try to detect if the server requires SSL by checking for HTTP->HTTPS
862 // redirects
863 String url =
864 String("http://") + server_address_ + ":" + server_port_ + "/signalk";
865
866 ESP_LOGD(__FILENAME__, "Probing for SSL redirect at %s", url.c_str());
867
868 String location;
869
870 esp_http_client_config_t config = {};
871 config.url = url.c_str();
872 config.disable_auto_redirect = true;
873 config.timeout_ms = 10000;
874 config.event_handler = detect_ssl_event_handler;
875 config.user_data = &location;
876
877 esp_http_client_handle_t client = esp_http_client_init(&config);
878 if (client == nullptr) {
879 ESP_LOGE(__FILENAME__, "Failed to initialize HTTP client");
880 return false;
881 }
882
883 esp_err_t err = esp_http_client_perform(client);
884 int status_code = esp_http_client_get_status_code(client);
885 esp_http_client_cleanup(client);
886
887 if (err != ESP_OK) {
888 ESP_LOGD(__FILENAME__, "HTTP request failed: %s", esp_err_to_name(err));
889 return false;
890 }
891
892 if ((status_code == 301 || status_code == 302 ||
893 status_code == 307 || status_code == 308) &&
894 location.startsWith("https://")) {
895 ESP_LOGI(__FILENAME__, "SSL redirect detected, enabling HTTPS/WSS");
896 ssl_enabled_ = true;
897 save();
898 return true;
899 }
900
901 return false;
902}
903
904
906 // A prior certificate rejection leaves the state in kSKWSCertificateError;
907 // treat it like a disconnect for retry purposes so the device keeps trying
908 // (and re-surfaces the cert error each failed attempt).
911 return;
912 }
913
914 // A connect attempt is already running on a worker task; let it finish.
915 if (auth_job_running_.load()) {
916 return;
917 }
918
919 // Reap any client left from a previous attempt before starting a new one,
920 // off this context. While the reap is in flight, defer bring-up so at most one
921 // client ever exists; state stays Disconnected, so a later cycle retries.
922 if (client_.load() != nullptr) {
924 }
925 if (teardown_in_progress_.load()) {
926 return;
927 }
928
929 // Discard any anchor candidate stashed by a previous attempt; it is only
930 // committed after a fully successful connection (see on_connected). Clearing
931 // here also prevents committing stale data if a resumed TLS session skips the
932 // certificate callback.
934
935 // Schedule next attempt with backoff in case this one fails.
936 // Will be reset on successful connection.
938
939 // Wait for the active network provisioner (WiFi, Ethernet, …) to be
940 // up before initiating the WS connection. The provisioner abstracts
941 // away whether we're on WiFi, Ethernet, or some other transport.
942 auto provisioner = SensESPApp::get()->get_network_provisioner();
943 if (!provisioner || !provisioner->is_connected()) {
944 ESP_LOGI(__FILENAME__,
945 "Network is not yet up. SignalK client connection will be "
946 "initiated when the link comes up.");
947 return;
948 }
949
951
952 // The rest of the attempt — mDNS resolve, SSL detect, and the access-request /
953 // poll / connect_ws leg — makes blocking HTTP/mDNS calls, so run it on a
954 // one-shot worker task. The SK/event-loop context stays responsive, and
955 // auth_job_running_ keeps at most one attempt in flight.
956 auth_job_running_.store(true);
957 if (xTaskCreate(&SKWSClient::connect_worker, "SKWSConnect",
958 kWsClientTaskStackSize, this, 1, nullptr) != pdPASS) {
959 ESP_LOGE(__FILENAME__, "connect worker spawn failed");
960 auth_job_running_.store(false);
962 }
963}
964
966 auto* self = static_cast<SKWSClient*>(arg);
968 self->auth_job_running_.store(false);
969 vTaskDelete(nullptr);
970}
971
973 ESP_LOGI(__FILENAME__, "Initiating websocket connection with server...");
974
975 if (use_mdns_) {
976 if (!get_mdns_service(this->server_address_, this->server_port_)) {
977 ESP_LOGE(__FILENAME__,
978 "No Signal K server found in network when using mDNS service!");
979 } else {
980 ESP_LOGI(__FILENAME__,
981 "Signal K server has been found at address %s:%d by mDNS.",
982 this->server_address_.c_str(), this->server_port_);
983 }
984 } else {
986 this->server_port_ = this->conf_server_port_;
987 }
988
989 if (!this->server_address_.isEmpty() && this->server_port_ > 0) {
990 ESP_LOGD(__FILENAME__,
991 "Websocket is connecting to Signal K server on address %s:%d",
992 this->server_address_.c_str(), this->server_port_);
993
994 // Detect if server requires SSL (check for HTTP->HTTPS redirects)
995 if (!ssl_enabled_) {
996 detect_ssl();
997 }
998 } else {
999 // host and port not defined - don't try to connect
1000 ESP_LOGD(__FILENAME__,
1001 "Websocket is not connecting to Signal K server because host and "
1002 "port are not defined.");
1004 return;
1005 }
1006
1007 if (this->polling_href_.length() > 0 && this->polling_href_.startsWith("/")) {
1008 // existing pending request
1010 this->polling_href_);
1011 return;
1012 }
1013
1014 if (this->auth_token_ == NULL_AUTH_TOKEN) {
1015 // initiate HTTP authentication
1016 ESP_LOGD(__FILENAME__, "No prior authorization token present.");
1018 return;
1019 }
1020
1021 // A token is already present. Validate it before streaming.
1022#ifdef SENSESP_SSL_SUPPORT
1023 // Connect the WebSocket directly rather than first probing the token over a
1024 // separate HTTPS request: on memory-constrained targets (e.g. ESP32-C3) the
1025 // back-to-back token-probe TLS handshake and the WebSocket TLS handshake
1026 // fragment the heap, and the second fails to allocate
1027 // (MBEDTLS_ERR_SSL_ALLOC_FAILED). The server validates the token on the
1028 // upgrade itself; a 401 there is handled in on_error() (clears the token and
1029 // re-requests access on the next reconnect).
1030 this->connect_ws(this->server_address_, this->server_port_);
1031#else
1032 // The bundled (non-SSL) esp_websocket_client reports no upgrade status, so
1033 // on_error() cannot tell a rejected token from a transient failure. Probe the
1034 // token over plain HTTP first -- there is no TLS handshake to fragment the
1035 // heap. A 401 there clears the token and re-requests access.
1036 this->test_token(this->server_address_, this->server_port_);
1037#endif
1038}
1039
1041 // No-op: esp_websocket_client handles data via event callbacks
1042}
1043
1044#ifndef SENSESP_SSL_SUPPORT
1045void SKWSClient::test_token(const String server_address,
1046 const uint16_t server_port) {
1047 String url = String("http://") + server_address + ":" + server_port +
1048 "/signalk/v1/stream";
1049 ESP_LOGD(__FILENAME__, "Testing token with url %s", url.c_str());
1050
1051 const String full_token = String("Bearer ") + auth_token_;
1052 ESP_LOGD(__FILENAME__, "Authorization: %.8s...[redacted]", full_token.c_str());
1053
1054 esp_http_client_config_t config = {};
1055 config.url = url.c_str();
1056 config.timeout_ms = 10000;
1057
1058 esp_http_client_handle_t client = esp_http_client_init(&config);
1059 if (client == nullptr) {
1060 ESP_LOGE(__FILENAME__, "Failed to initialize HTTP client");
1062 return;
1063 }
1064
1065 esp_http_client_set_header(client, "Authorization", full_token.c_str());
1066
1067 // Use streaming API for GET request
1068 esp_err_t err = esp_http_client_open(client, 0);
1069 if (err != ESP_OK) {
1070 ESP_LOGE(__FILENAME__, "Failed to open HTTP connection: %s",
1071 esp_err_to_name(err));
1072 esp_http_client_cleanup(client);
1074 return;
1075 }
1076
1077 int content_length = esp_http_client_fetch_headers(client);
1078 int http_code = esp_http_client_get_status_code(client);
1079
1080 ESP_LOGD(__FILENAME__, "Testing resulted in http status %d", http_code);
1081
1082 // Read response body
1083 String payload;
1084 if (content_length > 0 && content_length < 4096) {
1085 char* buffer = new char[content_length + 1];
1086 int read_len = esp_http_client_read(client, buffer, content_length);
1087 buffer[read_len > 0 ? read_len : 0] = '\0';
1088 payload = String(buffer);
1089 delete[] buffer;
1090 } else {
1091 // Chunked encoding or unknown/large content length - read in chunks
1092 char buffer[512];
1093 int read_len;
1094 while ((read_len = esp_http_client_read(client, buffer,
1095 sizeof(buffer) - 1)) > 0) {
1096 buffer[read_len] = '\0';
1097 payload += String(buffer);
1098 if (payload.length() > 4096) break;
1099 }
1100 }
1101
1102 esp_http_client_close(client);
1103 esp_http_client_cleanup(client);
1104
1105 if (payload.length() > 0) {
1106 ESP_LOGD(__FILENAME__, "Returned payload (%d bytes): %s", payload.length(),
1107 payload.c_str());
1108 }
1109
1110 if (http_code == 426) {
1111 // HTTP status 426 is "Upgrade Required", the expected response for a
1112 // websocket endpoint reached over plain HTTP: the token is valid.
1113 ESP_LOGD(__FILENAME__, "Attempting to connect to Signal K Websocket...");
1114 this->connect_ws(server_address, server_port);
1115 } else if (http_code == kHttpUnauthorized) {
1116 // Token is invalid/expired - clear it and request new access.
1117 // Keep client_id_ so we reuse the same device identity.
1118 ESP_LOGW(__FILENAME__, "Token rejected (401), requesting new access");
1119 this->auth_token_ = NULL_AUTH_TOKEN;
1120 this->save();
1121 this->send_access_request(server_address, server_port);
1122 } else if (http_code > 0) {
1124 } else {
1125 ESP_LOGE(__FILENAME__, "HTTP request failed with code %d", http_code);
1127 }
1128}
1129#endif // !SENSESP_SSL_SUPPORT
1130
1131void SKWSClient::send_access_request(const String server_address,
1132 const uint16_t server_port) {
1133 ESP_LOGD(__FILENAME__, "Sending access request (client_id=%s, ssl=%d)",
1134 client_id_.c_str(), ssl_enabled_);
1135 if (client_id_ == "") {
1136 // generate a client ID
1138 save();
1139 }
1140
1141 // create a new access request
1142 JsonDocument doc;
1143 doc["clientId"] = client_id_;
1144 doc["description"] =
1145 String("SensESP device: ") + SensESPBaseApp::get_hostname();
1146 doc["permissions"] = kRequestPermission;
1147 String json_req = "";
1148 serializeJson(doc, json_req);
1149
1150 ESP_LOGD(__FILENAME__, "Access request: %s", json_req.c_str());
1151
1152 String protocol = ssl_enabled_ ? "https://" : "http://";
1153 String url = protocol + server_address + ":" + server_port +
1154 "/signalk/v1/access/requests";
1155 ESP_LOGD(__FILENAME__, "Access request url: %s", url.c_str());
1156
1157 esp_http_client_config_t config = {};
1158 config.url = url.c_str();
1159 config.method = HTTP_METHOD_POST;
1160 config.timeout_ms = 10000;
1161#ifdef SENSESP_SSL_SUPPORT
1162 if (ssl_enabled_) {
1163 config.crt_bundle_attach = tofu_crt_bundle_attach;
1164 config.skip_cert_common_name_check = true;
1165 }
1166#endif
1167
1168 esp_http_client_handle_t client = esp_http_client_init(&config);
1169 if (client == nullptr) {
1170 ESP_LOGE(__FILENAME__, "Failed to initialize HTTP client");
1172 // Don't clear client_id_ - keep device identity for retry
1173 return;
1174 }
1175
1176 esp_http_client_set_header(client, "Content-Type", "application/json");
1177
1178 // Use streaming API: open -> write request -> fetch headers -> read response
1179 esp_err_t err = esp_http_client_open(client, json_req.length());
1180 if (err != ESP_OK) {
1181 ESP_LOGE(__FILENAME__, "Failed to open HTTP connection: %s", esp_err_to_name(err));
1182 esp_http_client_cleanup(client);
1184 return;
1185 }
1186
1187 int write_len = esp_http_client_write(client, json_req.c_str(), json_req.length());
1188 if (write_len < 0 || write_len != (int)json_req.length()) {
1189 ESP_LOGE(__FILENAME__, "Failed to write request body (wrote %d of %d bytes)",
1190 write_len, json_req.length());
1191 esp_http_client_close(client);
1192 esp_http_client_cleanup(client);
1194 return;
1195 }
1196
1197 int content_length = esp_http_client_fetch_headers(client);
1198 int http_code = esp_http_client_get_status_code(client);
1199
1200 ESP_LOGD(__FILENAME__, "HTTP response: code=%d, content_length=%d", http_code, content_length);
1201
1202 // Read response body
1203 String payload;
1204 char buffer[512];
1205 int read_len;
1206 while ((read_len = esp_http_client_read(client, buffer, sizeof(buffer) - 1)) > 0) {
1207 buffer[read_len] = '\0';
1208 payload += String(buffer);
1209 if (payload.length() > 4096) break;
1210 }
1211 ESP_LOGD(__FILENAME__, "Response payload (%d bytes): %s",
1212 payload.length(), payload.c_str());
1213
1214 esp_http_client_close(client);
1215 esp_http_client_cleanup(client);
1216
1217 // Parse JSON response for both 202 and 400 status codes
1218 deserializeJson(doc, payload.c_str());
1219 String state = doc["state"].is<const char*>() ? doc["state"].as<String>() : "";
1220 String href = doc["href"].is<const char*>() ? doc["href"].as<String>() : "";
1221 String message = doc["message"].is<const char*>() ? doc["message"].as<String>() : "";
1222
1223 ESP_LOGD(__FILENAME__, "Access request response: http=%d, state=%s, href=%s",
1224 http_code, state.c_str(), href.c_str());
1225 if (message.length() > 0) {
1226 ESP_LOGI(__FILENAME__, "Server message: %s", message.c_str());
1227 }
1228
1229 // HTTP 400 with href means "already requested" - save href for polling on
1230 // next connect() cycle (after backoff)
1231 if (http_code == 400 && href.length() > 0 && href.startsWith("/")) {
1232 ESP_LOGI(__FILENAME__, "Existing request found, will poll href: %s", href.c_str());
1233 polling_href_ = href;
1234 save();
1236 return;
1237 }
1238
1239 // HTTP 202 with href means new request pending - save href for polling on
1240 // next connect() cycle (after backoff)
1241 if (http_code == 202 && href.length() > 0 && href.startsWith("/")) {
1242 polling_href_ = href;
1243 save();
1245 return;
1246 }
1247
1248 // HTTP 404 means the server has no security enabled — access requests are
1249 // not available. Connect without a token.
1250 if (http_code == 404) {
1251 ESP_LOGI(__FILENAME__,
1252 "Server security disabled (404 on access request) — connecting "
1253 "without token");
1254 auth_token_ = NULL_AUTH_TOKEN;
1255 this->connect_ws(server_address, server_port);
1256 return;
1257 }
1258
1259 // Can't proceed - disconnect and retry later
1260 ESP_LOGW(__FILENAME__, "Cannot handle response: http=%d, state=%s", http_code, state.c_str());
1262}
1263
1264void SKWSClient::poll_access_request(const String server_address,
1265 const uint16_t server_port,
1266 const String href) {
1267 ESP_LOGD(__FILENAME__, "Polling SK Server for authentication token");
1268
1269 String protocol = ssl_enabled_ ? "https://" : "http://";
1270 String url = protocol + server_address + ":" + server_port + href;
1271
1272 esp_http_client_config_t config = {};
1273 config.url = url.c_str();
1274 config.timeout_ms = 10000;
1275#ifdef SENSESP_SSL_SUPPORT
1276 if (ssl_enabled_) {
1277 config.crt_bundle_attach = tofu_crt_bundle_attach;
1278 config.skip_cert_common_name_check = true;
1279 }
1280#endif
1281
1282 esp_http_client_handle_t client = esp_http_client_init(&config);
1283 if (client == nullptr) {
1284 ESP_LOGE(__FILENAME__, "Failed to initialize HTTP client");
1286 return;
1287 }
1288
1289 // Use streaming API for GET request
1290 esp_err_t err = esp_http_client_open(client, 0);
1291 if (err != ESP_OK) {
1292 ESP_LOGE(__FILENAME__, "Failed to open HTTP connection: %s", esp_err_to_name(err));
1293 esp_http_client_cleanup(client);
1295 return;
1296 }
1297
1298 int content_length = esp_http_client_fetch_headers(client);
1299 int http_code = esp_http_client_get_status_code(client);
1300
1301 // Read response body
1302 String payload;
1303 if (content_length > 0 && content_length < 4096) {
1304 char* buffer = new char[content_length + 1];
1305 int read_len = esp_http_client_read(client, buffer, content_length);
1306 buffer[read_len > 0 ? read_len : 0] = '\0';
1307 payload = String(buffer);
1308 delete[] buffer;
1309 } else {
1310 // Chunked encoding or unknown/large content length - read in chunks
1311 char buffer[512];
1312 int read_len;
1313 while ((read_len = esp_http_client_read(client, buffer, sizeof(buffer) - 1)) > 0) {
1314 buffer[read_len] = '\0';
1315 payload += String(buffer);
1316 if (payload.length() > 4096) break;
1317 }
1318 }
1319
1320 // An APPROVED poll response carries the access token in its body, so log only
1321 // the status and size, never the payload itself. The web log buffer exposes
1322 // captured log lines over HTTP, so a payload dump here would leak the token.
1323 ESP_LOGD(__FILENAME__, "Poll response: http=%d, %d bytes", http_code,
1324 static_cast<int>(payload.length()));
1325
1326 esp_http_client_close(client);
1327 esp_http_client_cleanup(client);
1328
1329 if (http_code == 200 || http_code == 202) {
1330 JsonDocument doc;
1331 auto error = deserializeJson(doc, payload.c_str());
1332 if (error) {
1333 ESP_LOGW(__FILENAME__, "WARNING: Could not deserialize http payload.");
1334 ESP_LOGW(__FILENAME__, "DeserializationError: %s", error.c_str());
1336 return;
1337 }
1338 String state = doc["state"];
1339 ESP_LOGD(__FILENAME__, "%s", state.c_str());
1340 if (state == "PENDING") {
1342 return;
1343 }
1344 if (state == "COMPLETED") {
1345 JsonObject access_req = doc["accessRequest"];
1346 String permission = access_req["permission"];
1347
1348 polling_href_ = "";
1349 save();
1350
1351 if (permission == "DENIED") {
1352 ESP_LOGW(__FILENAME__, "Permission denied");
1354 return;
1355 }
1356
1357 if (permission == "APPROVED") {
1358 ESP_LOGI(__FILENAME__, "Permission granted");
1359 String token = access_req["token"];
1360 auth_token_ = token;
1361 save();
1362 this->connect_ws(server_address, server_port);
1363 return;
1364 }
1365 }
1366 } else {
1367 if (http_code == 404 || http_code == 500) {
1368 // Server doesn't recognize this request (stale href after
1369 // server restart, different server, or security disabled).
1370 // Clear the polling href so the next connect cycle starts
1371 // a fresh access-request flow.
1372 ESP_LOGD(__FILENAME__,
1373 "Got %d polling access request — clearing stale href.",
1374 http_code);
1375 polling_href_ = "";
1376 save();
1378 return;
1379 }
1380 // any other HTTP status code
1381 ESP_LOGW(__FILENAME__,
1382 "Can't handle response %d to pending access request.\n",
1383 http_code);
1385 return;
1386 }
1387 // Catch-all: a 200/202 COMPLETED whose permission is neither APPROVED nor
1388 // DENIED, or any unexpected state, leaves no terminal state. With no live
1389 // client, no event will move us off Authorizing, so fall back to Disconnected
1390 // to retry rather than wedge.
1392}
1393
1394void SKWSClient::connect_ws(const String& host, const uint16_t port) {
1395 // connect() reaps any prior client before dispatching here, so client_ is
1396 // null and no teardown is in flight. Guard defensively against a leak.
1397 if (client_.load() != nullptr || teardown_in_progress_.load()) {
1398 ESP_LOGW(__FILENAME__, "connect_ws: prior client not reaped; deferring");
1401 return;
1402 }
1403
1404 // Discard any anchor candidate stashed by the earlier esp_http_client legs
1405 // (token check / access request). Only the WebSocket handshake -- the one
1406 // whose success reaches on_connected and proves the server holds the leaf's
1407 // private key -- may populate the anchor that gets committed.
1410
1411 String protocol = ssl_enabled_ ? "wss" : "ws";
1412 String path = "/signalk/v1/stream?subscribe=none";
1413 if (send_meta_enabled_) path += "&sendMeta=all";
1414 String url = protocol + "://" + host + ":" + String(port) + path;
1415
1416 ESP_LOGD(__FILENAME__, "Connecting WebSocket to %s", url.c_str());
1417
1418 // Build authorization header string (must persist through init call)
1419 String auth_header;
1420 if (auth_token_ != NULL_AUTH_TOKEN) {
1421 auth_header = String("Authorization: Bearer ") + auth_token_ + "\r\n";
1422 }
1423
1424 // Configure WebSocket client
1425 esp_websocket_client_config_t config = {};
1426 config.uri = url.c_str();
1427 config.task_stack = kWsTransportTaskStackSize;
1428 config.buffer_size = SENSESP_SK_WS_BUFFER_SIZE;
1429 if (auth_header.length() > 0) {
1430 config.headers = auth_header.c_str();
1431 }
1432
1433#ifdef SENSESP_SSL_SUPPORT
1434 if (ssl_enabled_) {
1435 // Use custom crt_bundle_attach to disable SSL verification
1436 // This directly configures mbedTLS to skip certificate verification
1437 config.crt_bundle_attach = tofu_crt_bundle_attach;
1438 config.skip_cert_common_name_check = true;
1439 }
1440#endif
1441
1442 esp_websocket_client_handle_t h = esp_websocket_client_init(&config);
1443 if (h == nullptr) {
1444 ESP_LOGE(__FILENAME__, "Failed to initialize WebSocket client");
1446 return;
1447 }
1448 client_.store(h);
1449
1450 // Register the event handler tagged with the current generation, so that any
1451 // late event from this client after it is later handed off for destruction is
1452 // dropped by websocket_event_handler (generation mismatch).
1453 esp_websocket_register_events(
1454 h, WEBSOCKET_EVENT_ANY, websocket_event_handler,
1455 reinterpret_cast<void*>(
1456 static_cast<uintptr_t>(client_generation_.load())));
1457
1458 // Start the client
1459 esp_err_t err = esp_websocket_client_start(h);
1460 if (err != ESP_OK) {
1461 ESP_LOGE(__FILENAME__, "Failed to start WebSocket client: %s",
1462 esp_err_to_name(err));
1463 // Null the shared handle before freeing it so a concurrent send sees null.
1464 client_.store(nullptr);
1465 esp_websocket_client_destroy(h);
1467 return;
1468 }
1469
1470 ESP_LOGD(__FILENAME__, "WebSocket client started, waiting for connection...");
1471}
1472
1476
1477namespace {
1478struct WsTeardownArg {
1479 esp_websocket_client_handle_t handle;
1480 SKWSClient* self;
1481};
1482} // namespace
1483
1485 auto* a = static_cast<WsTeardownArg*>(arg);
1486 // Blocking: esp_websocket_client_stop() waits for the client's task to exit
1487 // (up to one ~1 s poll cycle), destroy() frees its buffers. Run here, on a
1488 // throwaway task, so it never blocks the connect/event-loop context.
1489 esp_websocket_client_stop(a->handle);
1490 esp_websocket_client_destroy(a->handle);
1491 a->self->teardown_in_progress_.store(false);
1492 delete a;
1493 vTaskDelete(nullptr);
1494}
1495
1496void SKWSClient::reap_async(esp_websocket_client_handle_t old) {
1497 if (old == nullptr) {
1498 return;
1499 }
1500 auto* arg = new WsTeardownArg{old, this};
1501 if (xTaskCreate(&SKWSClient::teardown_task, "SKWSTeardown", 4096, arg, 1,
1502 nullptr) == pdPASS) {
1503 // The task now owns `old` and clears teardown_in_progress_ when done.
1504 pending_teardown_.store(nullptr);
1505 return;
1506 }
1507 // Spawn failed (OOM). Do NOT reap synchronously: a ~1 s stop()+destroy() on
1508 // the event loop would stall every consumer under the exact heap pressure
1509 // this path exists for. Stash the handle and retry on the next loop tick;
1510 // teardown_in_progress_ stays set so bring-up remains deferred (at most one
1511 // un-reaped client, no leak beyond it).
1512 delete arg;
1513 pending_teardown_.store(old);
1514 ESP_LOGW(__FILENAME__, "teardown task spawn failed; will retry next cycle");
1515}
1516
1518 // Single atomic check-and-null: if two contexts race here, only one gets the
1519 // handle, so it is stopped/destroyed exactly once (no double-free).
1520 esp_websocket_client_handle_t old = client_.exchange(nullptr);
1521 if (old == nullptr) {
1522 return;
1523 }
1524 // Invalidate the old client's generation so any late event it dispatches
1525 // while being reaped is dropped by websocket_event_handler.
1526 client_generation_.fetch_add(1);
1527 teardown_in_progress_.store(true);
1528 reap_async(old);
1529}
1530
1532 // Set state first so event handler callbacks and send callsites see the
1533 // disconnected state and skip operations on the client being destroyed.
1536}
1537
1540 if (sk_delta_queue_->data_available()) {
1541 std::vector<String> deltas;
1542 sk_delta_queue_->get_deltas(deltas);
1543 bool first = true;
1544 for (const auto& delta : deltas) {
1545 if (sk_delta_exceeds_ws_buffer(delta.length(),
1547 // Drop the delta to keep the connection alive (signalk_ws_delta_size.h
1548 // explains why an oversize delta would otherwise abort it). Unlike the
1549 // transient send failure below, an oversize delta is deterministic, so
1550 // do NOT re-arm metadata here: get_deltas() already bundles metadata
1551 // into the first delta and marks it sent, and re-arming would have
1552 // get_deltas() rebuild the same oversize first delta every cycle --
1553 // dropped and re-armed forever, never delivered. Leaving it sent lets
1554 // the next, metadata-free first delta fit and flow; metadata waits for
1555 // a reconnect or a larger SENSESP_SK_WS_BUFFER_SIZE.
1556 uint32_t now = millis();
1557 if (last_oversize_log_ms_ == 0 ||
1559 ESP_LOGW(__FILENAME__,
1560 "Delta too large (%u B > %u buffer); dropped to keep the "
1561 "connection alive -- raise SENSESP_SK_WS_BUFFER_SIZE",
1562 (unsigned)delta.length(),
1563 (unsigned)SENSESP_SK_WS_BUFFER_SIZE);
1565 }
1566 first = false;
1567 continue;
1568 }
1569 int send_result = esp_websocket_client_send_text(
1570 client_.load(), delta.c_str(), delta.length(), kWsDeltaSendTimeoutTicks);
1571 if (send_result < 0) {
1572 // Non-blocking send (0 timeout) did not complete: either brief
1573 // ws-client lock contention (retry next cycle) or the transport
1574 // backpressured and esp_websocket_client aborted the connection
1575 // internally -- its disconnect/error event drives reconnect. Never
1576 // block or tear the connection down from here. Deltas are
1577 // supersedable, so drop the rest of this batch. See SignalK/SensESP#1033.
1578 if (first) {
1579 // get_deltas() builds one-shot metadata (units, zones, ...) into the
1580 // first delta and marks it sent before it leaves the device. The
1581 // first delta is the only one that can carry that metadata, so if its
1582 // send is the one that fails, re-arm metadata for the next batch --
1583 // otherwise the server runs without it until the next reconnect.
1584 sk_delta_queue_->reset_meta_send();
1585 }
1586 ESP_LOGW(__FILENAME__,
1587 "Delta send incomplete (result=%d); dropping rest of batch",
1588 send_result);
1589 break;
1590 }
1592 first = false;
1593 }
1594 }
1595 }
1596}
1597
1598bool SKWSClient::to_json(JsonObject& root) {
1599 root["sk_address"] = this->conf_server_address_;
1600 root["sk_port"] = this->conf_server_port_;
1601 root["use_mdns"] = this->use_mdns_;
1602
1603 root["token"] = this->auth_token_;
1604 root["client_id"] = this->client_id_;
1605 root["polling_href"] = this->polling_href_;
1606
1607 root["ssl_enabled"] = this->ssl_enabled_;
1608 root["tofu_enabled"] = this->tofu_enabled_;
1609 // Persisted trust anchor: leaf fingerprint (legacy / leaf mode) and/or CA PEM.
1610 root["tofu_fingerprint"] = this->tofu_fingerprint_;
1611 root["tofu_ca_pem"] = this->tofu_ca_pem_;
1612 root["tofu_san"] = this->tofu_san_;
1613 // Read-only display fields for the pinned identity.
1614 root["tofu_pin_cn"] = this->tofu_pin_cn_;
1615 root["tofu_pin_is_ca"] = this->tofu_pin_is_ca_;
1616 root["send_meta_enabled"] = this->send_meta_enabled_;
1617 return true;
1618}
1619
1620bool SKWSClient::from_json(const JsonObject& config) {
1621 if (config["sk_address"].is<String>()) {
1622 this->conf_server_address_ = config["sk_address"].as<String>();
1623 }
1624 if (config["sk_port"].is<int>()) {
1625 this->conf_server_port_ = config["sk_port"].as<int>();
1626 }
1627 if (config["use_mdns"].is<bool>()) {
1628 this->use_mdns_ = config["use_mdns"].as<bool>();
1629 }
1630 if (config["token"].is<String>()) {
1631 this->auth_token_ = config["token"].as<String>();
1632 }
1633 if (config["client_id"].is<String>()) {
1634 this->client_id_ = config["client_id"].as<String>();
1635 }
1636 if (config["polling_href"].is<String>()) {
1637 String href = config["polling_href"].as<String>();
1638 // Only accept valid hrefs (must start with /)
1639 this->polling_href_ = href.startsWith("/") ? href : "";
1640 }
1641
1642 if (config["ssl_enabled"].is<bool>()) {
1643 this->ssl_enabled_ = config["ssl_enabled"].as<bool>();
1644 }
1645 if (config["tofu_enabled"].is<bool>()) {
1646 this->tofu_enabled_ = config["tofu_enabled"].as<bool>();
1647 }
1648 // A legacy config carries only tofu_fingerprint (loads as leaf mode — the
1649 // migration entry point); newer configs may also carry a pinned CA and the
1650 // display fields. Tolerate both.
1651 if (config["tofu_fingerprint"].is<String>()) {
1652 this->tofu_fingerprint_ = config["tofu_fingerprint"].as<String>();
1653 }
1654 if (config["tofu_ca_pem"].is<String>()) {
1655 this->tofu_ca_pem_ = config["tofu_ca_pem"].as<String>();
1656 }
1657 if (config["tofu_san"].is<String>()) {
1658 this->tofu_san_ = config["tofu_san"].as<String>();
1659 }
1660 if (config["tofu_pin_cn"].is<String>()) {
1661 this->tofu_pin_cn_ = config["tofu_pin_cn"].as<String>();
1662 }
1663 if (config["tofu_pin_is_ca"].is<bool>()) {
1664 this->tofu_pin_is_ca_ = config["tofu_pin_is_ca"].as<bool>();
1665 }
1666 if (config["send_meta_enabled"].is<bool>()) {
1667 this->send_meta_enabled_ = config["send_meta_enabled"].as<bool>();
1668 }
1669
1670 return true;
1671}
1672
1679 auto state = get_connection_state();
1680 switch (state) {
1682 return "Authorizing with SignalK";
1684 return "Connected";
1686 return "Connecting";
1688 return "Disconnected";
1690 return "Certificate verification failed";
1691 }
1692
1693 return "Unknown";
1694}
1695
1696} // namespace sensesp
virtual bool load() override
Load and populate the object from a persistent storage.
Definition saveable.cpp:8
virtual bool save() override
Save the object to a persistent storage.
Definition saveable.cpp:40
virtual void set(const C &input) override final
Definition integrator.h:34
int attach(std::function< void()> observer)
Attach an observer callback.
Definition observable.h:40
An Obervable class that listens for Signal K stream deltas and notifies any observers of value change...
static bool take_semaphore(uint64_t timeout_ms=0)
static void release_semaphore()
virtual void parse_value(const JsonObject &json)
virtual bool matches(const String &path) const
virtual void parse_meta(const std::shared_ptr< const JsonDocument > &meta_doc)
static const std::vector< SKListener * > & get_listeners()
virtual bool wants_meta() const
An Obervable class that listens for Signal K PUT requests coming over the websocket connection and no...
static const std::vector< SKPutListener * > & get_listeners()
virtual void parse_value(const JsonObject &put)=0
static void handle_response(JsonDocument &response)
The websocket connection to the Signal K server.
void commit_pending_tofu()
Persist a stashed anchor after a successful (authenticated) connection. Called from on_connected().
void poll_access_request(const String host, const uint16_t port, const String href)
SKWSClient(const String &config_path, std::shared_ptr< SKDeltaQueue > sk_delta_queue, const String &server_address, uint16_t server_port, bool use_mdns=true)
void on_receive_delta(uint8_t *payload, size_t length)
Called when the websocket receives a delta.
void process_received_updates()
Loop through the received updates and process them.
std::atomic< esp_websocket_client_handle_t > pending_teardown_
std::atomic< bool > auth_job_running_
void connect_ws(const String &host, const uint16_t port)
Integrator< int, int > delta_tx_count_producer_
TaskQueueProducer< SKWSConnectionState > connection_state_
SKWSConnectionState get_connection_state()
std::atomic< bool > teardown_in_progress_
void on_receive_put(JsonDocument &message)
Called when a PUT event is received.
uint32_t client_generation() const
Generation tag of the currently-valid client. The event handler compares the generation it was regist...
void release_received_updates_semaphore()
void run_connect_attempt()
The (blocking) connect attempt — mDNS resolve, SSL detect, and the access-request / poll / connect_ws...
void send_access_request(const String host, const uint16_t port)
std::atomic< uint32_t > client_generation_
void on_error(int handshake_status)
Called when the websocket connection encounters an error.
Integrator< int, int > delta_rx_count_producer_
std::list< ReceivedUpdate > received_updates_
virtual bool to_json(JsonObject &root) override final
void detach_teardown()
Hand client_ to a detached one-shot task that stops+destroys it, so the blocking teardown never runs ...
void enqueue_received_update(ReceivedUpdate &&update)
Push a received delta entry onto the queue, enforcing a per-kind budget so a metadata burst (one meta...
void sendTXT(String &payload)
Send some processed data to the websocket.
uint32_t last_oversize_log_ms_
millis() timestamp of the last oversize-delta-drop warning, used to rate-limit it when a device keeps...
bool get_mdns_service(String &server_address, uint16_t &server_port)
std::atomic< esp_websocket_client_handle_t > client_
void on_disconnected()
Called when the websocket connection is disconnected.
void reap_async(esp_websocket_client_handle_t old)
Spawn the detached reaper for old; on spawn failure (OOM) stash it in pending_teardown_ for a later r...
void subscribe_listeners()
Subscribes the SK delta paths to the websocket.
void on_receive_updates(JsonDocument &message)
Called when a delta update is received.
void test_token(const String host, const uint16_t port)
void set_connection_state(SKWSConnectionState state)
bool take_received_updates_semaphore(unsigned long int timeout_ms=0)
virtual bool from_json(const JsonObject &config) override final
String get_connection_status()
Get a String representation of the current connection state.
static void connect_worker(void *arg)
TaskQueueProducer< int > delta_tx_tick_producer_
Emits the number of deltas sent since last report.
static void teardown_task(void *arg)
Body of the detached teardown task (stop+destroy+self-delete).
std::shared_ptr< SKDeltaQueue > sk_delta_queue_
void restart()
Drop the current connection (detached, non-blocking teardown) and let the reconnect path rebuild it.
void on_connected()
Called when the websocket connection is established.
static std::shared_ptr< SensESPApp > get()
Get the singleton instance of the SensESPApp.
Definition sensesp_app.h:57
static String get_hostname()
Get the current hostname.
virtual void set(const T &value) override
virtual const T & get() const
void emit(const SKWSConnectionState &new_value)
std::enable_if< std::is_base_of< ValueConsumer< typenameVConsumer::input_type >, VConsumer >::value &&std::is_convertible< T, typenameVConsumer::input_type >::value, std::shared_ptr< VConsumer > >::type connect_to(std::shared_ptr< VConsumer > consumer)
Connect a producer to a transform with a different input type.
std::shared_ptr< reactesp::EventLoop > event_loop()
Definition sensesp.cpp:9
String generate_uuid4()
Generate a random UUIDv4 string.
Definition uuid.cpp:5
constexpr int kWsClientTaskStackSize
bool sk_delta_exceeds_ws_buffer(size_t delta_length, size_t buffer_size)
True if a delta is too large to hand to esp_websocket_client as a single tx chunk,...
SKWSClient * ws_client
constexpr TickType_t kWsSendTimeoutTicks
constexpr uint32_t kOversizeDropLogIntervalMs
constexpr int kWsTransportTaskStackSize
TofuCaptureDecision tofu_decide_capture(bool has_leaf_anchor, bool leaf_matches_anchor, bool ca_present, bool leaf_has_identity)
TofuCaptureDecision
Capture-mode decision for TOFU certificate pinning.
@ kCaptureLeaf
first use, no usable CA in the chain: pin the leaf
@ kAccept
leaf matches the stored fingerprint; nothing to capture
@ kCaptureCa
first use, usable CA present: pin the CA
@ kReject
leaf does not match the stored fingerprint
bool should_clear_token_on_status(int handshake_status)
Decide whether a failed WebSocket upgrade should clear the auth token.
constexpr TickType_t kWsDeltaSendTimeoutTicks
esp_websocket_client_handle_t handle
SKWSClient * self
#define SENSESP_SK_WS_BUFFER_SIZE
#define SENSESP_MAX_RECEIVED_META_UPDATES
#define SENSESP_MAX_RECEIVED_VALUE_UPDATES
A single received delta entry awaiting dispatch on the main task.