audio_stream_buffered.c 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278
  1. #include <errno.h>
  2. #include <netinet/in.h>
  3. #include <stdbool.h>
  4. #include <stdint.h>
  5. #include <stdlib.h>
  6. #include <string.h>
  7. #include <sys/socket.h>
  8. #include <unistd.h>
  9. #include "audio_receiver_internal.h"
  10. #include "esp_heap_caps.h"
  11. #include "esp_log.h"
  12. #include "audio_crypto.h"
  13. #include "network/socket_utils.h"
  14. #define BUFFERED_AUDIO_PACKET_SIZE 8192
  15. #define AUDIO_BUFFERED_STACK_SIZE 4096
  16. static const char *TAG = "audio_buf";
  17. // Read exact number of bytes, but keep waiting on timeout if paused
  18. // Returns: positive = bytes read, 0 = connection closed, -1 = error
  19. static ssize_t read_exact(audio_stream_t *stream, audio_receiver_state_t *state,
  20. int sock, uint8_t *buf, size_t len) {
  21. size_t total = 0;
  22. while (total < len && stream->running) {
  23. ssize_t n = recv(sock, buf + total, len - total, 0);
  24. if (n > 0) {
  25. total += (size_t)n;
  26. } else if (n == 0) {
  27. // Connection closed by peer
  28. ESP_LOGI(TAG, "Buffered audio connection closed by peer");
  29. return 0;
  30. } else {
  31. // n < 0: error or timeout
  32. if (errno == EAGAIN || errno == EWOULDBLOCK) {
  33. // Timeout - if we're paused, keep waiting for resume
  34. if (!state->timing.playing) {
  35. // Still paused, keep the connection alive
  36. vTaskDelay(pdMS_TO_TICKS(100));
  37. continue;
  38. }
  39. // Playing but timed out - connection may be dead
  40. ESP_LOGW(TAG, "Buffered audio timeout while playing");
  41. return -1;
  42. }
  43. ESP_LOGE(TAG, "Buffered audio recv error: %d", errno);
  44. return -1;
  45. }
  46. }
  47. return stream->running ? (ssize_t)total : -1;
  48. }
  49. static void buffered_audio_task(void *pvParameters) {
  50. audio_stream_t *stream = (audio_stream_t *)pvParameters;
  51. audio_receiver_state_t *state = audio_stream_state(stream);
  52. while (stream->running) {
  53. struct sockaddr_in client_addr;
  54. socklen_t addr_len = sizeof(client_addr);
  55. int client_sock = accept(state->buffered_listen_socket,
  56. (struct sockaddr *)&client_addr, &addr_len);
  57. if (client_sock < 0) {
  58. if (errno != EAGAIN && errno != EWOULDBLOCK && stream->running) {
  59. ESP_LOGE(TAG, "Buffered audio accept error: %d", errno);
  60. }
  61. vTaskDelay(pdMS_TO_TICKS(100));
  62. continue;
  63. }
  64. state->buffered_client_socket = client_sock;
  65. struct timeval tv = {.tv_sec = 30, .tv_usec = 0};
  66. setsockopt(client_sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
  67. // Socket receive buffer: match lwIP's TCP receive window so the kernel
  68. // buffer can hold exactly what the TCP window allows in flight. A larger
  69. // SO_RCVBUF (e.g. the old 65536) accumulates stale audio data that must
  70. // drain through the RTP gates on every track skip, adding transition
  71. // latency. Keeping it at TCP_WND ties both knobs to a single sdkconfig
  72. // value (CONFIG_LWIP_TCP_WND_DEFAULT).
  73. int rcvbuf = CONFIG_LWIP_TCP_WND_DEFAULT;
  74. setsockopt(client_sock, SOL_SOCKET, SO_RCVBUF, &rcvbuf, sizeof(rcvbuf));
  75. uint8_t *packet = state->buffered_recv_buffer;
  76. if (!packet) {
  77. packet = heap_caps_malloc(BUFFERED_AUDIO_PACKET_SIZE,
  78. MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT);
  79. if (!packet) {
  80. packet = malloc(BUFFERED_AUDIO_PACKET_SIZE);
  81. }
  82. if (!packet) {
  83. ESP_LOGE(TAG, "Failed to allocate buffered audio packet buffer");
  84. close(client_sock);
  85. state->buffered_client_socket = -1;
  86. continue;
  87. }
  88. state->buffered_recv_buffer = packet;
  89. }
  90. while (stream->running) {
  91. // Back-pressure: if buffer is nearly full, pause reading to let TCP
  92. // flow control slow down the sender. This prevents buffer overflow
  93. // and keeps frames in order.
  94. while (audio_buffer_is_nearly_full(&state->buffer) && stream->running) {
  95. vTaskDelay(pdMS_TO_TICKS(10));
  96. }
  97. uint8_t len_buf[2];
  98. if (read_exact(stream, state, client_sock, len_buf, 2) != 2) {
  99. break;
  100. }
  101. uint16_t data_len = (uint16_t)((len_buf[0] << 8) | len_buf[1]);
  102. if (data_len < 2 || data_len > BUFFERED_AUDIO_PACKET_SIZE) {
  103. ESP_LOGW(TAG, "Invalid buffered audio packet length: %u", data_len);
  104. break;
  105. }
  106. size_t packet_len = data_len - 2;
  107. if (read_exact(stream, state, client_sock, packet, packet_len) !=
  108. (ssize_t)packet_len) {
  109. break;
  110. }
  111. state->stats.packets_received++;
  112. uint32_t seq_no = (packet[1] << 16) | (packet[2] << 8) | packet[3];
  113. uint32_t timestamp =
  114. (packet[4] << 24) | (packet[5] << 16) | (packet[6] << 8) | packet[7];
  115. uint8_t *decrypted = state->decrypt_buffer;
  116. size_t decrypt_capacity = state->decrypt_buffer_size;
  117. if (!decrypted) {
  118. decrypted = packet + 12;
  119. decrypt_capacity = packet_len > 12 ? packet_len - 12 : 0;
  120. }
  121. int decrypted_len = audio_crypto_decrypt_buffered(
  122. &stream->encrypt, packet, packet_len, decrypted, decrypt_capacity);
  123. if (decrypted_len < 0) {
  124. state->stats.decrypt_errors++;
  125. state->stats.packets_dropped++;
  126. continue;
  127. }
  128. state->stats.last_seq = (uint16_t)(seq_no & 0xFFFF);
  129. state->stats.last_timestamp = timestamp;
  130. state->blocks_read++;
  131. state->blocks_read_in_sequence++;
  132. if (!audio_stream_process_frame(state, timestamp, decrypted,
  133. (size_t)decrypted_len)) {
  134. state->stats.packets_dropped++;
  135. }
  136. }
  137. close(client_sock);
  138. state->buffered_client_socket = -1;
  139. }
  140. state->buffered_task_handle = NULL;
  141. vTaskDelete(NULL);
  142. }
  143. static bool buffered_wait_for_task_stopped(audio_receiver_state_t *state,
  144. int timeout_ticks) {
  145. while (state->buffered_task_handle && timeout_ticks-- > 0) {
  146. vTaskDelay(pdMS_TO_TICKS(50));
  147. }
  148. return state->buffered_task_handle == NULL;
  149. }
  150. static esp_err_t buffered_start(audio_stream_t *stream, uint16_t port) {
  151. audio_receiver_state_t *state = audio_stream_state(stream);
  152. if (stream->running) {
  153. ESP_LOGI(TAG, "Buffered audio already running, continuing");
  154. return ESP_OK;
  155. }
  156. if (state->buffered_task_handle) {
  157. ESP_LOGW(TAG, "Buffered audio task still stopping, waiting");
  158. if (!buffered_wait_for_task_stopped(state, 20)) {
  159. ESP_LOGW(TAG, "Buffered audio task still active");
  160. return ESP_ERR_INVALID_STATE;
  161. }
  162. }
  163. uint16_t bound_port = port;
  164. state->buffered_listen_socket =
  165. socket_utils_bind_tcp_listener(port, 1, true, &bound_port);
  166. if (state->buffered_listen_socket < 0) {
  167. return ESP_FAIL;
  168. }
  169. state->buffered_port = bound_port;
  170. stream->running = true;
  171. state->buffered_task_handle = NULL;
  172. BaseType_t task_ret =
  173. xTaskCreate(buffered_audio_task, "buff_audio", AUDIO_BUFFERED_STACK_SIZE,
  174. stream, 5, &state->buffered_task_handle);
  175. if (task_ret != pdPASS || !state->buffered_task_handle) {
  176. ESP_LOGE(TAG, "Failed to create buffered audio task");
  177. close(state->buffered_listen_socket);
  178. state->buffered_listen_socket = -1;
  179. stream->running = false;
  180. return ESP_FAIL;
  181. }
  182. return ESP_OK;
  183. }
  184. static void buffered_stop(audio_stream_t *stream) {
  185. audio_receiver_state_t *state = audio_stream_state(stream);
  186. if (!stream->running && !state->buffered_task_handle) {
  187. return;
  188. }
  189. stream->running = false;
  190. if (state->buffered_client_socket > 0) {
  191. close(state->buffered_client_socket);
  192. state->buffered_client_socket = -1;
  193. }
  194. if (state->buffered_listen_socket > 0) {
  195. close(state->buffered_listen_socket);
  196. state->buffered_listen_socket = -1;
  197. }
  198. if (!buffered_wait_for_task_stopped(state, 20)) {
  199. ESP_LOGW(TAG, "Buffered audio task did not exit within timeout");
  200. return;
  201. }
  202. if (state->buffered_recv_buffer) {
  203. heap_caps_free(state->buffered_recv_buffer);
  204. state->buffered_recv_buffer = NULL;
  205. }
  206. state->buffered_port = 0;
  207. }
  208. static uint16_t buffered_get_port(audio_stream_t *stream) {
  209. audio_receiver_state_t *state = audio_stream_state(stream);
  210. return state->buffered_port;
  211. }
  212. static bool buffered_is_running(audio_stream_t *stream) {
  213. return stream->running;
  214. }
  215. static void buffered_destroy(audio_stream_t *stream) {
  216. if (!stream) {
  217. return;
  218. }
  219. buffered_stop(stream);
  220. audio_receiver_state_t *state = audio_stream_state(stream);
  221. if (state->buffered_task_handle) {
  222. ESP_LOGW(TAG, "Leaking buffered stream because task shutdown timed out");
  223. return;
  224. }
  225. free(stream);
  226. }
  227. const audio_stream_ops_t audio_stream_buffered_ops = {
  228. .start = buffered_start,
  229. .stop = buffered_stop,
  230. .receive_packet = NULL,
  231. .decrypt_payload = NULL,
  232. .get_port = buffered_get_port,
  233. .is_running = buffered_is_running,
  234. .destroy = buffered_destroy};