Skip to content
File

Blob: archive/legacy-c/firmware/main/signaling.c

c277 lines
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 
23static const char *TAG = "signaling";
24static QueueHandle_t events;
25static atomic_bool transport_lost;
26static RTC_DATA_ATTR unsigned recovery_attempt;
27static esp_http_client_handle_t client;
28static char *response_buffer;
29static size_t response_used;
30static bool response_overflow;
31#define RESPONSE_LIMIT 24576
32 
33bool 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, &copy, 0)) cJSON_Delete(copy);
50 }
51 return sdp || candidate;
52}
53 
54static void drain_events(void)
55{
56 cJSON *event;
57 while (xQueueReceive(events, &event, 0)) cJSON_Delete(event);
58}
59 
60static 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 
77static 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 
92static cJSON *command(const char *name)
93{
94 cJSON *object = cJSON_CreateObject();
95 cJSON_AddStringToObject(object, "cmd", name);
96 return object;
97}
98 
99static 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.
113static 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 
145static 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 
158static 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 
169static 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 
271void 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}