Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
113 changes: 88 additions & 25 deletions src/main.c
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -2188,51 +2204,59 @@ 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;
stream->response_pending = false;
}
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);
if (!item)
{
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);
if (xQueueSend(stream->receive_queue, &item, 0) != pdTRUE)
{
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)
{
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)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 (;;)
{
Expand All @@ -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)
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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();
Expand All @@ -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)
Expand Down Expand Up @@ -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)
{
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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++;
Expand All @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -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;
Expand Down
19 changes: 16 additions & 3 deletions src/status_led.c
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include "esp_log.h"

#include "freertos/FreeRTOS.h"
#include "freertos/semphr.h"
#include "freertos/task.h"

#include "status_led.h"
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand All @@ -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;
}
Loading