File
Blob: archive/legacy-c/firmware/main/signaling.c
| 1 | #include <stdatomic.h> |
| 2 | #include <stdio.h> |
| 3 | #include <stdlib.h> |
| 4 | #include <string.h> |
| 5 | #include <time.h> |
| 6 | #include "esp_attr.h" |
| 7 | #include "esp_crt_bundle.h" |
| 8 | #include "esp_heap_caps.h" |
| 9 | #include "esp_http_client.h" |
| 10 | #include "esp_log.h" |
| 11 | #include "esp_netif_sntp.h" |
| 12 | #include "esp_peer.h" |
| 13 | #include "esp_random.h" |
| 14 | #include "esp_system.h" |
| 15 | #include "esp_timer.h" |
| 16 | #include "freertos/FreeRTOS.h" |
| 17 | #include "freertos/queue.h" |
| 18 | #include "freertos/task.h" |
| 19 | #include "radio.h" |
| 20 | #include "signaling.h" |
| 21 | #include "signaling_private.h" |
| 22 | |
| 23 | static const char *TAG = "signaling"; |
| 24 | static QueueHandle_t events; |
| 25 | static atomic_bool transport_lost; |
| 26 | static RTC_DATA_ATTR unsigned recovery_attempt; |
| 27 | static esp_http_client_handle_t client; |
| 28 | static char *response_buffer; |
| 29 | static size_t response_used; |
| 30 | static bool response_overflow; |
| 31 | #define RESPONSE_LIMIT 24576 |
| 32 | |
| 33 | bool signaling_observe(const cJSON *message) |
| 34 | { |
| 35 | if (!SIGNALING_AUTONOMOUS || !events) return false; |
| 36 | const cJSON *kind = cJSON_GetObjectItemCaseSensitive(message, "event"); |
| 37 | if (!cJSON_IsString(kind)) return false; |
| 38 | bool sdp = !strcmp(kind->valuestring, "sdp"); |
| 39 | bool candidate = !strcmp(kind->valuestring, "candidate"); |
| 40 | bool state = !strcmp(kind->valuestring, "peer_state"); |
| 41 | if (state) { |
| 42 | const cJSON *value = cJSON_GetObjectItemCaseSensitive(message, "state"); |
| 43 | if (cJSON_IsNumber(value) && (value->valueint == ESP_PEER_STATE_DISCONNECTED || |
| 44 | value->valueint == ESP_PEER_STATE_CONNECT_FAILED || value->valueint == ESP_PEER_STATE_DATA_CHANNEL_DISCONNECTED)) |
| 45 | atomic_store(&transport_lost, true); |
| 46 | } |
| 47 | if (sdp || state || !strcmp(kind->valuestring, "command_result")) { |
| 48 | cJSON *copy = cJSON_Duplicate(message, true); |
| 49 | if (copy && !xQueueSend(events, ©, 0)) cJSON_Delete(copy); |
| 50 | } |
| 51 | return sdp || candidate; |
| 52 | } |
| 53 | |
| 54 | static void drain_events(void) |
| 55 | { |
| 56 | cJSON *event; |
| 57 | while (xQueueReceive(events, &event, 0)) cJSON_Delete(event); |
| 58 | } |
| 59 | |
| 60 | static cJSON *wait_event(const char *kind, const char *command, int state, unsigned timeout_ms) |
| 61 | { |
| 62 | int64_t deadline = esp_timer_get_time() + (int64_t)timeout_ms * 1000; |
| 63 | while (esp_timer_get_time() < deadline) { |
| 64 | cJSON *event = NULL; |
| 65 | if (!xQueueReceive(events, &event, pdMS_TO_TICKS(100))) continue; |
| 66 | const cJSON *name = cJSON_GetObjectItemCaseSensitive(event, "event"); |
| 67 | const cJSON *cmd = cJSON_GetObjectItemCaseSensitive(event, "cmd"); |
| 68 | const cJSON *value = cJSON_GetObjectItemCaseSensitive(event, "state"); |
| 69 | if (cJSON_IsString(name) && !strcmp(name->valuestring, kind) && |
| 70 | (!command || (cJSON_IsString(cmd) && !strcmp(cmd->valuestring, command))) && |
| 71 | (state < 0 || (cJSON_IsNumber(value) && value->valueint == state))) return event; |
| 72 | cJSON_Delete(event); |
| 73 | } |
| 74 | return NULL; |
| 75 | } |
| 76 | |
| 77 | static bool command_ok(cJSON *command) |
| 78 | { |
| 79 | const cJSON *name = cJSON_GetObjectItemCaseSensitive(command, "cmd"); |
| 80 | char expected[32]; |
| 81 | if (!cJSON_IsString(name) || strlen(name->valuestring) >= sizeof(expected)) { cJSON_Delete(command); return false; } |
| 82 | strcpy(expected, name->valuestring); |
| 83 | radio_command(command); |
| 84 | cJSON_Delete(command); |
| 85 | cJSON *reply = wait_event("command_result", expected, -1, 15000); |
| 86 | const cJSON *result = cJSON_GetObjectItemCaseSensitive(reply, "result"); |
| 87 | bool ok = cJSON_IsNumber(result) && result->valueint == 0; |
| 88 | cJSON_Delete(reply); |
| 89 | return ok; |
| 90 | } |
| 91 | |
| 92 | static cJSON *command(const char *name) |
| 93 | { |
| 94 | cJSON *object = cJSON_CreateObject(); |
| 95 | cJSON_AddStringToObject(object, "cmd", name); |
| 96 | return object; |
| 97 | } |
| 98 | |
| 99 | static esp_err_t http_event(esp_http_client_event_t *event) |
| 100 | { |
| 101 | if (event->event_id == HTTP_EVENT_ON_DATA && event->data_len > 0) { |
| 102 | if (response_used + event->data_len >= RESPONSE_LIMIT) response_overflow = true; |
| 103 | else if (!response_overflow) { |
| 104 | memcpy(response_buffer + response_used, event->data, event->data_len); |
| 105 | response_used += event->data_len; |
| 106 | } |
| 107 | } |
| 108 | return ESP_OK; |
| 109 | } |
| 110 | |
| 111 | // Positive values are HTTP statuses, negative values are local transport errors. |
| 112 | // The same client owns all requests, on this task only. It reuses verified TLS. |
| 113 | static int post(const char *path, const cJSON *body, cJSON **reply) |
| 114 | { |
| 115 | *reply = NULL; |
| 116 | if (!radio_network_ready()) return -1; |
| 117 | char url[256]; |
| 118 | if (snprintf(url, sizeof(url), "%s/api/device/%s", SIGNALING_URL, path) >= sizeof(url)) return -1; |
| 119 | char *json = cJSON_PrintUnformatted(body); |
| 120 | if (!json) return -1; |
| 121 | response_used = 0; |
| 122 | response_overflow = false; |
| 123 | esp_err_t err = esp_http_client_set_url(client, url); |
| 124 | if (err == ESP_OK) err = esp_http_client_set_post_field(client, json, strlen(json)); |
| 125 | if (err == ESP_OK) err = esp_http_client_perform(client); |
| 126 | int status = esp_http_client_get_status_code(client); |
| 127 | esp_http_client_set_post_field(client, NULL, 0); |
| 128 | free(json); |
| 129 | if (err != ESP_OK || response_overflow) { |
| 130 | ESP_LOGW(TAG, "HTTPS %s: transport=%s oversized=%d", path, esp_err_to_name(err), response_overflow); |
| 131 | esp_http_client_close(client); |
| 132 | return -1; |
| 133 | } |
| 134 | response_buffer[response_used] = 0; |
| 135 | *reply = cJSON_ParseWithLength(response_buffer, response_used); |
| 136 | if (!cJSON_IsObject(*reply)) { |
| 137 | cJSON_Delete(*reply); |
| 138 | *reply = NULL; |
| 139 | return status >= 400 ? status : -1; |
| 140 | } |
| 141 | ESP_LOGI(TAG, "HTTPS %s: status=%d", path, status); |
| 142 | return status; |
| 143 | } |
| 144 | |
| 145 | static cJSON *post_retry(const char *path, const cJSON *body) |
| 146 | { |
| 147 | for (unsigned attempt = 0; attempt < 3; attempt++) { |
| 148 | cJSON *reply = NULL; |
| 149 | int status = post(path, body, &reply); |
| 150 | if (status >= 200 && status < 300 && reply) return reply; |
| 151 | cJSON_Delete(reply); |
| 152 | if (status >= 400 && status < 500 && status != 429) break; |
| 153 | vTaskDelay(pdMS_TO_TICKS(1000u << attempt)); |
| 154 | } |
| 155 | return NULL; |
| 156 | } |
| 157 | |
| 158 | static void restart_later(const char *reason) |
| 159 | { |
| 160 | unsigned delay = 5u << (recovery_attempt < 3 ? recovery_attempt : 3); |
| 161 | recovery_attempt++; |
| 162 | ESP_LOGW(TAG, "%s; restarting in %u seconds", reason, delay); |
| 163 | vTaskDelay(pdMS_TO_TICKS(delay * 1000 + esp_random() % 1000)); |
| 164 | // The vendor peer teardown currently hangs. A restart safely resets its |
| 165 | // tasks and sockets; the Worker closes the previous publisher generation. |
| 166 | esp_restart(); |
| 167 | } |
| 168 | |
| 169 | static void run(void *unused) |
| 170 | { |
| 171 | esp_sntp_config_t clock = ESP_NETIF_SNTP_DEFAULT_CONFIG_MULTIPLE(2, |
| 172 | ESP_SNTP_SERVER_LIST("time.cloudflare.com", "pool.ntp.org")); |
| 173 | ESP_ERROR_CHECK(esp_netif_sntp_init(&clock)); |
| 174 | while (time(NULL) < 1700000000) { |
| 175 | if (esp_netif_sntp_sync_wait(pdMS_TO_TICKS(20000)) != ESP_OK) |
| 176 | ESP_LOGW(TAG, "Waiting for time synchronization before verified HTTPS"); |
| 177 | } |
| 178 | response_buffer = heap_caps_malloc(RESPONSE_LIMIT, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); |
| 179 | if (!response_buffer) { restart_later("No response buffer"); return; } |
| 180 | esp_http_client_config_t config = {.url = SIGNALING_URL, |
| 181 | .crt_bundle_attach = esp_crt_bundle_attach, .timeout_ms = 12000, |
| 182 | .event_handler = http_event, .buffer_size = 4096, .buffer_size_tx = 4096, |
| 183 | .disable_auto_redirect = true, .keep_alive_enable = true}; |
| 184 | client = esp_http_client_init(&config); |
| 185 | if (!client) { restart_later("HTTP initialization failed"); return; } |
| 186 | esp_http_client_set_method(client, HTTP_METHOD_POST); |
| 187 | esp_http_client_set_header(client, "Content-Type", "application/json"); |
| 188 | esp_http_client_set_header(client, "Authorization", "Bearer " SIGNALING_DEVICE_TOKEN); |
| 189 | esp_http_client_set_header(client, "User-Agent", "PocketRadio-S3/1.0"); |
| 190 | ESP_LOGI(TAG, "Starting autonomous signaling at %s", SIGNALING_URL); |
| 191 | if (!command_ok(command("peer_init"))) { restart_later("Radio initialization failed"); return; } |
| 192 | cJSON *offer = wait_event("sdp", NULL, -1, 15000); |
| 193 | const cJSON *sdp = cJSON_GetObjectItemCaseSensitive(offer, "text"); |
| 194 | if (!cJSON_IsString(sdp)) { cJSON_Delete(offer); restart_later("Missing radio offer"); return; } |
| 195 | cJSON *start = cJSON_CreateObject(); |
| 196 | cJSON *description = cJSON_AddObjectToObject(start, "sessionDescription"); |
| 197 | cJSON_AddStringToObject(description, "type", "offer"); |
| 198 | cJSON_AddStringToObject(description, "sdp", sdp->valuestring); |
| 199 | char boot_id[33]; |
| 200 | snprintf(boot_id, sizeof(boot_id), "%08lx%08lx%08lx%08lx", (unsigned long)esp_random(), |
| 201 | (unsigned long)esp_random(), (unsigned long)esp_random(), (unsigned long)esp_random()); |
| 202 | cJSON_AddStringToObject(start, "bootId", boot_id); |
| 203 | cJSON_Delete(offer); |
| 204 | cJSON *started = post_retry("start", start); |
| 205 | cJSON_Delete(start); |
| 206 | const cJSON *generation = cJSON_GetObjectItemCaseSensitive(started, "generation"); |
| 207 | const cJSON *answer = cJSON_GetObjectItemCaseSensitive(started, "sessionDescription"); |
| 208 | const cJSON *answer_text = cJSON_GetObjectItemCaseSensitive(answer, "sdp"); |
| 209 | if (!cJSON_IsString(generation) || strlen(generation->valuestring) > 64 || !cJSON_IsString(answer_text)) { |
| 210 | cJSON_Delete(started); restart_later("Publisher setup failed"); return; |
| 211 | } |
| 212 | cJSON *identity = cJSON_CreateObject(); |
| 213 | cJSON_AddStringToObject(identity, "generation", generation->valuestring); |
| 214 | cJSON *remote = command("sdp"); |
| 215 | cJSON_AddStringToObject(remote, "text", answer_text->valuestring); |
| 216 | cJSON_Delete(started); |
| 217 | if (!command_ok(remote)) { restart_later("Answer rejected"); return; } |
| 218 | cJSON *connected = wait_event("peer_state", NULL, ESP_PEER_STATE_DATA_CHANNEL_CONNECTED, 30000); |
| 219 | if (!connected) { restart_later("WebRTC connection timed out"); return; } |
| 220 | cJSON_Delete(connected); |
| 221 | cJSON *created = post_retry("channels", identity); |
| 222 | const cJSON *channels = cJSON_GetObjectItemCaseSensitive(created, "channels"); |
| 223 | int ids[2] = {-1, -1}; |
| 224 | const cJSON *channel; |
| 225 | cJSON_ArrayForEach(channel, channels) { |
| 226 | const cJSON *label = cJSON_GetObjectItemCaseSensitive(channel, "dataChannelName"); |
| 227 | const cJSON *id = cJSON_GetObjectItemCaseSensitive(channel, "id"); |
| 228 | if (cJSON_IsString(label) && cJSON_IsNumber(id)) { |
| 229 | if (!strcmp(label->valuestring, "robot")) ids[0] = id->valueint; |
| 230 | if (!strcmp(label->valuestring, "spectrum")) ids[1] = id->valueint; |
| 231 | } |
| 232 | } |
| 233 | cJSON_Delete(created); |
| 234 | if (ids[0] != 2 || ids[1] != 4) { restart_later("Unexpected data channel allocation"); return; } |
| 235 | for (unsigned i = 0; i < 2; i++) { |
| 236 | if (!command_ok(command("create_channel"))) { restart_later("Local channel creation failed"); return; } |
| 237 | } |
| 238 | cJSON *play = command("start"); |
| 239 | cJSON_AddNumberToObject(play, "robot_id", ids[0]); |
| 240 | cJSON_AddNumberToObject(play, "spectrum_id", ids[1]); |
| 241 | if (!command_ok(play)) { restart_later("Playback could not start"); return; } |
| 242 | cJSON *ready = post_retry("ready", identity); |
| 243 | if (!ready) { restart_later("Publisher readiness was not acknowledged"); return; } |
| 244 | cJSON_Delete(ready); |
| 245 | recovery_attempt = 0; |
| 246 | ESP_LOGI(TAG, "ON AIR: autonomous Wi-Fi signaling, audio and data; USB is optional"); |
| 247 | unsigned failures = 0; |
| 248 | int64_t last_success = esp_timer_get_time(); |
| 249 | for (;;) { |
| 250 | drain_events(); |
| 251 | for (unsigned i = 0; i < 50; i++) { |
| 252 | if (atomic_load(&transport_lost)) { restart_later("WebRTC transport disconnected"); return; } |
| 253 | vTaskDelay(pdMS_TO_TICKS(100)); |
| 254 | } |
| 255 | cJSON *reply = NULL; |
| 256 | int status = post("heartbeat", identity, &reply); |
| 257 | cJSON_Delete(reply); |
| 258 | if (status >= 200 && status < 300) { |
| 259 | failures = 0; |
| 260 | last_success = esp_timer_get_time(); |
| 261 | } else { |
| 262 | ESP_LOGW(TAG, "Heartbeat unavailable: status=%d attempt=%u", status, ++failures); |
| 263 | if ((status >= 400 && status < 500 && status != 429) || |
| 264 | esp_timer_get_time() - last_success > 55000000) { |
| 265 | restart_later("Signaling session needs recovery"); return; |
| 266 | } |
| 267 | } |
| 268 | } |
| 269 | } |
| 270 | |
| 271 | void signaling_start(void) |
| 272 | { |
| 273 | if (!SIGNALING_AUTONOMOUS) return; |
| 274 | events = xQueueCreate(16, sizeof(cJSON *)); |
| 275 | if (!events || xTaskCreate(run, "signaling", 16384, NULL, 4, NULL) != pdPASS) abort(); |
| 276 | } |