File
Blob: archive/fft-benchmark/firmware/platform/peer.c
| 1 | #include <stdatomic.h> |
| 2 | #include <stdlib.h> |
| 3 | #include <string.h> |
| 4 | #include "radio_bridge.h" |
| 5 | #include "spectrum_benchmark.h" |
| 6 | #include "esp_peer.h" |
| 7 | #include "esp_peer_default.h" |
| 8 | #include "freertos/FreeRTOS.h" |
| 9 | #include "freertos/queue.h" |
| 10 | #include "freertos/semphr.h" |
| 11 | |
| 12 | #define SDP_LIMIT 16000 |
| 13 | #define COMMAND_LIMIT 512 |
| 14 | typedef struct { uint32_t length; uint8_t bytes[COMMAND_LIMIT]; } command; |
| 15 | |
| 16 | /* Process-lifetime ownership: the vendor close path hangs on the pinned SDK. |
| 17 | * No C callback calls Rust or touches Rust memory. Queue elements are copies. |
| 18 | */ |
| 19 | struct radio_peer { |
| 20 | esp_peer_handle_t handle; |
| 21 | esp_peer_cfg_t config; |
| 22 | esp_peer_default_cfg_t extra; |
| 23 | esp_peer_data_channel_cfg_t channels[2]; |
| 24 | QueueHandle_t commands; |
| 25 | SemaphoreHandle_t offer_lock; |
| 26 | atomic_int state; |
| 27 | atomic_uint dropped; |
| 28 | bool channels_created; |
| 29 | size_t offer_length; |
| 30 | char offer[SDP_LIMIT + 1]; |
| 31 | char answer[SDP_LIMIT + 1]; |
| 32 | uint8_t audio[1275]; |
| 33 | uint8_t data[1024]; |
| 34 | }; |
| 35 | static atomic_bool taken; |
| 36 | |
| 37 | static int on_state(esp_peer_state_t state, void *ctx) |
| 38 | { |
| 39 | radio_peer *p = ctx; |
| 40 | if (state == ESP_PEER_STATE_DATA_CHANNEL_CONNECTED) atomic_store(&p->state, 1); |
| 41 | if (state == ESP_PEER_STATE_DISCONNECTED || state == ESP_PEER_STATE_CLOSED || |
| 42 | state == ESP_PEER_STATE_CONNECT_FAILED || state == ESP_PEER_STATE_DATA_CHANNEL_DISCONNECTED) |
| 43 | atomic_store(&p->state, -1); |
| 44 | return 0; |
| 45 | } |
| 46 | |
| 47 | static int on_signal(esp_peer_msg_t *msg, void *ctx) |
| 48 | { |
| 49 | radio_peer *p = ctx; |
| 50 | if (msg->type != ESP_PEER_MSG_TYPE_SDP) return 0; |
| 51 | if (!msg->data || msg->size <= 0 || msg->size > SDP_LIMIT) { atomic_store(&p->state, -1); return -1; } |
| 52 | xSemaphoreTake(p->offer_lock, portMAX_DELAY); |
| 53 | memcpy(p->offer, msg->data, msg->size); |
| 54 | p->offer[msg->size] = 0; |
| 55 | p->offer_length = msg->size; |
| 56 | xSemaphoreGive(p->offer_lock); |
| 57 | return 0; |
| 58 | } |
| 59 | |
| 60 | static int on_data(esp_peer_data_frame_t *frame, void *ctx) |
| 61 | { |
| 62 | radio_peer *p = ctx; |
| 63 | if (frame->stream_id != 2 || !frame->data || frame->size <= 0 || frame->size > COMMAND_LIMIT) return 0; |
| 64 | command msg = {.length = frame->size}; |
| 65 | memcpy(msg.bytes, frame->data, frame->size); |
| 66 | if (!xQueueSend(p->commands, &msg, 0)) atomic_fetch_add(&p->dropped, 1); |
| 67 | return 0; |
| 68 | } |
| 69 | |
| 70 | radio_peer *radio_peer_open(void) |
| 71 | { |
| 72 | if (atomic_exchange(&taken, true)) return NULL; |
| 73 | spectrum_benchmark_start(); |
| 74 | radio_peer *p = calloc(1, sizeof(*p)); |
| 75 | if (!p) return NULL; |
| 76 | p->commands = xQueueCreate(16, sizeof(command)); |
| 77 | p->offer_lock = xSemaphoreCreateMutex(); |
| 78 | if (!p->commands || !p->offer_lock) goto early_error; |
| 79 | atomic_init(&p->state, 0); |
| 80 | atomic_init(&p->dropped, 0); |
| 81 | /* 1 ms caused DTLS receive timeouts to be interpreted as a peer close. |
| 82 | * Retain the measured 10 ms setting and the validated transport buffers. */ |
| 83 | p->extra = (esp_peer_default_cfg_t){.agent_recv_timeout = 10, |
| 84 | .data_ch_cfg = {.send_cache_size = 16384, .recv_cache_size = 16384}, |
| 85 | .rtp_cfg = {.send_pool_size = 32768, .send_queue_num = 32}}; |
| 86 | p->config = (esp_peer_cfg_t){.role = ESP_PEER_ROLE_CONTROLLING, |
| 87 | .audio_info = {.codec = ESP_PEER_AUDIO_CODEC_OPUS, .sample_rate = 48000, .channel = 2}, |
| 88 | .audio_dir = ESP_PEER_MEDIA_DIR_SEND_ONLY, .video_dir = ESP_PEER_MEDIA_DIR_NONE, |
| 89 | .enable_data_channel = true, .manual_ch_create = true, .no_auto_reconnect = true, |
| 90 | .on_state = on_state, .on_msg = on_signal, .on_data = on_data, .ctx = p, |
| 91 | .extra_cfg = &p->extra, .extra_size = sizeof(p->extra)}; |
| 92 | /* All configuration and context remain live even if initialization fails. */ |
| 93 | if (esp_peer_open(&p->config, esp_peer_get_default_impl(), &p->handle) || esp_peer_new_connection(p->handle)) return NULL; |
| 94 | return p; |
| 95 | early_error: |
| 96 | if (p->commands) vQueueDelete(p->commands); |
| 97 | if (p->offer_lock) vSemaphoreDelete(p->offer_lock); |
| 98 | free(p); |
| 99 | return NULL; |
| 100 | } |
| 101 | |
| 102 | int32_t radio_peer_poll(radio_peer *p) { return esp_peer_main_loop(p->handle); } |
| 103 | int32_t radio_peer_state(radio_peer *p) { return atomic_load(&p->state); } |
| 104 | uint32_t radio_peer_dropped(radio_peer *p) { return atomic_load(&p->dropped); } |
| 105 | int32_t radio_peer_offer(radio_peer *p, uint8_t *out, size_t capacity) |
| 106 | { |
| 107 | xSemaphoreTake(p->offer_lock, portMAX_DELAY); |
| 108 | int32_t size = p->offer_length; |
| 109 | if ((size_t)size > capacity) size = -1; |
| 110 | else if (size) { memcpy(out, p->offer, size); p->offer_length = 0; } |
| 111 | xSemaphoreGive(p->offer_lock); |
| 112 | return size; |
| 113 | } |
| 114 | int32_t radio_peer_answer(radio_peer *p, const uint8_t *bytes, size_t length) |
| 115 | { |
| 116 | if (!length || length > SDP_LIMIT || p->answer[0]) return -1; |
| 117 | memcpy(p->answer, bytes, length); |
| 118 | p->answer[length] = 0; |
| 119 | esp_peer_msg_t msg = {.type = ESP_PEER_MSG_TYPE_SDP, .data = (uint8_t *)p->answer, .size = length}; |
| 120 | return esp_peer_send_msg(p->handle, &msg); |
| 121 | } |
| 122 | int32_t radio_peer_channels(radio_peer *p) |
| 123 | { |
| 124 | if (p->channels_created || radio_peer_state(p) != 1) return -1; |
| 125 | p->channels[0] = (esp_peer_data_channel_cfg_t){.label = "robot", .type = ESP_PEER_DATA_CHANNEL_RELIABLE, .ordered = true}; |
| 126 | p->channels[1] = (esp_peer_data_channel_cfg_t){.label = "spectrum", .type = ESP_PEER_DATA_CHANNEL_PARTIAL_RELIABLE_RETX, |
| 127 | .ordered = false, .max_retransmit_count = 0}; |
| 128 | for (unsigned i = 0; i < 2; i++) { int err = esp_peer_create_data_channel(p->handle, &p->channels[i]); if (err) return err; } |
| 129 | p->channels_created = true; |
| 130 | return 0; |
| 131 | } |
| 132 | int32_t radio_peer_command(radio_peer *p, uint8_t *out, size_t capacity) |
| 133 | { |
| 134 | command msg; |
| 135 | if (!xQueueReceive(p->commands, &msg, 0)) return 0; |
| 136 | if (msg.length > capacity) return -1; |
| 137 | memcpy(out, msg.bytes, msg.length); |
| 138 | return msg.length; |
| 139 | } |
| 140 | int32_t radio_peer_audio(radio_peer *p, uint32_t pts, const uint8_t *bytes, size_t length) |
| 141 | { |
| 142 | if (!length || length > sizeof(p->audio)) return -1; |
| 143 | spectrum_benchmark_push(pts, bytes, length); |
| 144 | memcpy(p->audio, bytes, length); |
| 145 | esp_peer_audio_frame_t frame = {.pts = pts, .data = p->audio, .size = length}; |
| 146 | return esp_peer_send_audio(p->handle, &frame); |
| 147 | } |
| 148 | int32_t radio_peer_data(radio_peer *p, uint16_t stream, int32_t binary, const uint8_t *bytes, size_t length) |
| 149 | { |
| 150 | if (!length || length > sizeof(p->data) || (stream != 2 && stream != 4)) return -1; |
| 151 | memcpy(p->data, bytes, length); |
| 152 | esp_peer_data_frame_t frame = {.type = binary ? ESP_PEER_DATA_CHANNEL_DATA : ESP_PEER_DATA_CHANNEL_STRING, |
| 153 | .stream_id = stream, .data = p->data, .size = length}; |
| 154 | return esp_peer_send_data(p->handle, &frame); |
| 155 | } |