5#include <ArduinoJson.h>
7#include <esp_http_client.h>
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>
26#include "elapsedMillis.h"
27#include "esp_arduino_version.h"
49static const char* kRequestPermission =
"readwrite";
51#ifdef SENSESP_SSL_SUPPORT
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]);
61static String cert_fingerprint(
const mbedtls_x509_crt* crt) {
63 mbedtls_sha256_context ctx;
64 mbedtls_sha256_init(&ctx);
65 mbedtls_sha256_starts(&ctx, 0);
66 mbedtls_sha256_update(&ctx, crt->raw.p, crt->raw.len);
67 mbedtls_sha256_finish(&ctx, sha256);
68 mbedtls_sha256_free(&ctx);
70 sha256_to_hex(sha256, hex);
75static constexpr size_t kMaxPinCnLen = 64;
80static String cert_common_name(
const mbedtls_x509_crt* crt) {
82 int len = mbedtls_x509_dn_gets(dn,
sizeof(dn), &crt->subject);
86 const char* cn = strstr(dn,
"CN=");
92 for (
size_t i = 0; i < kMaxPinCnLen && cn[i] !=
'\0' && cn[i] !=
','; i++) {
94 if (c >= 0x20 && c < 0x7f && c !=
'"' && c !=
'\\') {
106static String cert_to_pem(
const mbedtls_x509_crt* crt) {
113 constexpr size_t kPemBufSize = 4096;
114 std::unique_ptr<unsigned char[]> pem_buf(
115 new (std::nothrow)
unsigned char[kPemBufSize]);
117 ESP_LOGE(
"SKWSClient",
"TOFU: PEM buffer allocation failed");
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);
125 ESP_LOGE(
"SKWSClient",
"TOFU: PEM encode failed (-0x%x)", -r);
128 return String(
reinterpret_cast<const char*
>(pem_buf.get()));
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) {
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;
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'));
163 mbedtls_x509_free_subject_alt_name(&san);
166 for (
const String& s : names) {
167 if (!out.isEmpty()) {
178static int tofu_verify_callback(
void* ctx, mbedtls_x509_crt* crt,
int depth,
180 SKWSClient* client =
static_cast<SKWSClient*
>(ctx);
181 if (client ==
nullptr) {
182 ESP_LOGW(
"SKWSClient",
"TOFU: no client context, allowing connection");
187 if (!client->is_tofu_enabled()) {
196 if (client->has_tofu_ca()) {
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;
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;
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);
236 if (!ca_pem.isEmpty()) {
237 client->stash_pending_ca(ca_pem, cert_common_name(crt));
245 String leaf_fp = cert_fingerprint(crt);
246 String leaf_san = cert_dns_sans(crt);
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());
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));
263 ESP_LOGI(
"SKWSClient",
"TOFU: first use, pinning issuing CA (identity %s)",
265 client->set_pending_san(leaf_san);
271 client->clear_pending_tofu();
282static esp_err_t tofu_crt_bundle_attach(
void* conf) {
283 mbedtls_ssl_config* ssl_conf =
static_cast<mbedtls_ssl_config*
>(conf);
286 if (client !=
nullptr && client->is_tofu_enabled() && client->has_tofu_ca()) {
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);
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()),
300 mbedtls_ssl_conf_ca_chain(ssl_conf, &pinned_ca,
nullptr);
302 ESP_LOGD(
"SKWSClient",
"TOFU: pinned CA installed as trust anchor");
308 ESP_LOGE(
"SKWSClient",
309 "TOFU: stored CA failed to parse (-0x%x); connections will fail "
315 mbedtls_ssl_conf_authmode(ssl_conf, MBEDTLS_SSL_VERIFY_OPTIONAL);
318 mbedtls_ssl_conf_verify(ssl_conf, tofu_verify_callback, client);
323static void websocket_event_handler(
void* handler_args,
324 esp_event_base_t base,
325 int32_t event_id,
void* event_data) {
330 static_cast<uint32_t
>(
reinterpret_cast<uintptr_t
>(handler_args)) !=
334 esp_websocket_event_data_t* data = (esp_websocket_event_data_t*)event_data;
336 case WEBSOCKET_EVENT_CONNECTED:
339 case WEBSOCKET_EVENT_DISCONNECTED:
342 case WEBSOCKET_EVENT_DATA:
345 if (data->op_code <= 0x2) {
349 case WEBSOCKET_EVENT_ERROR:
356#ifdef SENSESP_SSL_SUPPORT
366 std::shared_ptr<SKDeltaQueue> sk_delta_queue,
367 const String& server_address, uint16_t server_port,
370 conf_server_address_{server_address},
371 conf_server_port_{server_port},
373 sk_delta_queue_{sk_delta_queue} {
393 ESP_LOGD(__FILENAME__,
"Starting SKWSClient");
406 MDNS.addService(
"signalk-sensesp",
"tcp", 80);
421 return handshake_status == kHttpUnauthorized;
442 ESP_LOGW(__FILENAME__,
"Token rejected (%d), requesting new access",
448 ESP_LOGW(__FILENAME__,
"Websocket client error.");
465 ESP_LOGI(__FILENAME__,
"Subscribing to Signal K listeners...");
476 bool output_available =
false;
477 JsonDocument subscription;
478 subscription[
"context"] =
"vessels.self";
483 if (listeners.size() > 0) {
484 output_available =
true;
485 JsonArray subscribe = subscription[
"subscribe"].to<JsonArray>();
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;
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);
513 if (output_available &&
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(),
524 ESP_LOGE(__FILENAME__,
"Subscription send failed (result=%d)", result);
538 constexpr size_t kMaxWsMessageSize = 4096;
539 if (length > kMaxWsMessageSize) {
540 ESP_LOGW(__FILENAME__,
"WebSocket message too large (%u bytes), dropping",
544 std::unique_ptr<char[]> buf(
new char[length + 1]);
545 memcpy(buf.get(), payload, length);
548#ifdef SIGNALK_PRINT_RCV_DELTA
549 ESP_LOGD(__FILENAME__,
"Websocket payload received: %s", buf.get());
552 JsonDocument message;
553 auto error = deserializeJson(message, buf.get());
556 if (message[
"updates"].is<JsonVariant>()) {
560 if (message[
"put"].is<JsonVariant>()) {
565 if (message[
"requestId"].is<JsonVariant>() &&
566 !message[
"put"].is<JsonVariant>()) {
570 ESP_LOGE(__FILENAME__,
"deserializeJson error: %s", error.c_str());
583 JsonArray updates = message[
"updates"];
589 bool has_meta_listener =
false;
592 if (listener->wants_meta()) {
593 has_meta_listener =
true;
600 for (
size_t i = 0; i < updates.size(); i++) {
601 JsonObject update = updates[i];
603 JsonArray values = update[
"values"];
605 for (
size_t vi = 0; vi < values.size(); vi++) {
611 ru.
doc.set(values[vi]);
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;
652 const bool is_meta = update.is_meta;
653 const size_t cap = is_meta ? kMaxReceivedMeta : kMaxReceivedValues;
655 size_t kind_count = 0;
657 if (ru.is_meta == is_meta) kind_count++;
660 while (kind_count >= cap) {
664 if (it->is_meta == is_meta) {
670 ESP_LOGW(__FILENAME__,
"Dropping oldest received %s update (queue full)",
671 is_meta ?
"meta" :
"value");
687 const std::vector<SKPutListener*>& put_listeners =
700 String path = ru.
doc[
"path"].as<String>();
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++) {
714 const char* path = ru.
doc[
"path"];
715 JsonObject value = ru.
doc.as<JsonObject>();
717 for (
size_t i = 0; i < listeners.size(); i++) {
724 for (
size_t i = 0; i < put_listeners.size(); i++) {
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;
756 const std::vector<SKPutListener*>& listeners =
758 for (
size_t j = 0; j < listeners.size(); j++) {
779 JsonDocument put_response;
780 put_response[
"requestId"] = message[
"requestId"];
782 put_response[
"state"] =
"COMPLETED";
783 put_response[
"statusCode"] = 200;
785 put_response[
"state"] =
"FAILED";
786 put_response[
"statusCode"] = 405;
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(),
794 ESP_LOGE(__FILENAME__,
"PUT response send failed (result=%d)", result);
808 int result = esp_websocket_client_send_text(
811 ESP_LOGE(__FILENAME__,
"sendTXT failed (result=%d)", result);
817 uint16_t& server_port) {
820 int num = MDNS.queryService(
"signalk-wss",
"tcp");
824 ESP_LOGI(__FILENAME__,
"Found Signal K server via mDNS (signalk-wss)");
827 num = MDNS.queryService(
"signalk-ws",
"tcp");
834 ESP_LOGI(__FILENAME__,
"Found Signal K server via mDNS (signalk-ws)");
837#if ESP_ARDUINO_VERSION_MAJOR < 3
838 server_address = MDNS.IP(0).toString();
840 server_address = MDNS.address(0).toString();
842 server_port = MDNS.port(0);
843 ESP_LOGI(__FILENAME__,
"Found server %s (port %d)", server_address.c_str(),
849static esp_err_t detect_ssl_event_handler(esp_http_client_event_t* evt) {
850 if (evt->event_id == HTTP_EVENT_ON_HEADER) {
852 if (strcasecmp(evt->header_key,
"Location") == 0) {
853 String* location =
static_cast<String*
>(evt->user_data);
854 *location = evt->header_value;
866 ESP_LOGD(__FILENAME__,
"Probing for SSL redirect at %s", url.c_str());
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;
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");
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);
888 ESP_LOGD(__FILENAME__,
"HTTP request failed: %s", esp_err_to_name(err));
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");
922 if (
client_.load() !=
nullptr) {
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.");
959 ESP_LOGE(__FILENAME__,
"connect worker spawn failed");
968 self->auth_job_running_.store(
false);
969 vTaskDelete(
nullptr);
973 ESP_LOGI(__FILENAME__,
"Initiating websocket connection with server...");
977 ESP_LOGE(__FILENAME__,
978 "No Signal K server found in network when using mDNS service!");
980 ESP_LOGI(__FILENAME__,
981 "Signal K server has been found at address %s:%d by mDNS.",
990 ESP_LOGD(__FILENAME__,
991 "Websocket is connecting to Signal K server on address %s:%d",
1000 ESP_LOGD(__FILENAME__,
1001 "Websocket is not connecting to Signal K server because host and "
1002 "port are not defined.");
1007 if (this->
polling_href_.length() > 0 && this->polling_href_.startsWith(
"/")) {
1016 ESP_LOGD(__FILENAME__,
"No prior authorization token present.");
1022#ifdef SENSESP_SSL_SUPPORT
1044#ifndef SENSESP_SSL_SUPPORT
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());
1051 const String full_token = String(
"Bearer ") +
auth_token_;
1052 ESP_LOGD(__FILENAME__,
"Authorization: %.8s...[redacted]", full_token.c_str());
1054 esp_http_client_config_t config = {};
1055 config.url = url.c_str();
1056 config.timeout_ms = 10000;
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");
1065 esp_http_client_set_header(client,
"Authorization", full_token.c_str());
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);
1077 int content_length = esp_http_client_fetch_headers(client);
1078 int http_code = esp_http_client_get_status_code(client);
1080 ESP_LOGD(__FILENAME__,
"Testing resulted in http status %d", http_code);
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);
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;
1102 esp_http_client_close(client);
1103 esp_http_client_cleanup(client);
1105 if (payload.length() > 0) {
1106 ESP_LOGD(__FILENAME__,
"Returned payload (%d bytes): %s", payload.length(),
1110 if (http_code == 426) {
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) {
1118 ESP_LOGW(__FILENAME__,
"Token rejected (401), requesting new access");
1122 }
else if (http_code > 0) {
1125 ESP_LOGE(__FILENAME__,
"HTTP request failed with code %d", http_code);
1132 const uint16_t server_port) {
1133 ESP_LOGD(__FILENAME__,
"Sending access request (client_id=%s, ssl=%d)",
1144 doc[
"description"] =
1146 doc[
"permissions"] = kRequestPermission;
1147 String json_req =
"";
1148 serializeJson(doc, json_req);
1150 ESP_LOGD(__FILENAME__,
"Access request: %s", json_req.c_str());
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());
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
1163 config.crt_bundle_attach = tofu_crt_bundle_attach;
1164 config.skip_cert_common_name_check =
true;
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");
1176 esp_http_client_set_header(client,
"Content-Type",
"application/json");
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);
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);
1197 int content_length = esp_http_client_fetch_headers(client);
1198 int http_code = esp_http_client_get_status_code(client);
1200 ESP_LOGD(__FILENAME__,
"HTTP response: code=%d, content_length=%d", http_code, content_length);
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;
1211 ESP_LOGD(__FILENAME__,
"Response payload (%d bytes): %s",
1212 payload.length(), payload.c_str());
1214 esp_http_client_close(client);
1215 esp_http_client_cleanup(client);
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>() :
"";
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());
1231 if (http_code == 400 && href.length() > 0 && href.startsWith(
"/")) {
1232 ESP_LOGI(__FILENAME__,
"Existing request found, will poll href: %s", href.c_str());
1241 if (http_code == 202 && href.length() > 0 && href.startsWith(
"/")) {
1250 if (http_code == 404) {
1251 ESP_LOGI(__FILENAME__,
1252 "Server security disabled (404 on access request) — connecting "
1255 this->
connect_ws(server_address, server_port);
1260 ESP_LOGW(__FILENAME__,
"Cannot handle response: http=%d, state=%s", http_code, state.c_str());
1265 const uint16_t server_port,
1266 const String href) {
1267 ESP_LOGD(__FILENAME__,
"Polling SK Server for authentication token");
1269 String protocol =
ssl_enabled_ ?
"https://" :
"http://";
1270 String url = protocol + server_address +
":" + server_port + href;
1272 esp_http_client_config_t config = {};
1273 config.url = url.c_str();
1274 config.timeout_ms = 10000;
1275#ifdef SENSESP_SSL_SUPPORT
1277 config.crt_bundle_attach = tofu_crt_bundle_attach;
1278 config.skip_cert_common_name_check =
true;
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");
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);
1298 int content_length = esp_http_client_fetch_headers(client);
1299 int http_code = esp_http_client_get_status_code(client);
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);
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;
1323 ESP_LOGD(__FILENAME__,
"Poll response: http=%d, %d bytes", http_code,
1324 static_cast<int>(payload.length()));
1326 esp_http_client_close(client);
1327 esp_http_client_cleanup(client);
1329 if (http_code == 200 || http_code == 202) {
1331 auto error = deserializeJson(doc, payload.c_str());
1333 ESP_LOGW(__FILENAME__,
"WARNING: Could not deserialize http payload.");
1334 ESP_LOGW(__FILENAME__,
"DeserializationError: %s", error.c_str());
1338 String state = doc[
"state"];
1339 ESP_LOGD(__FILENAME__,
"%s", state.c_str());
1340 if (state ==
"PENDING") {
1344 if (state ==
"COMPLETED") {
1345 JsonObject access_req = doc[
"accessRequest"];
1346 String permission = access_req[
"permission"];
1351 if (permission ==
"DENIED") {
1352 ESP_LOGW(__FILENAME__,
"Permission denied");
1357 if (permission ==
"APPROVED") {
1358 ESP_LOGI(__FILENAME__,
"Permission granted");
1359 String token = access_req[
"token"];
1362 this->
connect_ws(server_address, server_port);
1367 if (http_code == 404 || http_code == 500) {
1372 ESP_LOGD(__FILENAME__,
1373 "Got %d polling access request — clearing stale href.",
1381 ESP_LOGW(__FILENAME__,
1382 "Can't handle response %d to pending access request.\n",
1398 ESP_LOGW(__FILENAME__,
"connect_ws: prior client not reaped; deferring");
1412 String path =
"/signalk/v1/stream?subscribe=none";
1414 String url = protocol +
"://" + host +
":" + String(port) + path;
1416 ESP_LOGD(__FILENAME__,
"Connecting WebSocket to %s", url.c_str());
1421 auth_header = String(
"Authorization: Bearer ") +
auth_token_ +
"\r\n";
1425 esp_websocket_client_config_t config = {};
1426 config.uri = url.c_str();
1429 if (auth_header.length() > 0) {
1430 config.headers = auth_header.c_str();
1433#ifdef SENSESP_SSL_SUPPORT
1437 config.crt_bundle_attach = tofu_crt_bundle_attach;
1438 config.skip_cert_common_name_check =
true;
1442 esp_websocket_client_handle_t h = esp_websocket_client_init(&config);
1444 ESP_LOGE(__FILENAME__,
"Failed to initialize WebSocket client");
1453 esp_websocket_register_events(
1454 h, WEBSOCKET_EVENT_ANY, websocket_event_handler,
1455 reinterpret_cast<void*
>(
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));
1465 esp_websocket_client_destroy(h);
1470 ESP_LOGD(__FILENAME__,
"WebSocket client started, waiting for connection...");
1478struct WsTeardownArg {
1485 auto* a =
static_cast<WsTeardownArg*
>(arg);
1489 esp_websocket_client_stop(a->handle);
1490 esp_websocket_client_destroy(a->handle);
1491 a->self->teardown_in_progress_.store(
false);
1493 vTaskDelete(
nullptr);
1497 if (old ==
nullptr) {
1500 auto* arg =
new WsTeardownArg{old,
this};
1502 nullptr) == pdPASS) {
1514 ESP_LOGW(__FILENAME__,
"teardown task spawn failed; will retry next cycle");
1520 esp_websocket_client_handle_t old =
client_.exchange(
nullptr);
1521 if (old ==
nullptr) {
1541 std::vector<String> deltas;
1544 for (
const auto& delta : deltas) {
1556 uint32_t now = millis();
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(),
1569 int send_result = esp_websocket_client_send_text(
1571 if (send_result < 0) {
1586 ESP_LOGW(__FILENAME__,
1587 "Delta send incomplete (result=%d); dropping rest of batch",
1621 if (config[
"sk_address"].is<String>()) {
1624 if (config[
"sk_port"].is<int>()) {
1627 if (config[
"use_mdns"].is<bool>()) {
1628 this->
use_mdns_ = config[
"use_mdns"].as<
bool>();
1630 if (config[
"token"].is<String>()) {
1633 if (config[
"client_id"].is<String>()) {
1634 this->
client_id_ = config[
"client_id"].as<String>();
1636 if (config[
"polling_href"].is<String>()) {
1637 String href = config[
"polling_href"].as<String>();
1642 if (config[
"ssl_enabled"].is<bool>()) {
1645 if (config[
"tofu_enabled"].is<bool>()) {
1651 if (config[
"tofu_fingerprint"].is<String>()) {
1654 if (config[
"tofu_ca_pem"].is<String>()) {
1655 this->
tofu_ca_pem_ = config[
"tofu_ca_pem"].as<String>();
1657 if (config[
"tofu_san"].is<String>()) {
1658 this->
tofu_san_ = config[
"tofu_san"].as<String>();
1660 if (config[
"tofu_pin_cn"].is<String>()) {
1661 this->
tofu_pin_cn_ = config[
"tofu_pin_cn"].as<String>();
1663 if (config[
"tofu_pin_is_ca"].is<bool>()) {
1666 if (config[
"send_meta_enabled"].is<bool>()) {
1682 return "Authorizing with SignalK";
1686 return "Connecting";
1688 return "Disconnected";
1690 return "Certificate verification failed";
virtual bool load() override
Load and populate the object from a persistent storage.
virtual bool save() override
Save the object to a persistent storage.
virtual void set(const C &input) override final
int attach(std::function< void()> observer)
Attach an observer callback.
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 schedule_reconnect()
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_
bool is_connect_due() const
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_
void clear_pending_tofu()
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.
uint16_t conf_server_port_
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).
void reset_reconnect_interval()
std::shared_ptr< SKDeltaQueue > sk_delta_queue_
String conf_server_address_
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.
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()
String generate_uuid4()
Generate a random UUIDv4 string.
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,...
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
#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.