Skip to content
File

Blob: archive/sfu-bringup/firmware/main/peer_probe.c

c190 lines
1#include <stdlib.h>
2#include <string.h>
3#include "esp_peer.h"
4#include "esp_peer_default.h"
5#include "freertos/FreeRTOS.h"
6#include "freertos/queue.h"
7#include "freertos/task.h"
8#include "peer_probe.h"
9 
10static esp_peer_handle_t peer;
11static QueueHandle_t commands;
12static char channel_labels[4][33];
13static unsigned channel_count;
14 
15static cJSON *event(const char *name)
16{
17 cJSON *msg = cJSON_CreateObject();
18 cJSON_AddStringToObject(msg, "event", name);
19 return msg;
20}
21 
22static int on_state(esp_peer_state_t state, void *ctx)
23{
24 cJSON *msg = event("peer_state");
25 cJSON_AddNumberToObject(msg, "state", state);
26 probe_emit(msg);
27 return 0;
28}
29 
30static int on_signal(esp_peer_msg_t *signal, void *ctx)
31{
32 char *text = calloc(1, signal->size + 1);
33 if (!text) return -1;
34 memcpy(text, signal->data, signal->size);
35 cJSON *msg = event(signal->type == ESP_PEER_MSG_TYPE_SDP ? "sdp" : "candidate");
36 cJSON_AddStringToObject(msg, "text", text);
37 probe_emit(msg);
38 free(text);
39 return 0;
40}
41 
42static int on_data(esp_peer_data_frame_t *frame, void *ctx)
43{
44 cJSON *msg = event("data_received");
45 cJSON_AddNumberToObject(msg, "stream_id", frame->stream_id);
46 cJSON_AddNumberToObject(msg, "bytes", frame->size);
47 char *text = calloc(1, frame->size + 1);
48 if (text) {
49 memcpy(text, frame->data, frame->size);
50 cJSON_AddStringToObject(msg, "text", text);
51 cJSON *command = cJSON_Parse(text);
52 const cJSON *color = command ? cJSON_GetObjectItemCaseSensitive(command, "led") : NULL;
53 if (cJSON_IsArray(color) && cJSON_GetArraySize(color) == 3) {
54 cJSON *r = cJSON_GetArrayItem(color, 0);
55 cJSON *g = cJSON_GetArrayItem(color, 1);
56 cJSON *b = cJSON_GetArrayItem(color, 2);
57 if (cJSON_IsNumber(r) && cJSON_IsNumber(g) && cJSON_IsNumber(b))
58 probe_set_led(r->valueint, g->valueint, b->valueint);
59 }
60 cJSON_Delete(command);
61 free(text);
62 }
63 probe_emit(msg);
64 return 0;
65}
66 
67static int on_channel(esp_peer_data_channel_info_t *channel, void *ctx)
68{
69 cJSON *msg = event("channel_open");
70 cJSON_AddNumberToObject(msg, "stream_id", channel->stream_id);
71 cJSON_AddStringToObject(msg, "label", channel->label ? channel->label : "");
72 probe_emit(msg);
73 return 0;
74}
75 
76static int on_channel_close(esp_peer_data_channel_info_t *channel, void *ctx)
77{
78 cJSON *msg = event("channel_closed");
79 cJSON_AddNumberToObject(msg, "stream_id", channel->stream_id);
80 cJSON_AddStringToObject(msg, "label", channel->label ? channel->label : "");
81 probe_emit(msg);
82 return 0;
83}
84 
85static void execute(cJSON *command)
86{
87 const cJSON *name = cJSON_GetObjectItemCaseSensitive(command, "cmd");
88 if (!cJSON_IsString(name)) return;
89 int result = ESP_PEER_ERR_INVALID_ARG;
90 if (!strcmp(name->valuestring, "peer_init") && !peer) {
91 esp_peer_default_cfg_t extra = {
92 .agent_recv_timeout = 10,
93 .data_ch_cfg = {.send_cache_size = 16384, .recv_cache_size = 16384},
94 .rtp_cfg = {.send_pool_size = 8192, .send_queue_num = 16},
95 };
96 esp_peer_cfg_t config = {
97 .role = cJSON_IsTrue(cJSON_GetObjectItemCaseSensitive(command, "initiator"))
98 ? ESP_PEER_ROLE_CONTROLLING : ESP_PEER_ROLE_CONTROLLED,
99 .audio_dir = ESP_PEER_MEDIA_DIR_NONE,
100 .video_dir = ESP_PEER_MEDIA_DIR_NONE,
101 .enable_data_channel = true,
102 .manual_ch_create = true,
103 .no_auto_reconnect = true,
104 .on_state = on_state,
105 .on_msg = on_signal,
106 .on_data = on_data,
107 .on_channel_open = on_channel,
108 .on_channel_close = on_channel_close,
109 .extra_cfg = &extra,
110 .extra_size = sizeof(extra),
111 };
112 result = esp_peer_open(&config, esp_peer_get_default_impl(), &peer);
113 channel_count = 0;
114 if (!result) result = esp_peer_new_connection(peer);
115 } else if (!strcmp(name->valuestring, "sdp") && peer) {
116 const cJSON *text = cJSON_GetObjectItemCaseSensitive(command, "text");
117 if (cJSON_IsString(text)) {
118 esp_peer_msg_t signal = {.type = ESP_PEER_MSG_TYPE_SDP,
119 .data = (uint8_t *)text->valuestring, .size = strlen(text->valuestring)};
120 result = esp_peer_send_msg(peer, &signal);
121 }
122 } else if (!strcmp(name->valuestring, "send") && peer) {
123 const cJSON *text = cJSON_GetObjectItemCaseSensitive(command, "text");
124 const cJSON *id = cJSON_GetObjectItemCaseSensitive(command, "stream_id");
125 if (cJSON_IsString(text) && cJSON_IsNumber(id) && id->valueint >= 0 && id->valueint < 65535) {
126 esp_peer_data_frame_t frame = {.type = ESP_PEER_DATA_CHANNEL_STRING,
127 .stream_id = id->valueint, .data = (uint8_t *)text->valuestring, .size = strlen(text->valuestring)};
128 result = esp_peer_send_data(peer, &frame);
129 }
130 } else if (!strcmp(name->valuestring, "create_channel") && peer) {
131 const cJSON *label = cJSON_GetObjectItemCaseSensitive(command, "label");
132 const cJSON *ordered = cJSON_GetObjectItemCaseSensitive(command, "ordered");
133 const cJSON *retransmits = cJSON_GetObjectItemCaseSensitive(command, "max_retransmits");
134 const char *text = cJSON_IsString(label) ? label->valuestring : "robot";
135 if (channel_count < 4 && strlen(text) <= 32 &&
136 (!retransmits || (cJSON_IsNumber(retransmits) && retransmits->valueint >= 0 && retransmits->valueint <= 65535))) {
137 strcpy(channel_labels[channel_count], text);
138 esp_peer_data_channel_cfg_t channel = {
139 .type = retransmits ? ESP_PEER_DATA_CHANNEL_PARTIAL_RELIABLE_RETX : ESP_PEER_DATA_CHANNEL_RELIABLE,
140 .ordered = !cJSON_IsFalse(ordered),
141 .label = channel_labels[channel_count],
142 .max_retransmit_count = retransmits ? retransmits->valueint : 0,
143 };
144 result = esp_peer_create_data_channel(peer, &channel);
145 if (!result) channel_count++;
146 }
147 } else if (!strcmp(name->valuestring, "close_channel") && peer) {
148 const cJSON *label = cJSON_GetObjectItemCaseSensitive(command, "label");
149 if (cJSON_IsString(label)) result = esp_peer_close_data_channel(peer, label->valuestring);
150 } else if (!strcmp(name->valuestring, "peer_close") && peer) {
151 result = esp_peer_close(peer);
152 peer = NULL;
153 } else if (!strcmp(name->valuestring, "led")) {
154 const cJSON *r = cJSON_GetObjectItemCaseSensitive(command, "r");
155 const cJSON *g = cJSON_GetObjectItemCaseSensitive(command, "g");
156 const cJSON *b = cJSON_GetObjectItemCaseSensitive(command, "b");
157 if (cJSON_IsNumber(r) && cJSON_IsNumber(g) && cJSON_IsNumber(b)) {
158 probe_set_led(r->valueint, g->valueint, b->valueint);
159 result = 0;
160 }
161 } else if (!strcmp(name->valuestring, "ping")) result = 0;
162 cJSON *reply = event("command_result");
163 cJSON_AddStringToObject(reply, "cmd", name->valuestring);
164 cJSON_AddNumberToObject(reply, "result", result);
165 probe_emit(reply);
166}
167 
168static void run(void *arg)
169{
170 for (;;) {
171 cJSON *command = NULL;
172 if (xQueueReceive(commands, &command, peer ? 0 : pdMS_TO_TICKS(20))) {
173 execute(command);
174 cJSON_Delete(command);
175 }
176 if (peer) esp_peer_main_loop(peer);
177 vTaskDelay(1);
178 }
179}
180 
181void probe_command(const cJSON *command)
182{
183 if (!commands) {
184 commands = xQueueCreate(8, sizeof(cJSON *));
185 if (!commands || xTaskCreate(run, "peer_probe", 16384, NULL, 5, NULL) != pdPASS) abort();
186 }
187 cJSON *copy = cJSON_Duplicate(command, true);
188 if (copy && !xQueueSend(commands, &copy, pdMS_TO_TICKS(100))) cJSON_Delete(copy);
189}