diff --git a/src/main.c b/src/main.c index dad476e..cba576e 100644 --- a/src/main.c +++ b/src/main.c @@ -143,6 +143,13 @@ typedef struct QueueHandle_t receive_queue; } singmux_tcp_stream_t; +typedef enum +{ + SINGMUX_TCP_DELIVER_OK, + SINGMUX_TCP_DELIVER_BLOCKED, + SINGMUX_TCP_DELIVER_ERROR, +} singmux_tcp_deliver_result_t; + typedef enum { SINGMUX_CONTROL_OPEN, @@ -205,9 +212,14 @@ static uint32_t s_singmux_next_stream_id = 3; static uint8_t *s_singmux_rx_buffer; /* Only the smux manager calls recv(), so this does not need per-task storage. * Keeping it static avoids consuming most of that task's 6 KiB stack. */ -static uint8_t s_singmux_socket_read_buffer[SINGMUX_SOCKET_READ_SIZE]; -static size_t s_singmux_rx_length; -static bool s_singmux_vless_response_pending; +static uint8_t s_singmux_socket_read_buffer[SINGMUX_SOCKET_READ_SIZE]; +static size_t s_singmux_rx_length; +static bool s_singmux_vless_response_pending; +/* A full per-stream queue leaves its frame in the shared reassembly buffer. + * The manager waits + * for the owning relay to dequeue an item before retrying. */ +static bool s_singmux_rx_blocked; +static TaskHandle_t s_singmux_manager_task; static singmux_tcp_stream_t s_singmux_tcp_streams[SINGMUX_TCP_STREAM_MAX]; static singmux_stats_t s_singmux_stats; static uint32_t s_upload_bps; @@ -1923,6 +1935,10 @@ static void relay_singmux_tcp_stream(int client, int slot) singmux_tcp_data_t *item = NULL; while (xQueueReceive(receive_queue, &item, 0) == pdTRUE) { + if (s_singmux_manager_task) + { + xTaskNotifyGive(s_singmux_manager_task); + } if (!item) { detached = singmux_tcp_close(slot, stream_id); @@ -2188,18 +2204,19 @@ static void singmux_tcp_stream_discard(int slot) memset(&s_singmux_tcp_streams[slot], 0, sizeof(s_singmux_tcp_streams[slot])); } -static bool singmux_tcp_deliver(int slot, const uint8_t *payload, size_t length) +static singmux_tcp_deliver_result_t singmux_tcp_deliver(int slot, const uint8_t *payload, + size_t length) { singmux_tcp_stream_t *stream = &s_singmux_tcp_streams[slot]; if (stream->response_pending) { if (!length) { - return true; + return SINGMUX_TCP_DELIVER_OK; } if (payload[0] != 0) { - return false; + return SINGMUX_TCP_DELIVER_ERROR; } ++payload; --length; @@ -2207,7 +2224,15 @@ static bool singmux_tcp_deliver(int slot, const uint8_t *payload, size_t length) } if (!length) { - return true; + return SINGMUX_TCP_DELIVER_OK; + } + /* The manager is the sole producer. Checking before allocating avoids a + * needless PSRAM + * copy when back-pressure is already required. */ + if (!stream->receive_queue || uxQueueSpacesAvailable(stream->receive_queue) == 0) + { + s_singmux_stats.tcp_rx_queue_full++; + return SINGMUX_TCP_DELIVER_BLOCKED; } singmux_tcp_data_t *item = heap_caps_malloc(sizeof(*item) + length, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); @@ -2215,7 +2240,7 @@ static bool singmux_tcp_deliver(int slot, const uint8_t *payload, size_t length) { s_singmux_stats.tcp_rx_alloc_failures++; stream->peer_closed = true; - return true; + return SINGMUX_TCP_DELIVER_OK; } item->length = length; memcpy(item->data, payload, length); @@ -2223,8 +2248,7 @@ static bool singmux_tcp_deliver(int slot, const uint8_t *payload, size_t length) { free(item); s_singmux_stats.tcp_rx_queue_full++; - stream->peer_closed = true; - return true; + return SINGMUX_TCP_DELIVER_BLOCKED; } UBaseType_t depth = uxQueueMessagesWaiting(stream->receive_queue); if (depth > s_singmux_stats.tcp_rx_queue_high_water) @@ -2232,7 +2256,7 @@ static bool singmux_tcp_deliver(int slot, const uint8_t *payload, size_t length) s_singmux_stats.tcp_rx_queue_high_water = (uint16_t)depth; } bandwidth_record_download(length); - return true; + return SINGMUX_TCP_DELIVER_OK; } static bool singmux_tcp_notify_closed(int slot) @@ -2454,6 +2478,7 @@ static void singmux_process_control(void) s_singmux_tunnel = open_singmux_session(); s_singmux_next_stream_id = 3; s_singmux_rx_length = 0; + s_singmux_rx_blocked = false; if (!s_singmux_rx_buffer) { s_singmux_rx_buffer = heap_caps_malloc(SINGMUX_RX_BUFFER_MAX, @@ -2551,6 +2576,7 @@ static void singmux_process_control(void) static void transparent_udp_manager_task(void *arg) { (void)arg; + s_singmux_manager_task = xTaskGetCurrentTaskHandle(); uint8_t smux_receive_burst = 0; for (;;) { @@ -2563,6 +2589,7 @@ static void transparent_udp_manager_task(void *arg) } s_singmux_rx_length = 0; s_singmux_vless_response_pending = false; + s_singmux_rx_blocked = false; for (int i = 0; i < UDP_ASSOCIATION_MAX; ++i) { if (s_udp_associations[i].in_use) @@ -2636,6 +2663,7 @@ static void transparent_udp_manager_task(void *arg) s_singmux_tunnel = open_singmux_session(); s_singmux_next_stream_id = 3; s_singmux_rx_length = 0; + s_singmux_rx_blocked = false; if (!s_singmux_rx_buffer) { s_singmux_rx_buffer = heap_caps_malloc( @@ -2754,10 +2782,12 @@ static void transparent_udp_manager_task(void *arg) } } struct timeval poll_timeout = {.tv_sec = 0, .tv_usec = 0}; - if (max_fd >= 0 && select(max_fd + 1, &reads, NULL, NULL, &poll_timeout) > 0) + if (s_singmux_rx_blocked || + (max_fd >= 0 && select(max_fd + 1, &reads, NULL, NULL, &poll_timeout) > 0)) { if (use_vless && s_config.singmux_enabled && s_singmux_tunnel >= 0 && - FD_ISSET(s_singmux_tunnel, &reads) && !singmux_receive_available()) + (s_singmux_rx_blocked || FD_ISSET(s_singmux_tunnel, &reads)) && + !singmux_receive_available()) { s_singmux_stats.sessions_closed++; s_singmux_stats.last_session_close_ms = now_ms(); @@ -2769,8 +2799,9 @@ static void transparent_udp_manager_task(void *arg) (unsigned)s_singmux_stats.close_protocol_error, s_singmux_stats.last_socket_errno); close(s_singmux_tunnel); - s_singmux_tunnel = -1; - s_singmux_rx_length = 0; + s_singmux_tunnel = -1; + s_singmux_rx_length = 0; + s_singmux_rx_blocked = false; for (int i = 0; i < UDP_ASSOCIATION_MAX; ++i) { if (s_udp_associations[i].in_use) @@ -2813,7 +2844,17 @@ static void transparent_udp_manager_task(void *arg) * task runnable forever and starves same-priority HTTP/relay tasks. * Four 4 KiB passes per tick preserves responsiveness without imposing * a roughly 400 KiB/s ceiling on a sustained download. */ - if (smux_receive_pending) + if (s_singmux_rx_blocked) + { + /* A relay task gives this notification after it frees one queue + * item. + * The timeout still services UDP/control work if the AP + * client has + * stopped accepting data. */ + smux_receive_burst = 0; + ulTaskNotifyTake(pdTRUE, pdMS_TO_TICKS(20)); + } + else if (smux_receive_pending) { if (++smux_receive_burst >= 4) { @@ -3156,8 +3197,9 @@ static esp_err_t config_post(httpd_req_t *req) if (section == CONFIG_SECTION_VLESS && s_singmux_tunnel >= 0) { close(s_singmux_tunnel); - s_singmux_tunnel = -1; - s_singmux_rx_length = 0; + s_singmux_tunnel = -1; + s_singmux_rx_length = 0; + s_singmux_rx_blocked = false; for (int i = 0; i < UDP_ASSOCIATION_MAX; ++i) { if (s_udp_associations[i].in_use) @@ -4091,8 +4133,9 @@ static void serial_console_task(void *arg) if (changed_vless && s_singmux_tunnel >= 0) { close(s_singmux_tunnel); - s_singmux_tunnel = -1; - s_singmux_rx_length = 0; + s_singmux_tunnel = -1; + s_singmux_rx_length = 0; + s_singmux_rx_blocked = false; singmux_tcp_notify_all_closed(); } if (apply_ap) @@ -4171,8 +4214,14 @@ static bool singmux_deliver_udp(udp_association_t *association, const uint8_t *p static bool singmux_receive_available(void) { - int bytes = recv(s_singmux_tunnel, s_singmux_socket_read_buffer, + bool retry_blocked_frame = s_singmux_rx_blocked; + s_singmux_rx_blocked = false; + int bytes = 0; + if (!retry_blocked_frame) + { + bytes = recv(s_singmux_tunnel, s_singmux_socket_read_buffer, sizeof(s_singmux_socket_read_buffer), MSG_DONTWAIT); + } if (bytes > 0) { s_singmux_stats.rx_reads++; @@ -4190,12 +4239,12 @@ static bool singmux_receive_available(void) } s_singmux_stats.rx_wire_bytes += (size_t)bytes; } - else if (bytes == 0) + else if (!retry_blocked_frame && bytes == 0) { s_singmux_stats.close_eof++; return false; } - else if (errno != EAGAIN && errno != EWOULDBLOCK) + else if (!retry_blocked_frame && errno != EAGAIN && errno != EWOULDBLOCK) { s_singmux_stats.close_socket_error++; s_singmux_stats.last_socket_errno = errno; @@ -4252,12 +4301,22 @@ static bool singmux_receive_available(void) int tcp_slot = singmux_tcp_stream_find(stream_id); if (tcp_slot >= 0) { - if (command == SMUX_CMD_PSH && - !singmux_tcp_deliver(tcp_slot, frame + SMUX_HEADER_SIZE, payload_length)) + singmux_tcp_deliver_result_t delivered = SINGMUX_TCP_DELIVER_OK; + if (command == SMUX_CMD_PSH) + { + delivered = + singmux_tcp_deliver(tcp_slot, frame + SMUX_HEADER_SIZE, payload_length); + } + if (delivered == SINGMUX_TCP_DELIVER_ERROR) { s_singmux_stats.close_protocol_error++; return false; } + if (delivered == SINGMUX_TCP_DELIVER_BLOCKED) + { + s_singmux_rx_blocked = true; + break; + } if (command == SMUX_CMD_FIN) { singmux_tcp_notify_closed(tcp_slot); @@ -4278,6 +4337,10 @@ static bool singmux_receive_available(void) static bool singmux_socket_readable(void) { + if (s_singmux_rx_blocked) + { + return true; + } if (s_singmux_tunnel < 0) { return false; diff --git a/src/status_led.c b/src/status_led.c index dbe7a69..c008493 100644 --- a/src/status_led.c +++ b/src/status_led.c @@ -7,6 +7,7 @@ #include "esp_log.h" #include "freertos/FreeRTOS.h" +#include "freertos/semphr.h" #include "freertos/task.h" #include "status_led.h" @@ -26,9 +27,10 @@ typedef struct rmt_symbol_word_t reset_code; } ws2812_encoder_t; -static const char *s_tag; +static const char *s_tag; static rmt_channel_handle_t s_rgb_channel; static rmt_encoder_handle_t s_rgb_encoder; +static SemaphoreHandle_t s_rgb_lock; static bool s_rgb_ready; RMT_ENCODER_FUNC_ATTR @@ -138,11 +140,16 @@ void status_led_set_rgb(uint8_t red, uint8_t green, uint8_t blue) #else const uint8_t rgb[] = {green, red, blue}; #endif + if (!s_rgb_lock || xSemaphoreTake(s_rgb_lock, pdMS_TO_TICKS(25)) != pdTRUE) + { + return; + } rmt_transmit_config_t transmit = {0}; if (rmt_transmit(s_rgb_channel, s_rgb_encoder, rgb, sizeof(rgb), &transmit) == ESP_OK) { - rmt_tx_wait_all_done(s_rgb_channel, pdMS_TO_TICKS(20)); + (void)rmt_tx_wait_all_done(s_rgb_channel, pdMS_TO_TICKS(20)); } + xSemaphoreGive(s_rgb_lock); } void status_led_show_mode(transparent_mode_t mode, bool has_upstream) @@ -182,7 +189,7 @@ void status_led_boot_indicator(transparent_mode_t mode, bool has_upstream) void status_led_init(const char *tag) { - s_tag = tag; + s_tag = tag; rmt_tx_channel_config_t channel = { .gpio_num = RGB_LED_GPIO, .clk_src = RMT_CLK_SRC_DEFAULT, @@ -204,5 +211,11 @@ void status_led_init(const char *tag) ESP_LOGW(s_tag, "RGB LED unavailable on GPIO%d: %s", RGB_LED_GPIO, esp_err_to_name(err)); return; } + s_rgb_lock = xSemaphoreCreateMutex(); + if (!s_rgb_lock) + { + ESP_LOGW(s_tag, "RGB LED lock unavailable"); + return; + } s_rgb_ready = true; }