File
Blob: archive/sfu-bringup/firmware/main/peer_probe.c
| 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 | |
| 10 | static esp_peer_handle_t peer; |
| 11 | static QueueHandle_t commands; |
| 12 | static char channel_labels[4][33]; |
| 13 | static unsigned channel_count; |
| 14 | |
| 15 | static cJSON *event(const char *name) |
| 16 | { |
| 17 | cJSON *msg = cJSON_CreateObject(); |
| 18 | cJSON_AddStringToObject(msg, "event", name); |
| 19 | return msg; |
| 20 | } |
| 21 | |
| 22 | static 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 | |
| 30 | static 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 | |
| 42 | static 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 | |
| 67 | static 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 | |
| 76 | static 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 | |
| 85 | static 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 | |
| 168 | static 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 | |
| 181 | void 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, ©, pdMS_TO_TICKS(100))) cJSON_Delete(copy); |
| 189 | } |