From f47ad7d87f4c4bb975fd97799b9b22cd180ed0d0 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 3 Sep 2026 13:49:09 +0300 Subject: [PATCH] tcp_proxy: fix throughput (tx_queue backpressure) + connect timeout + integration test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Router: unified service delivery format [svc_id][src][dst][payload] (end-to-end src/dst from SVC_ROUTE header) - tcp_proxy_client: replace single tx_buf drop-on-backpressure with bounded tx_queue + send_q waiter + retry timer (no data loss, bulk window reopen — removes 62.5ms delayed-ACK stall) - tcp_io: configurable connect timeout (tcp_conn_set_connect_timeout) - tcp_proxy_server: connect_timeout config option (ms, default 2000, 0=off) - check.sh: run tcp_proxy_full integration test with sudo (skip if passwordless sudo unavailable) - test_tcp_io: connect timeout unit test; docs (lwip_nuances.md) --- AGENTS.md | 3 + check.sh | 11 ++ doc/TASK_tcp_proxy_throughput.md | 115 ++++++++++++ doc/lwip_nuances.md | 66 +++++++ lib/tcp_io.c | 27 +++ lib/tcp_io.h | 9 + src/config_parser.c | 9 +- src/config_parser.h | 1 + src/media_delivery/media_delivery.c | 12 +- src/nat_transport.c | 20 +- src/proxy/icmp_proxy.c | 59 +++--- src/proxy/icmp_proxy.h | 10 +- src/proxy/tcp_proxy_client.c | 244 +++++++++++++++++-------- src/proxy/tcp_proxy_client.h | 10 +- src/proxy/tcp_proxy_server.c | 55 +++--- src/proxy/tcp_proxy_server.h | 4 + src/proxy/tcp_proxy_server_doc.md | 1 + src/proxy/udp_proxy.c | 56 +++--- src/proxy/udp_proxy.h | 10 +- src/routing_layer/conn_mgr_core.c | 4 +- src/routing_layer/etcp_router.c | 70 +++++-- src/routing_layer/etcp_router.h | 8 + src/routing_layer/routing.c | 15 +- tests/Makefile.am | 5 + tests/tcp_proxy_full/client.conf | 10 +- tests/tcp_proxy_full/exit.conf | 4 +- tests/tcp_proxy_full/intermediate.conf | 28 +++ tests/tcp_proxy_full/run_test.sh | 32 +++- tests/test_etcp_router.c | 18 +- tests/test_etcp_router_reconnect.c | 12 +- tests/test_etcp_router_unit.c | 8 +- tests/test_icmp_proxy.c | 12 +- tests/test_nat_transport.c | 10 +- tests/test_tcp_io.c | 47 +++++ tests/test_udp_proxy.c | 9 +- tools/chatgui/transport/utun_node.cpp | 7 +- 36 files changed, 755 insertions(+), 266 deletions(-) create mode 100644 doc/TASK_tcp_proxy_throughput.md create mode 100644 doc/lwip_nuances.md create mode 100644 tests/tcp_proxy_full/intermediate.conf diff --git a/AGENTS.md b/AGENTS.md index 94088946..c0657df2 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -25,6 +25,9 @@ This file contains essential information for AI coding agents working in the uTu При поиске багов всегда проверяй достаточно ли отладочной информации. Если не достаточно - фокусируйся на том чтобы добавить нужную информацию в вывод. Помни что отладка управляется фильтрами по категормяи, которые нужно проверить и при необходимости правильно настроить (в конфиге или в коде). +Для поиска выстараивай информативные логи. По которым максимально чётко понятно что происходит. Если какие-то места которые могут иметь отношение к ошибке в логах пропущены - это повод доработать логи и прогнать еще раз. +Всегда отладку выстраивай через доработку логов, пока точно и однозначно не будет видно в каком точно месте проблема. + ## Quick Reference **Repository:** uTun - Secure VPN tunnel with ETCP protocol diff --git a/check.sh b/check.sh index ce0e53ab..19127923 100755 --- a/check.sh +++ b/check.sh @@ -16,6 +16,17 @@ timeout $TMO make check -j4 > "$TMP" 2>&1 || { grep -E '\[PASS|\[FAIL|\[SKIP|Test Results|Failed tests|Log directory|^====' "$TMP" rm -f "$TMP" +# Интеграционный тест tcp_proxy_full (требует sudo; без passwordless sudo — пропускаем) +echo "" +if [ "$(id -u)" -eq 0 ]; then + make -C tests check-proxy +elif command -v sudo >/dev/null 2>&1 && sudo -n true 2>/dev/null; then + echo "=== tcp_proxy_full integration test (sudo) ===" + sudo -n make -C tests check-proxy +else + echo "[SKIP] tcp_proxy_full integration test (requires passwordless sudo)" +fi + ELAPSED=$(( $(date +%s) - START )) echo "Time: ${ELAPSED}s" echo "$ELAPSED" > .check_time diff --git a/doc/TASK_tcp_proxy_throughput.md b/doc/TASK_tcp_proxy_throughput.md new file mode 100644 index 00000000..e9bd377d --- /dev/null +++ b/doc/TASK_tcp_proxy_throughput.md @@ -0,0 +1,115 @@ +# Задача: продолжить исправление TCP-proxy (крупные файлы через транзит) + +Продолжение незавершённой работы. Читай целиком, прежде чем что-то менять. + +## 1. Что уже сделано (НЕ переделывать) + +Унифицирован формат доставки сервисных кодограмм роутером. Раньше сервисы определяли +источник через `conn->peer_node_id`, что ломало транзит (при промежуточном узле это был +id промежуточного, а не реального клиента). + +Теперь роутер доставляет сервисам единый формат: + +``` +[svc_id(1)][src_node_id(8)][dst_node_id(8)][payload...] +``` + +где src/dst — реальные end-to-end узлы из SVC_ROUTE заголовка. + +Ключевые константы в `src/routing_layer/etcp_router.h`: +`ROUTER_SVC_SRC_OFF=1`, `ROUTER_SVC_DST_OFF=9`, `ROUTER_SVC_PAYLOAD_OFF=17`, `ROUTER_SVC_HDR_SIZE=17`. + +Изменённые файлы (все в некоммиченном состоянии, `git status`): +- `src/routing_layer/etcp_router.{h,c}` — `router_deliver` + `router_deliver_loopback` + CLOSE/RESTART +- `src/proxy/tcp_proxy_server.{h,c}` — recv +8→+16 сдвиг, `src_node_id` из `dgram[1..8]` +- `src/proxy/tcp_proxy_client.{h,c}` — recv сдвиг + счётчики `bytes_to_exit/from_exit/bp_count` +- `src/proxy/udp_proxy.{h,c}`, `src/proxy/icmp_proxy.{h,c}` — сдвиг + убран встроенный node_id +- `src/nat_transport.c` — убран встроенный node_id +- `src/routing_layer/routing.c`, `conn_mgr_core.c` — сдвиг +- `src/media_delivery/media_delivery.c` — сдвиг + `from_node` из `dgram[1..8]` +- `tools/chatgui/transport/utun_node.cpp` — сдвиг +- тесты `tests/test_etcp_router*.c`, `test_udp_proxy.c`, `test_icmp_proxy.c`, `test_nat_transport.c` +- `tests/tcp_proxy_full/` — добавлен `intermediate.conf` (3-узловая топология), обновлены `client.conf`/`exit.conf`/`run_test.sh` + +**Статус проверки: `./check.sh` = 73 passed, 0 failed, 1 skipped.** Форматная часть ВЕРНА. + +## 2. Что осталось (собственно задача) + +Интеграционный тест прокси `tests/tcp_proxy_full/run_test.sh` (нужен sudo) — крупные +передачи НЕ проходят. Проверено: проблема **НЕ в формате и НЕ в транзите** — прямой +2-узловой прогон падает идентично. Это предсуществующая проблема прокси/lwIP. + +Результат `sudo ./tests/tcp_proxy_full/run_test.sh`: +- `idle` (32 КБ) — PASS +- `half_close` (64 КБ) — 64280/65536 (почти, теряется хвост) +- `basic_1mb` (1 МБ), `concurrent_*`, `stress` — FAIL (timeout) + +## 3. Диагноз (уже установлен по логам) + +Пропускная способность ~24 КБ/с (1 пакет 1460 байт за ~60 мс). Для 1 МБ это ~43 с → таймаут 10 с. + +Цепочка (видна в добавленных логах `PROXY BP *` и `SEND_Q_ACKED`): +1. Клиент быстро шлёт 320 пакетов (initial burst ~33 МБ/с). +2. Роутерный `send_q` переполняется (`router_send_q: FULL count=64`) → `etcp_route_send` возвращает -1. +3. Прокси в `tcp_proxy_client_recv_cb` ставит одиночный `tx_buf` и **сбрасывает** дальнейшие данные + (`if (pc->tx_buf) { pbuf_free(p); return; }`) — лог `PROXY BP drop`, без `tcp_recved` → окно lwIP схлопывается. +4. Клиент вырождается в режим «1 пакет за ~60 мс». +5. ~60 мс = таймер lwIP: `ctx->tmr_interval_ms = TCP_TMR_INTERVAL / 4` = 62.5 мс (`src/lwip_tcp/lwip_tcp.c:106`). + В режиме «пакет-за-пакетом» lwIP не шлёт ACK сразу (порог `TCP_WND_UPDATE_THRESHOLD = TCP_WND/4 = 2920` + не набирается за 1 пакет), и отложенный ACK уходит только по `tcp_fasttmr` раз в 62.5 мс. + +Триггер — сброс данных при backpressure в прокси (`tcp_proxy_client.c`), НЕ lwIP. + +## 4. Ограничения (ВАЖНО) + +- **В `src/lwip_tcp/` ничего НЕ править без явного согласования пользователя. Разрешён только debug-лог.** +- Правки исходников только через Edit tool (никакого sed для массовых замен). +- Ошибки/нештатные ветки — обязательно `DEBUG_ERROR/DEBUG_WARN`. +- Логи информативные, без спама. Debug-категории: `proxy`, `etcp_route`. +- Перед сборкой `make clean`. + +## 5. План продолжения + +Шаг 1 (диагностика, закрыть пробел — по желанию, уже почти доказано): +- Добавить в lwIP только DEBUG-лог (разрешено): в `tcp_fasttmr` — счётчик отложенных ACK; в `tcp_recved` — `rcv_wnd`. +- В прокси (`tcp_proxy_client.c`, НЕ lwIP) в момент `BP drop` лог `pcb->rcv_wnd`. +- Подтвердить, что 62.5 мс — это отложенный ACK. + +Шаг 2 (фикс, требуется согласование — предложи пользователю): +- **Вариант A (предпочтительно, НЕ lwIP):** в прокси заменить одиночный `tx_buf` + сброс на очередь с backpressure + (буферизовать все пакеты во время backpressure, не сбрасывать; `tcp_recved` вызывать по мере отправки). + Тогда окно закрывается плавно, OS TCP не сваливается в congestion avoidance, деградации нет. + Смотри `struct tcp_proxy_client_conn` в `tcp_proxy_client.h` (сейчас: `tx_buf`/`tx_len`/`tx_waiter`). + Для образца очереди с backpressure смотри `src/proxy/tcp_proxy_server.c` (там `read_queue` + `pause_waiter`). +- **Вариант B (lwIP, только по согласованию):** уменьшить `tmr_interval_ms` (`TCP_TMR_INTERVAL/4` → чаще), + либо повысить `TCP_WND_UPDATE_THRESHOLD`/`TCP_WND`, чтобы ACK слался сразу. + +Шаг 3 (почистить тестовую обвязку): +- В `tests/tcp_proxy_full/*.conf` сейчас включены отладочные уровни (`etcp_route=trace`, `proxy=trace`) — + перед коммитом вернуть к разумным (`etcp_route=debug`/`proxy=debug` или `error`). +- Решить: оставить добавленные в `tcp_proxy_client.c` счётчики/BP-логи или вычистить. +- В `run_test.sh` тайминги `elapsed` через `date +%s%3N` дают мусор в этой среде — поправить (например `date +%s` + без миллисекунд, или `SECONDS`). Stress-таймаут `60s` можно уменьшить до ~15-20s. + +## 6. Сборка и тест + +```bash +./build.sh --full -j4 # или make -j4 (после make clean) +./check.sh # unit-тесты: должны остаться 73 passed +cp src/utun utun # обновить корневой бинарник для run_test.sh! +sudo ./tests/tcp_proxy_full/run_test.sh # интеграционный (нужен root/TUN/iptables) +``` + +Важно: `run_test.sh` использует `UTUN_BIN="$SCRIPT_DIR/../../utun"` (корневой `utun`), а `make` собирает +`src/utun`. После каждой пересборки делать `cp src/utun utun`. + +## 7. Как быстро убедиться, что фикс помог + +Прогнать `basic_1mb` и смотреть лог клиента `tests/tcp_proxy_full/log/client_utun.log`: +- ДО фикса: `PROXY BP drop` + `SEND_Q_ACKED` с интервалом ~60 мс. +- ПОСЛЕ фикса: нет `BP drop`, интервал ACK ~10 мс, `basic_1mb` PASS. + +Полезные grep-и по логу клиента: +``` +grep -E "PROXY BP|SEND_Q_ACKED|PROXY FIN" client_utun.log +``` diff --git a/doc/lwip_nuances.md b/doc/lwip_nuances.md new file mode 100644 index 00000000..ef36d176 --- /dev/null +++ b/doc/lwip_nuances.md @@ -0,0 +1,66 @@ +# Важные нюансы работы встроенного lwIP TCP-стека + +Краткая памятка по тонким местам `src/lwip_tcp/`, критичным для прокси. +Полное описание модуля — в `src/lwip_tcp/lwip_tcp_doc.md`. + +## Ключевые константы (`lwip_tcp_opts.h`) + +| Константа | Значение | Смысл | +|-----------|----------|-------| +| `TCP_MSS` | 1460 | максимальный размер сегмента | +| `TCP_WND` | 8×MSS = 11680 | окно приёма lwIP (объявляется пиру) | +| `TCP_SND_BUF` | 16×MSS = 23360 | буфер отправки | +| `TCP_TMR_INTERVAL` | 250 мс | базовый интервал таймера | +| `tmr_interval_ms` | `TCP_TMR_INTERVAL/4` = 62.5 мс | фактический шаг `tcp_fasttmr` | +| `TCP_WND_UPDATE_THRESHOLD` | `TCP_WND/4` = 2920 | порог немедленного ACK окна | + +**Важно:** таймер uasync использует timebase 0.1 мс — `uasync_set_timeout(ua, N, …)` +задаёт N×0.1 мс. Например, 5000 = 500 мс. + +## Контракт recv_cb — главный источник багов + +lwIP вызывает `recv_cb(arg, pcb, pbuf, err)` для данных, **уже извлечённых из pcb**. +Дальше приложение само решает судьбу данных и окна: + +| Действие | Результат | +|----------|-----------| +| скопировать + `tcp_recved(pcb, len)` | данные приняты, окно восстановлено (норма) | +| `return != LERR_OK` | lwIP кладёт pbuf в `refused_data` и отдаст позже — **без потери**, окно закрывается | +| `pbuf_free(p)` без `tcp_recved` | данные потеряны, окно не восстановлено → схлопывание окна | + +`refused_data` — штатный backpressure lwIP: одноканальный (один pbuf), ре-доставка +через `tcp_fasttmr` (62.5 мс) или при следующем входящем сегменте. + +**Ошибка, из-за которой падал `tcp_proxy_full`:** при backpressure прокси делал +`pbuf_free(p)` без `tcp_recved`. Каждый такой сброс уменьшал `rcv_wnd` на MSS; после +~8 пакетов окно падало в 0, OS-TCP переставал слать, и передача умирала. + +## Управление окном и «деградация 62.5 мс» + +`tcp_recved(pcb, len)` увеличивает `rcv_wnd` и через `tcp_update_rcv_ann_wnd` решает, +слать ли ACK окна немедленно. Немедленный ACK идёт только если прирост окна +`wnd_inflation ≥ TCP_WND_UPDATE_THRESHOLD` (=2920 = 2×MSS); иначе ACK откладывается +до `tcp_fasttmr` (62.5 мс). + +Отсюда классическая ловушка: если после backpressure окно восстанавливать **по одному +пакету** (`tcp_recved(1460)` за раз), прирост 1460 < 2920 → ACK уходит раз в 62.5 мс → +пир шлёт 1 пакет за 62.5 мс ≈ 23 КБ/с. Это выглядит как «пропускная способность +застряла», хотя сеть свободна. + +**Правильно:** восстанавливать окно **пачкой** (несколько `tcp_recved` подряд либо один +`tcp_recved(TCP_WND_MAX - rcv_wnd)`), тогда прирост ≥ 2920 → ACK немедленный → пир +сразу возобновляет бурст. + +## Рекомендуемый паттерн для relay-потребителя (как в `tcp_proxy_client.c`) + +1. При получении данных — копировать в ограниченную очередь, **не вызывать** + `tcp_recved` (окно само плавно закрывается; очередь ограничена `TCP_WND`). +2. Дрейн очереди по сигналу освобождения нижележащего канала + (`etcp_router_on_send_ready`) + retry-таймер (force=1) как safety-net. +3. `tcp_recved(len)` вызывать **только после успешной отправки** — окно + восстанавливается пачкой по мере дрена, ACK уходит немедленно. +4. Никогда не сбрасывать pbuf без `tcp_recved`; не держать данные в `refused_data` + дольше одного пакета (single-slot). + +См. реализацию: `tcp_proxy_client_tx_queue_drain_cb` / `tcp_proxy_client_fin_flush` +в `src/proxy/tcp_proxy_client.c`. diff --git a/lib/tcp_io.c b/lib/tcp_io.c index bea103c9..884cabb6 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -25,6 +25,7 @@ static void fin_deferred_cb(struct ll_queue* q, void* arg); static void write_queue_fetch_cb(struct ll_queue* q, void* arg); static int flush_write_buf(struct tcp_conn* tc); static void tcp_conn_handle_error(struct tcp_conn* tc, int err); +static void connect_timeout_cb(void* arg); static const uint8_t tcp_fin_sentinel; struct tcp_conn* tcp_conn_create( @@ -121,6 +122,8 @@ void tcp_conn_destroy(struct tcp_conn* tc) { DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_conn_destroy: fd=%d connected=%d error=%d fin_remote=%d fin_local=%d closed=%d", (int)tc->sock, tc->connected, tc->error, tc->fin_remote, tc->fin_local, tc->closed); + if (tc->connect_timer) { uasync_cancel_timeout(tc->ua, tc->connect_timer); tc->connect_timer = NULL; } + if (tc->socket_id) { uasync_remove_socket_t(tc->ua, tc->sock); tc->socket_id = NULL; @@ -383,6 +386,7 @@ static void write_cb(socket_t sock, void* arg) { socklen_t len = sizeof(err); if (getsockopt(tc->sock, SOL_SOCKET, SO_ERROR, (char*)&err, &len) == 0 && err == 0) { tc->connected = 1; + if (tc->connect_timer) { uasync_cancel_timeout(tc->ua, tc->connect_timer); tc->connect_timer = NULL; } DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock); } else { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err); @@ -453,3 +457,26 @@ void tcp_conn_pause_read(struct tcp_conn* tc) { queue_waiter_cancel(tc->read_queue, &tc->read_waiter); tc->read_paused = 1; } + +// ==================================================================== +// Таймаут установления соединения +// ==================================================================== + +static void connect_timeout_cb(void* arg) { + struct tcp_conn* tc = (struct tcp_conn*)arg; + if (!tc) return; + tc->connect_timer = NULL; + if (!tc->connected && tc->sock != SOCKET_INVALID) { + DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: connect timeout fd=%d (limit %d ms)", (int)tc->sock, tc->connect_timeout_ms); + tcp_conn_handle_error(tc, ETIMEDOUT); + } +} + +void tcp_conn_set_connect_timeout(struct tcp_conn* tc, int timeout_ms) { + if (!tc) return; + if (timeout_ms < 0) timeout_ms = 0; + tc->connect_timeout_ms = timeout_ms; + if (!tc->connected && timeout_ms > 0) { + tc->connect_timer = uasync_set_timeout(tc->ua, timeout_ms * 10, tc, connect_timeout_cb, "tcp_conn_ct"); + } +} diff --git a/lib/tcp_io.h b/lib/tcp_io.h index a8f609e4..39480036 100644 --- a/lib/tcp_io.h +++ b/lib/tcp_io.h @@ -53,6 +53,9 @@ struct tcp_conn { uint8_t closed; // сокет полностью закрыт (close) uint8_t destroyed; // 1 = tcp_conn_destroy вызван, tc ожидает отложенного free + int connect_timeout_ms; // таймаут установления соединения, 0 = без таймаута + void* connect_timer; // handle uasync таймера connect + // Частичная отправка (из data_pool, не в очереди — досылается первой) uint8_t* write_buf; size_t write_len; @@ -113,6 +116,12 @@ void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_con // Принудительная остановка чтения: убирает EPOLLIN и отменяет read_waiter. void tcp_conn_pause_read(struct tcp_conn* tc); +// Таймаут установления соединения (мс, 0 = без таймаута). Если соединение ещё не +// установлено (connected=0), взводит таймер; по истечении вызывает +// tcp_conn_handle_error(ETIMEDOUT) → on_error. Таймер гасится при успешном connect +// и при tcp_conn_destroy. +void tcp_conn_set_connect_timeout(struct tcp_conn* tc, int timeout_ms); + #ifdef __cplusplus } diff --git a/src/config_parser.c b/src/config_parser.c index 7e69d4c3..71b025b4 100644 --- a/src/config_parser.c +++ b/src/config_parser.c @@ -870,6 +870,7 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename) cfg->global.auto_sockets = 0; // Default: disabled, use manual [server] sockets cfg->global.auto_socket_skip_no_default_route = 1; // Default: в авто-режиме пропускать интерфейсы без default-маршрута cfg->global.tcp_recv_buf = 0; // 0 = OS default + cfg->global.tcp_proxy_server_connect_timeout_ms = 2000; // 2 с, 0 = без таймаута cfg->global.firewall_rules = NULL; cfg->global.firewall_rule_count = 0; cfg->global.firewall_bypass_all = 0; @@ -1053,8 +1054,10 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename) } else if (strcmp(key, "bind_ip") == 0) { strncpy(cfg->global.tcp_proxy_server_bind_ip, value, sizeof(cfg->global.tcp_proxy_server_bind_ip) - 1); cfg->global.tcp_proxy_server_bind_ip[sizeof(cfg->global.tcp_proxy_server_bind_ip) - 1] = '\0'; + } else if (strcmp(key, "connect_timeout") == 0) { + cfg->global.tcp_proxy_server_connect_timeout_ms = atoi(value); } else { - DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Unknown tcp_proxy_server option '%s'. Valid: tcp_recv_buf, fib, so_mark, netif, bind_ip", filename, line_num, key); + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Unknown tcp_proxy_server option '%s'. Valid: tcp_recv_buf, fib, so_mark, netif, bind_ip, connect_timeout", filename, line_num, key); } break; case SECTION_NETWORK: @@ -1234,8 +1237,10 @@ struct utun_config* parse_config_from_buf(const char *buf, size_t len, const cha } else if (strcmp(key, "bind_ip") == 0) { strncpy(cfg->global.tcp_proxy_server_bind_ip, value, sizeof(cfg->global.tcp_proxy_server_bind_ip) - 1); cfg->global.tcp_proxy_server_bind_ip[sizeof(cfg->global.tcp_proxy_server_bind_ip) - 1] = '\0'; + } else if (strcmp(key, "connect_timeout") == 0) { + cfg->global.tcp_proxy_server_connect_timeout_ms = atoi(value); } else { - DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Unknown tcp_proxy_server option '%s'. Valid: tcp_recv_buf, fib, so_mark, netif, bind_ip", filename, line_num, key); + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "%s:%d: Unknown tcp_proxy_server option '%s'. Valid: tcp_recv_buf, fib, so_mark, netif, bind_ip, connect_timeout", filename, line_num, key); } break; case SECTION_NETWORK: if (cur_network) parse_network(key, value, cur_network, filename, line_num); break; diff --git a/src/config_parser.h b/src/config_parser.h index 06b9c4d5..a2f33cda 100644 --- a/src/config_parser.h +++ b/src/config_parser.h @@ -196,6 +196,7 @@ struct global_config { int tcp_proxy_server_so_mark; // Linux SO_MARK, 0=не менять uint32_t tcp_proxy_server_netif_index; // Linux SO_BINDTODEVICE, 0=нет char tcp_proxy_server_bind_ip[64]; // IP для bind() перед connect() + int tcp_proxy_server_connect_timeout_ms; // таймаут connect(), мс (0 = без таймаута) // NTP configuration ([ntp] section) int ntp_enabled; diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 8a4bdd34..2199c30a 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -1195,12 +1195,12 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { const uint8_t* data = entry->dgram; size_t len = entry->len; - /* router_deliver prepends svc_id byte → data[0]=svc_id, data[1]=наш subcmd. - conn_mgr uses two-tier (cmd+subcmd) hence +2; we have single-tier so +1. */ - if (len < 2) return; - uint8_t subcmd = data[1]; - data += 1; len -= 1; - uint64_t from_node = conn->peer_node_id; + /* Единый формат router_deliver: [svc_id][src_node_id(8)][dst_node_id(8)][subcmd][...] */ + if (len < ROUTER_SVC_PAYLOAD_OFF + 1) return; + uint64_t from_node; + memcpy(&from_node, data + ROUTER_SVC_SRC_OFF, 8); + uint8_t subcmd = data[ROUTER_SVC_PAYLOAD_OFF]; + data += ROUTER_SVC_PAYLOAD_OFF; len -= ROUTER_SVC_PAYLOAD_OFF; DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: recv %s(%02x) from 0x%016llx len=%zu", MD_ID, md_subcmd_name(subcmd), subcmd, (unsigned long long)from_node, len); diff --git a/src/nat_transport.c b/src/nat_transport.c index 43da63b0..727db025 100644 --- a/src/nat_transport.c +++ b/src/nat_transport.c @@ -12,7 +12,7 @@ #include "../lib/ll_queue.h" #include -#define NAT_SVC_HDR_SIZE 9 // svc_id(1) + src_node_id(8) +#define NAT_SVC_HDR_SIZE 1 // send: svc_id(1) + ip_data (src/dst добавляет роутер) static ip_str_t ip_host_to_str(uint32_t ip_host) { struct in_addr a; a.s_addr = htonl(ip_host); @@ -37,8 +37,7 @@ static void nat_transport_client_tun_out_cb(struct ll_queue* q, void* arg) { if (!new_dgram) { queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; } new_dgram[0] = ETCP_RT_ID_NAT; - memcpy(new_dgram + 1, &tr->self_node_id, 8); - memcpy(new_dgram + 9, pkt->dgram + 1, ip_len); + memcpy(new_dgram + 1, pkt->dgram + 1, ip_len); struct ll_entry* new_entry = queue_entry_new(0); if (!new_entry) { u_free(new_dgram); queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; } @@ -80,8 +79,7 @@ static void nat_transport_provider_tun_out_cb(struct ll_queue* q, void* arg) { if (!new_dgram) { queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; } new_dgram[0] = ETCP_RT_ID_NAT; - memcpy(new_dgram + 1, &inst->nat_tr.self_node_id, 8); - memcpy(new_dgram + 9, pkt->dgram + 1, ip_len); + memcpy(new_dgram + 1, pkt->dgram + 1, ip_len); struct ll_entry* new_entry = queue_entry_new(0); if (!new_entry) { u_free(new_dgram); queue_dgram_free(pkt); queue_entry_free(pkt); queue_resume_callback(q); return; } @@ -101,7 +99,7 @@ static void nat_transport_provider_tun_out_cb(struct ll_queue* q, void* arg) { // ETCP_RT_ID_NAT receive via etcp_router: CLIENT gets response, PROVIDER gets request static void nat_transport_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!conn || !entry || !entry->dgram || entry->len < NAT_SVC_HDR_SIZE) { + if (!conn || !entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } @@ -112,13 +110,13 @@ static void nat_transport_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* if (!ctx->initialized) { queue_dgram_free(entry); queue_entry_free(entry); return; } uint64_t src_node_id; - memcpy(&src_node_id, entry->dgram + 1, 8); - uint8_t* ip_data = entry->dgram + NAT_SVC_HDR_SIZE; - size_t ip_len = entry->len - NAT_SVC_HDR_SIZE; + memcpy(&src_node_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + uint8_t* ip_data = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; + size_t ip_len = entry->len - ROUTER_SVC_PAYLOAD_OFF; if (tr->nat_via_node_id != 0) { if (tr->nat_tun) { - tun_write(tr->nat_tun, entry->dgram + NAT_SVC_HDR_SIZE - 1, ip_len + 1); + tun_write(tr->nat_tun, entry->dgram + ROUTER_SVC_PAYLOAD_OFF - 1, ip_len + 1); DEBUG_DEBUG(DEBUG_CATEGORY_NAT, "NAT client: received %zu bytes from provider", ip_len); } } else { @@ -126,7 +124,7 @@ static void nat_transport_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* if (ret < 0) { DEBUG_WARN(DEBUG_CATEGORY_NAT, "Egress NAT failed"); } else if (ret == 0 && tr->nat_tun) { - tun_write(tr->nat_tun, entry->dgram + NAT_SVC_HDR_SIZE - 1, ip_len + 1); + tun_write(tr->nat_tun, entry->dgram + ROUTER_SVC_PAYLOAD_OFF - 1, ip_len + 1); DEBUG_DEBUG(DEBUG_CATEGORY_NAT, "NAT provider: sent %zu bytes to internet", ip_len); } } diff --git a/src/proxy/icmp_proxy.c b/src/proxy/icmp_proxy.c index c4fc716d..bd9eda0f 100644 --- a/src/proxy/icmp_proxy.c +++ b/src/proxy/icmp_proxy.c @@ -137,11 +137,10 @@ static void raw_read_cb(socket_t sock, void* arg) { if (!e->dgram) { queue_entry_free(e); return; } e->dgram[0] = ETCP_RT_ID_ICMP_PROXY; e->dgram[1] = ICMP_PROXY_SUBCMD_REPLY; - memcpy(e->dgram + 2, &r->client_node_id, 8); - memcpy(e->dgram + 10, &r->dst_ip, 4); - memcpy(e->dgram + 14, &r->orig_src_ip, 4); - memcpy(e->dgram + 18, &icmp_hdr->icmp_id, 2); - memcpy(e->dgram + 20, &icmp_hdr->icmp_seq, 2); + memcpy(e->dgram + 2, &r->dst_ip, 4); + memcpy(e->dgram + 6, &r->orig_src_ip, 4); + memcpy(e->dgram + 10, &icmp_hdr->icmp_id, 2); + memcpy(e->dgram + 12, &icmp_hdr->icmp_seq, 2); if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len); e->len = ICMP_PROXY_HDR_SIZE + payload_len; int ret = etcp_route_send(g_icmp_ctx->inst, TOPO_GROUP_UTUN, r->client_node_id, e, 0); @@ -155,15 +154,15 @@ static void raw_read_cb(socket_t sock, void* arg) { // ==================================================================== static void exit_handle_request(struct ETCP_CONN* conn, struct ll_entry* entry) { struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_icmp_ctx ? g_icmp_ctx->inst : NULL); - if (!inst || !g_icmp_ctx || entry->len < ICMP_PROXY_HDR_SIZE + 1) goto drop; + if (!inst || !g_icmp_ctx || entry->len < ICMP_PROXY_RECV_HDR_SIZE + 1) goto drop; - uint64_t client_node_id; memcpy(&client_node_id, entry->dgram + 2, 8); - uint32_t dst_ip; memcpy(&dst_ip, entry->dgram + 10, 4); - uint32_t orig_src_ip; memcpy(&orig_src_ip, entry->dgram + 14, 4); - uint16_t echo_id; memcpy(&echo_id, entry->dgram + 18, 2); - uint16_t echo_seq; memcpy(&echo_seq, entry->dgram + 20, 2); - uint8_t* payload = entry->dgram + ICMP_PROXY_HDR_SIZE; - size_t payload_len = entry->len - ICMP_PROXY_HDR_SIZE; + uint64_t client_node_id; memcpy(&client_node_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + uint32_t dst_ip; memcpy(&dst_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); + uint32_t orig_src_ip; memcpy(&orig_src_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 5, 4); + uint16_t echo_id; memcpy(&echo_id, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 9, 2); + uint16_t echo_seq; memcpy(&echo_seq, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 11, 2); + uint8_t* payload = entry->dgram + ICMP_PROXY_RECV_HDR_SIZE; + size_t payload_len = entry->len - ICMP_PROXY_RECV_HDR_SIZE; if (g_icmp_ctx->raw_sock != SOCKET_INVALID) { exit_send_echo(inst, client_node_id, dst_ip, orig_src_ip, echo_id, echo_seq, payload, payload_len); @@ -175,11 +174,10 @@ static void exit_handle_request(struct ETCP_CONN* conn, struct ll_entry* entry) if (e->dgram) { e->dgram[0] = ETCP_RT_ID_ICMP_PROXY; e->dgram[1] = ICMP_PROXY_SUBCMD_REPLY; - memcpy(e->dgram + 2, &client_node_id, 8); - memcpy(e->dgram + 10, &dst_ip, 4); - memcpy(e->dgram + 14, &orig_src_ip, 4); - memcpy(e->dgram + 18, &echo_id, 2); - memcpy(e->dgram + 20, &echo_seq, 2); + memcpy(e->dgram + 2, &dst_ip, 4); + memcpy(e->dgram + 6, &orig_src_ip, 4); + memcpy(e->dgram + 10, &echo_id, 2); + memcpy(e->dgram + 12, &echo_seq, 2); if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len); e->len = ICMP_PROXY_HDR_SIZE + payload_len; etcp_route_send(inst, TOPO_GROUP_UTUN, client_node_id, e, 0); @@ -197,19 +195,19 @@ drop: // Сторона клиента: принять REPLY, доставить echo ответ в TUN // ==================================================================== static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (entry->len < ICMP_PROXY_HDR_SIZE + 1) { queue_dgram_free(entry); queue_entry_free(entry); return; } + if (entry->len < ICMP_PROXY_RECV_HDR_SIZE + 1) { queue_dgram_free(entry); queue_entry_free(entry); return; } uint32_t src_ip; uint32_t orig_src_ip; uint16_t echo_id, echo_seq; - memcpy(&src_ip, entry->dgram + 10, 4); - memcpy(&orig_src_ip, entry->dgram + 14, 4); - memcpy(&echo_id, entry->dgram + 18, 2); - memcpy(&echo_seq, entry->dgram + 20, 2); + memcpy(&src_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); + memcpy(&orig_src_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 5, 4); + memcpy(&echo_id, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 9, 2); + memcpy(&echo_seq, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 11, 2); DEBUG_INFO(DEBUG_CATEGORY_PROXY, "icmp_proxy: client got reply id=0x%04x seq=%u from=0x%08x dst=0x%08x", echo_id, echo_seq, src_ip, orig_src_ip); struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_icmp_ctx ? g_icmp_ctx->inst : NULL); icmp_proxy_deliver_reply(inst, orig_src_ip, src_ip, echo_id, echo_seq, - entry->dgram + ICMP_PROXY_HDR_SIZE, entry->len - ICMP_PROXY_HDR_SIZE); + entry->dgram + ICMP_PROXY_RECV_HDR_SIZE, entry->len - ICMP_PROXY_RECV_HDR_SIZE); queue_dgram_free(entry); queue_entry_free(entry); } @@ -217,11 +215,11 @@ static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry) // Единый etcp_router коллбэк // ==================================================================== void icmp_proxy_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || !entry->dgram || entry->len < 2) { + if (!entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - uint8_t subcmd = entry->dgram[1]; + uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF]; if (subcmd == ICMP_PROXY_SUBCMD_REQUEST) { exit_handle_request(conn, entry); return; } if (subcmd == ICMP_PROXY_SUBCMD_REPLY) { client_handle_reply(conn, entry); return; } queue_dgram_free(entry); queue_entry_free(entry); @@ -240,11 +238,10 @@ int icmp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id, if (!e->dgram) { queue_entry_free(e); return -1; } e->dgram[0] = ETCP_RT_ID_ICMP_PROXY; e->dgram[1] = ICMP_PROXY_SUBCMD_REQUEST; - memcpy(e->dgram + 2, &inst->node_id, 8); - memcpy(e->dgram + 10, &dst_ip, 4); - memcpy(e->dgram + 14, &orig_src_ip, 4); - memcpy(e->dgram + 18, &echo_id, 2); - memcpy(e->dgram + 20, &echo_seq, 2); + memcpy(e->dgram + 2, &dst_ip, 4); + memcpy(e->dgram + 6, &orig_src_ip, 4); + memcpy(e->dgram + 10, &echo_id, 2); + memcpy(e->dgram + 12, &echo_seq, 2); if (payload_len > 0) memcpy(e->dgram + ICMP_PROXY_HDR_SIZE, payload, payload_len); e->len = ICMP_PROXY_HDR_SIZE + payload_len; return etcp_route_send(inst, TOPO_GROUP_UTUN, exit_node_id, e, 0); diff --git a/src/proxy/icmp_proxy.h b/src/proxy/icmp_proxy.h index 2e2cd392..d195cc78 100644 --- a/src/proxy/icmp_proxy.h +++ b/src/proxy/icmp_proxy.h @@ -9,6 +9,7 @@ extern "C" { #include #include "../lib/socket_compat.h" +#include "../routing_layer/etcp_router.h" struct UTUN_INSTANCE; struct UASYNC; @@ -19,9 +20,12 @@ struct ETCP_CONN; #define ICMP_PROXY_SUBCMD_REQUEST 0x01 // client→exit: пингани dst_ip #define ICMP_PROXY_SUBCMD_REPLY 0x02 // exit→client: echo reply -// Заголовок сообщения (не включая байт svc_id) -// svc_id(1) + subcmd(1) + sender_node_id(8) + dst_ip(4) + orig_src_ip(4) + icmp_id(2) + icmp_seq(2) + payload -#define ICMP_PROXY_HDR_SIZE 22 +// Заголовок сообщения на отправке (включая байт svc_id; src/dst node_id добавляет роутер): +// svc_id(1) + subcmd(1) + dst_ip(4) + orig_src_ip(4) + icmp_id(2) + icmp_seq(2) + payload +#define ICMP_PROXY_HDR_SIZE 14 +// Полный заголовок на приёме (после router_deliver): +// svc_id(1) + src_node_id(8) + dst_node_id(8) + subcmd(1) + dst_ip(4) + orig_src_ip(4) + icmp_id(2) + icmp_seq(2) +#define ICMP_PROXY_RECV_HDR_SIZE (ROUTER_SVC_PAYLOAD_OFF + 13) // 30 // Активный ICMP echo запрос (сторона exit узла) struct icmp_request { diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index bc3fc31d..e09dbb56 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -45,8 +45,10 @@ static struct ll_entry* tcp_proxy_client_entry_from_data(struct memory_pool* poo static void tcp_proxy_client_conn_free(struct tcp_proxy_client_conn *pc); static void tcp_proxy_client_conn_finish(struct tcp_proxy_client_conn *pc); static int tcp_proxy_client_send_msg(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst, uint8_t subcmd, uint32_t sid, const uint8_t* data, size_t len, int force); -static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len); -static void tcp_proxy_client_tx_waiter_cb(struct ll_queue* q, void* arg); +static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len, int force); +static void tcp_proxy_client_tx_queue_drain_cb(struct ll_queue* q, void* arg); +static void tcp_proxy_client_pause_resume_cb(struct ll_queue* q, void* arg); +static void tcp_proxy_client_retry_timer_cb(void* arg); // ==================================================================== // Помощник ll_entry @@ -90,10 +92,17 @@ static int tcp_proxy_client_send_connect(struct tcp_proxy_client_conn* pc) { TCP_PROXY_SUBCMD_CONNECT, pc->stream_id, buf, 6, 1); } -static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len) { - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY SEND sid=%08x len=%u", pc->stream_id, len); - return tcp_proxy_client_send_msg(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, - TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len, 0); +static int tcp_proxy_client_send_data(struct tcp_proxy_client_conn* pc, const uint8_t* data, uint16_t len, int force) { + int ret = tcp_proxy_client_send_msg(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, + TCP_PROXY_SUBCMD_DATA, pc->stream_id, data, len, force); + if (ret == 0) pc->bytes_to_exit += len; + else { + pc->bp_count++; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "PROXY BP send-fail sid=%08x len=%u ret=%d tot=%u bp#=%u rcv_wnd=%u", + pc->stream_id, len, ret, pc->bytes_to_exit, pc->bp_count, + pc->pcb ? pc->pcb->rcv_wnd : 0); + } + return ret; } static int tcp_proxy_client_send_close(struct tcp_proxy_client_conn* pc) { @@ -238,6 +247,13 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t tcp_nagle_disable(newpcb); pc->to_lwip = queue_new(p->ua, 0, 0, 0, "to_lwip"); + pc->tx_queue = queue_new(p->ua, 0, 0, 0, "tx_queue"); + if (!pc->tx_queue) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy client: tx_queue alloc failed sid=%08x", pc->stream_id); + queue_free(pc->to_lwip); pc->to_lwip = NULL; + tcp_arg(newpcb, NULL); tcp_abort(newpcb); u_free(pc); return LERR_MEM; + } + queue_set_callback(pc->tx_queue, tcp_proxy_client_tx_queue_drain_cb, pc); if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); tcp_arg(newpcb, NULL); tcp_abort(newpcb); queue_free(pc->to_lwip); u_free(pc); return LERR_MEM; } @@ -248,23 +264,83 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t } // ==================================================================== -// Backpressure: waiter callback при освобождении normalizer очереди +// Backpressure: буфер lwIP→ETCP (tx_queue) с пачковым восстановлением окна. +// Зеркало tcp_proxy_server.c read_queue_drain_cb + pause_resume_cb + retry_timer_cb. // ==================================================================== -static void tcp_proxy_client_tx_waiter_cb(struct ll_queue* q, void* arg) { + +// После отправки последних данных (очередь пуста): релеить FIN/CLOSE в exit. +// Возвращает 1, если соединение завершено и pc освобождён (дальше pc/q использовать нельзя). +static int tcp_proxy_client_fin_flush(struct tcp_proxy_client_conn* pc) { + if (pc->fin_local && !pc->close_sent && !pc->close_pending) { + if (pc->fin_remote || pc->rem_closed) + tcp_proxy_client_send_close(pc); + else + tcp_proxy_client_send_fin(pc); + } + if (pc->fin_local && pc->fin_remote) { + tcp_proxy_client_conn_finish(pc); + return 1; + } + return 0; +} + +// Дрейн tx_queue → ETCP. tcp_recved вызывается только после успешной отправки: +// окно закрывается плавно (штатная backpressure), а при возобновлении открывается +// пачкой (за 2 пакета wnd_inflation ≥ TCP_WND_UPDATE_THRESHOLD → немедленный ACK). +static void tcp_proxy_client_tx_queue_drain_cb(struct ll_queue* q, void* arg) { + struct tcp_proxy_client_conn* pc = (struct tcp_proxy_client_conn*)arg; + if (!pc || !pc->proxy || pc->rem_closed) return; + + struct ll_entry* e = queue_data_get(q); + if (!e) { queue_resume_callback(q); return; } + + int ret = tcp_proxy_client_send_data(pc, e->dgram, e->len, 0); + if (ret == 0) { + if (pc->pcb) tcp_recved(pc->pcb, e->len); + queue_dgram_free(e); queue_entry_free(e); + if (queue_entry_count(q) == 0 && tcp_proxy_client_fin_flush(pc)) return; + queue_resume_callback(q); + } else { + queue_data_put_first(q, e); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "PROXY BP stall sid=%08x q=%d rcv_wnd=%u", + pc->stream_id, queue_entry_count(q), pc->pcb ? pc->pcb->rcv_wnd : 0); + etcp_router_on_send_ready(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, + ETCP_RT_ID_TCP_PROXY_SERVER, &pc->tx_waiter, + tcp_proxy_client_pause_resume_cb, pc); + if (!pc->tx_retry_timer) + pc->tx_retry_timer = uasync_set_timeout(pc->proxy->ua, 5000, pc, tcp_proxy_client_retry_timer_cb, "tpc_retry"); + } +} + +// send_q освободился — возобновить дрейн tx_queue +static void tcp_proxy_client_pause_resume_cb(struct ll_queue* q, void* arg) { (void)q; struct tcp_proxy_client_conn* pc = (struct tcp_proxy_client_conn*)arg; - if (!pc || !pc->proxy || pc->rem_closed || !pc->tx_buf) return; - int ret = tcp_proxy_client_send_data(pc, pc->tx_buf, pc->tx_len); + if (!pc || !pc->proxy || pc->rem_closed || !pc->tx_queue) return; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "PROXY BP resume sid=%08x q=%d", + pc->stream_id, queue_entry_count(pc->tx_queue)); + queue_resume_callback(pc->tx_queue); +} + +// Safety-net: форсированный ретрай, если waiter не сработал +static void tcp_proxy_client_retry_timer_cb(void* arg) { + struct tcp_proxy_client_conn* pc = (struct tcp_proxy_client_conn*)arg; + if (!pc || !pc->proxy) return; + pc->tx_retry_timer = NULL; + if (pc->rem_closed || !pc->tx_queue) return; + struct ll_entry* e = queue_data_get(pc->tx_queue); + if (!e) { queue_resume_callback(pc->tx_queue); return; } + int ret = tcp_proxy_client_send_data(pc, e->dgram, e->len, 1); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "PROXY BP retry(force) sid=%08x len=%u ret=%d", + pc->stream_id, e->len, ret); if (ret == 0) { - if (pc->pcb) tcp_recved(pc->pcb, pc->tx_len); - u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; - if (pc->fin_local && !pc->close_sent && !pc->close_pending) { - if (pc->fin_remote || pc->rem_closed) - tcp_proxy_client_send_close(pc); - else - tcp_proxy_client_send_fin(pc); - } - if (pc->fin_local && pc->fin_remote) tcp_proxy_client_conn_finish(pc); + if (pc->pcb) tcp_recved(pc->pcb, e->len); + queue_dgram_free(e); queue_entry_free(e); + if (queue_entry_count(pc->tx_queue) == 0 && tcp_proxy_client_fin_flush(pc)) return; + queue_resume_callback(pc->tx_queue); + } else { + queue_data_put_first(pc->tx_queue, e); + pc->tx_retry_timer = uasync_set_timeout(pc->proxy->ua, 5000, pc, tcp_proxy_client_retry_timer_cb, "tpc_retry"); } } @@ -275,16 +351,13 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu if (p == NULL || err != LERR_OK) { pc->fin_local = 1; if (p == NULL && err == ERR_OK) { - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN sid=%08x pcb_state=%u sndbuf=%u cwnd=%u unsent=%p unacked=%p", - pc->stream_id, pcb->state, pcb->snd_buf, pcb->cwnd, - (void*)pcb->unsent, (void*)pcb->unacked); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY FIN sid=%08x to_exit=%u from_exit=%u bp#=%u pcb_state=%u sndbuf=%u cwnd=%u", + pc->stream_id, pc->bytes_to_exit, pc->bytes_from_exit, pc->bp_count, + pcb->state, pcb->snd_buf, pcb->cwnd); { uint16_t wnd_gap = TCP_WND_MAX(pcb) - pcb->rcv_wnd; if (wnd_gap > 0) tcp_recved(pcb, wnd_gap); } - if (!pc->tx_buf) { - tcp_proxy_client_send_fin(pc); - if (pc->fin_remote && !pc->close_sent && !pc->close_pending) - tcp_proxy_client_send_close(pc); - if (pc->fin_remote) { tcp_proxy_client_conn_finish(pc); return LERR_OK; } + if (pc->tx_queue && queue_entry_count(pc->tx_queue) == 0) { + if (tcp_proxy_client_fin_flush(pc)) return LERR_OK; } } else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY ERROR recv sid=%08x pcb_state=%u err=%d", @@ -300,20 +373,19 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu tcp_recved(pcb, len); pbuf_free(p); return LERR_OK; } - if (pc->tx_buf) { pbuf_free(p); return LERR_OK; } - uint8_t *data = u_malloc(len); - if (data) { - pbuf_copy_partial(p, data, len, 0); - DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY RECV client->exit sid=%08x len=%u", pc->stream_id, len); - int ret = tcp_proxy_client_send_data(pc, data, len); - if (ret == 0) { tcp_recved(pcb, len); u_free(data); } - else { - pc->tx_buf = data; pc->tx_len = len; - etcp_router_waiter_register(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, &pc->tx_waiter, - tcp_proxy_client_tx_waiter_cb, pc); - } - } else DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY recv malloc(%u) failed sid=%08x", len, pc->stream_id); + if (!data) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY recv malloc(%u) failed sid=%08x", len, pc->stream_id); + pbuf_free(p); return LERR_OK; + } + pbuf_copy_partial(p, data, len, 0); + struct ll_entry* e = queue_entry_new_from_pool(pc->proxy->entry_pool); + if (!e) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "PROXY recv entry pool exhausted sid=%08x", pc->stream_id); + u_free(data); pbuf_free(p); return LERR_OK; + } + e->dgram = data; e->len = len; + queue_data_put(pc->tx_queue, e); pbuf_free(p); return LERR_OK; } @@ -435,8 +507,13 @@ static void tcp_proxy_client_conn_free(struct tcp_proxy_client_conn *pc) { while ((e = queue_data_get(pc->to_lwip))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->to_lwip); pc->to_lwip = NULL; } - if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; } - if (pc->proxy && pc->proxy->inst) etcp_router_waiter_cancel(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, &pc->tx_waiter); + if (pc->tx_queue) { + struct ll_entry *e; + while ((e = queue_data_get(pc->tx_queue))) { queue_dgram_free(e); queue_entry_free(e); } + queue_free(pc->tx_queue); pc->tx_queue = NULL; + } + if (pc->tx_retry_timer) { uasync_cancel_timeout(p->ua, pc->tx_retry_timer); pc->tx_retry_timer = NULL; } + if (p && p->inst) etcp_router_cancel_send_ready(p->inst, TOPO_GROUP_UTUN, p->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &pc->tx_waiter); u_free(pc); } @@ -471,10 +548,11 @@ static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY DATA sid=%08x — no conn/closed, dropping", stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; } - size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; + size_t data_len = entry->len - TCP_PROXY_RECV_HDR_SIZE; DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY DATA <- sid=%08x len=%zu", stream_id, data_len); if (data_len > 0) { - struct ll_entry* e = tcp_proxy_client_entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_HDR_SIZE, (uint16_t)data_len); + pc->bytes_from_exit += (uint32_t)data_len; + struct ll_entry* e = tcp_proxy_client_entry_from_data(pc->proxy->entry_pool, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, (uint16_t)data_len); if (e) queue_data_put(pc->to_lwip, e); tcp_proxy_client_feed_from_transport(pc); } @@ -486,8 +564,6 @@ static void tcp_proxy_client_handle_close(struct tcp_proxy_client* p, uint32_t s if (!pc) { DEBUG_WARN(DEBUG_CATEGORY_PROXY, "PROXY CLOSE sid=%08x — no conn, dropping", stream_id); return; } DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY REM_CLOSED sid=%08x fin_local=%d pcb_state=%u", stream_id, pc->fin_local, pc->pcb ? pc->pcb->state : 0); pc->rem_closed = 1; - if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; } - etcp_router_waiter_cancel(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, &pc->tx_waiter); if (pc->pcb) { tcp_arg(pc->pcb, NULL); tcp_recv(pc->pcb, NULL); @@ -514,8 +590,6 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s pc->pcb = NULL; } pc->rem_closed = 1; - if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; } - etcp_router_waiter_cancel(pc->proxy->inst, TOPO_GROUP_UTUN, pc->proxy->via_node_id, &pc->tx_waiter); tcp_proxy_client_conn_free(pc); } @@ -541,49 +615,56 @@ void tcp_proxy_client_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* en struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL; struct tcp_proxy_client* proxy = inst ? inst->tcp_proxy_client : NULL; - if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) { - if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY_CLIENT) { - uint64_t peer_id; - memcpy(&peer_id, entry->dgram + 1, 8); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY CLOSE_ALL from %016llx — clearing client conns for peer", - (unsigned long long)peer_id); - if (proxy) { - struct tcp_proxy_client_conn *pc, *next; - for (pc = proxy->conns; pc; pc = next) { - next = pc->next; - if (pc->pcb) { - tcp_arg(pc->pcb, NULL); - tcp_recv(pc->pcb, NULL); - tcp_sent(pc->pcb, NULL); - tcp_err(pc->pcb, NULL); - tcp_poll(pc->pcb, NULL, 0); - tcp_abort(pc->pcb); - pc->pcb = NULL; - } - tcp_proxy_client_conn_free(pc); + if (!entry || !entry->dgram) { + if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + return; + } + // CLOSE_ALL / restart-уведомление: [svc_id][src][dst] без payload + if (entry->len == ROUTER_SVC_HDR_SIZE && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY_CLIENT) { + uint64_t peer_id; + memcpy(&peer_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY CLOSE_ALL from %016llx — clearing client conns for peer", + (unsigned long long)peer_id); + if (proxy) { + struct tcp_proxy_client_conn *pc, *next; + for (pc = proxy->conns; pc; pc = next) { + next = pc->next; + if (pc->pcb) { + tcp_arg(pc->pcb, NULL); + tcp_recv(pc->pcb, NULL); + tcp_sent(pc->pcb, NULL); + tcp_err(pc->pcb, NULL); + tcp_poll(pc->pcb, NULL, 0); + tcp_abort(pc->pcb); + pc->pcb = NULL; } - proxy->conns = NULL; proxy->conn_count = 0; - socks_proxy_conn_free_all(&proxy->socks_conns, &proxy->socks_conn_count); - socks_proxy_conn_free_all(&proxy->http_conns, &proxy->http_conn_count); + tcp_proxy_client_conn_free(pc); } + proxy->conns = NULL; proxy->conn_count = 0; + socks_proxy_conn_free_all(&proxy->socks_conns, &proxy->socks_conn_count); + socks_proxy_conn_free_all(&proxy->http_conns, &proxy->http_conn_count); } - if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + queue_dgram_free(entry); queue_entry_free(entry); return; } - uint8_t subcmd = entry->dgram[1]; - uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4); + if (entry->len < TCP_PROXY_RECV_HDR_SIZE) { + queue_dgram_free(entry); queue_entry_free(entry); + return; + } + uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF]; + uint32_t stream_id; memcpy(&stream_id, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY CLIENT RECV subcmd=%02x sid=%08x pkt_len=%u", subcmd, stream_id, entry->len); if (proxy) { - size_t data_len = entry->len > TCP_PROXY_HDR_SIZE ? entry->len - TCP_PROXY_HDR_SIZE : 0; + size_t data_len = entry->len > TCP_PROXY_RECV_HDR_SIZE ? entry->len - TCP_PROXY_RECV_HDR_SIZE : 0; // Пробуем SOCKS/HTTP first (у них приоритет — могут быть без TUN) if (proxy->socks_enabled && socks_proxy_handle_etcp(&proxy->socks_conns, &proxy->socks_conn_count, - stream_id, subcmd, entry->dgram + TCP_PROXY_HDR_SIZE, data_len)) { + stream_id, subcmd, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len)) { queue_dgram_free(entry); queue_entry_free(entry); return; } if (proxy->http_proxy_enabled && socks_proxy_handle_etcp(&proxy->http_conns, &proxy->http_conn_count, - stream_id, subcmd, entry->dgram + TCP_PROXY_HDR_SIZE, data_len)) { + stream_id, subcmd, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len)) { queue_dgram_free(entry); queue_entry_free(entry); return; } // Существующие lwIP conns @@ -707,8 +788,13 @@ void tcp_proxy_client_destroy(struct tcp_proxy_client* p) { while ((e = queue_data_get(pc->to_lwip))) { queue_dgram_free(e); queue_entry_free(e); } queue_free(pc->to_lwip); pc->to_lwip = NULL; } - if (pc->tx_buf) { u_free(pc->tx_buf); pc->tx_buf = NULL; pc->tx_len = 0; } - if (p->inst) etcp_router_waiter_cancel(p->inst, TOPO_GROUP_UTUN, p->via_node_id, &pc->tx_waiter); + if (pc->tx_queue) { + struct ll_entry *e; + while ((e = queue_data_get(pc->tx_queue))) { queue_dgram_free(e); queue_entry_free(e); } + queue_free(pc->tx_queue); pc->tx_queue = NULL; + } + if (pc->tx_retry_timer) { uasync_cancel_timeout(p->ua, pc->tx_retry_timer); pc->tx_retry_timer = NULL; } + if (p->inst) etcp_router_cancel_send_ready(p->inst, TOPO_GROUP_UTUN, p->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &pc->tx_waiter); u_free(pc); pc = next; } p->conns = NULL; diff --git a/src/proxy/tcp_proxy_client.h b/src/proxy/tcp_proxy_client.h index 9b8da4fb..141bebb6 100644 --- a/src/proxy/tcp_proxy_client.h +++ b/src/proxy/tcp_proxy_client.h @@ -43,9 +43,13 @@ struct tcp_proxy_client_conn { uint8_t close_sent; // отправили CLOSE/ERROR в exit uint8_t close_pending; // CLOSE/ERROR не доставлен, ждём повтора - uint8_t* tx_buf; // буфер при backpressure (lwIP→ETCP send fail) - uint16_t tx_len; - struct queue_waiter_handle tx_waiter; + struct ll_queue* tx_queue; // буфер lwIP→ETCP при backpressure (аналог server read_queue) + struct queue_waiter_handle tx_waiter; // waiter на send_q (etcp_router_on_send_ready) + void* tx_retry_timer; // retry safety-net (force=1), аналог server retry_timer + + uint32_t bytes_to_exit; // байт отправлено в exit (из lwIP/TUN) + uint32_t bytes_from_exit; // байт получено от exit (в lwIP/TUN) + uint32_t bp_count; // счётчик backpressure (send fail) uint8_t dest_ip[4]; uint16_t dest_port; diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index 9cd2c868..dfb6d899 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -352,8 +352,8 @@ struct tcp_proxy_server_conn* tcp_proxy_server_find_conn(struct tcp_proxy_server int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* entry, uint32_t stream_id, uint64_t src_node_id) { if (!inst || !inst->tcp_proxy_server.enabled) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return -1; } struct tcp_proxy_server* ctx = &inst->tcp_proxy_server; - if (entry->len < TCP_PROXY_CONNECT_HDR_SIZE) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy server: CONNECT too short len=%u", entry->len); queue_dgram_free(entry); queue_entry_free(entry); return -1; } - uint8_t* dest_ip = entry->dgram + TCP_PROXY_HDR_SIZE; + if (entry->len < TCP_PROXY_RECV_HDR_SIZE + 6) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy server: CONNECT too short len=%u", entry->len); queue_dgram_free(entry); queue_entry_free(entry); return -1; } + uint8_t* dest_ip = entry->dgram + TCP_PROXY_RECV_HDR_SIZE; uint16_t dest_port = 0; memcpy(&dest_port, dest_ip + 4, 2); struct tcp_proxy_server_conn* rc = u_calloc(1, sizeof(struct tcp_proxy_server_conn)); @@ -365,8 +365,9 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* socket_t sock = socket(AF_INET, SOCK_STREAM, 0); if (sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } ctx->conn_count++; - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "SOCK:NEW fd=%d sid=%08x dest=%d.%d.%d.%d:%d total=%d", - (int)sock, stream_id, dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "SOCK:NEW fd=%d sid=%08x src=%016llx dest=%d.%d.%d.%d:%d total=%d", + (int)sock, stream_id, (unsigned long long)src_node_id, + dest_ip[0],dest_ip[1],dest_ip[2],dest_ip[3],ntohs(dest_port), ctx->conn_count); socket_set_nonblocking(sock); struct global_config* g = &inst->config->global; @@ -413,6 +414,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* } rc->next = ctx->conns; ctx->conns = rc; + tcp_conn_set_connect_timeout(rc->tc, g->tcp_proxy_server_connect_timeout_ms); rc->diag_timer = uasync_set_timeout(rc->ua, 10000, rc, diag_timer_cb, "tps_diag"); queue_dgram_free(entry); queue_entry_free(entry); @@ -428,7 +430,7 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c DEBUG_INFO(DEBUG_CATEGORY_PROXY, "TPS handle_data: no/closed conn sid=%08x, dropping", stream_id); queue_dgram_free(entry); queue_entry_free(entry); return -1; } - size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; + size_t data_len = entry->len - TCP_PROXY_RECV_HDR_SIZE; rc->bytes_sent += (uint32_t)data_len; rc->data_count++; DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "TPS DATA #%u sid=%08x len=%zu", rc->data_count, stream_id, data_len); @@ -441,7 +443,7 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c struct ll_entry* e = queue_entry_new_from_pool(rc->tc->entry_pool); uint8_t* buf = memory_pool_alloc(rc->tc->data_pool); if (e && buf) { - memcpy(buf, entry->dgram + TCP_PROXY_HDR_SIZE, data_len); + memcpy(buf, entry->dgram + TCP_PROXY_RECV_HDR_SIZE, data_len); e->dgram = buf; e->len = (uint16_t)data_len; queue_data_put(rc->tc->write_queue, e); } else { @@ -503,30 +505,39 @@ void tcp_proxy_server_handle_fin(struct UTUN_INSTANCE* inst, uint32_t stream_id) void tcp_proxy_server_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL; - if (!entry || !entry->dgram || entry->len < TCP_PROXY_HDR_SIZE) { - if (entry && entry->len == 9 && !conn && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY_SERVER) { - uint64_t peer_id; - memcpy(&peer_id, entry->dgram + 1, 8); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY CLOSE_ALL from %016llx — clearing server conns for peer", - (unsigned long long)peer_id); - if (inst && inst->tcp_proxy_server.enabled) { - struct tcp_proxy_server_conn *rc, *next; - for (rc = inst->tcp_proxy_server.conns; rc; rc = next) { - next = rc->next; - if (rc->peer_node_id == peer_id) tcp_proxy_server_conn_free(rc); - } + if (!entry || !entry->dgram) { + if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + return; + } + // CLOSE_ALL / restart-уведомление: [svc_id][src][dst] без payload + if (entry->len == ROUTER_SVC_HDR_SIZE && entry->dgram[0] == ETCP_RT_ID_TCP_PROXY_SERVER) { + uint64_t peer_id; + memcpy(&peer_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "PROXY CLOSE_ALL from %016llx — clearing server conns for peer", + (unsigned long long)peer_id); + if (inst && inst->tcp_proxy_server.enabled) { + struct tcp_proxy_server_conn *rc, *next; + for (rc = inst->tcp_proxy_server.conns; rc; rc = next) { + next = rc->next; + if (rc->peer_node_id == peer_id) tcp_proxy_server_conn_free(rc); } } - if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } + queue_dgram_free(entry); queue_entry_free(entry); + return; + } + if (entry->len < TCP_PROXY_RECV_HDR_SIZE) { + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "TCP proxy server: short packet len=%u", entry->len); + queue_dgram_free(entry); queue_entry_free(entry); return; } - uint8_t subcmd = entry->dgram[1]; - uint32_t stream_id; memcpy(&stream_id, entry->dgram + 2, 4); + uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF]; + uint32_t stream_id; memcpy(&stream_id, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "PROXY SERVER RECV subcmd=%02x sid=%08x pkt_len=%u", subcmd, stream_id, entry->len); if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { - uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0); + uint64_t src_node_id; + memcpy(&src_node_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); tcp_proxy_server_handle_connect(inst, entry, stream_id, src_node_id); return; } diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index b9016e84..3c3291a1 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -11,6 +11,7 @@ extern "C" { #include "../lib/socket_compat.h" #include "../lib/ll_queue.h" #include "../lib/tcp_io.h" +#include "../routing_layer/etcp_router.h" struct UTUN_INSTANCE; struct UASYNC; @@ -26,6 +27,9 @@ struct ETCP_CONN; #define TCP_PROXY_HDR_SIZE 6 // svc_id(1)+subcmd(1)+stream_id(4) #define TCP_PROXY_CONNECT_HDR_SIZE 12 // HDR_SIZE + dest_ip(4)+dest_port(2) +// recv-формат (после router_deliver): [svc_id][src][dst][subcmd][stream_id][data] +#define TCP_PROXY_RECV_HDR_SIZE (ROUTER_SVC_PAYLOAD_OFF + 5) // 22: до data (subcmd+stream_id) + struct tcp_proxy_server_conn { struct tcp_proxy_server_conn* next; struct tcp_proxy_server* ctx; diff --git a/src/proxy/tcp_proxy_server_doc.md b/src/proxy/tcp_proxy_server_doc.md index bf521b76..435ee9a4 100644 --- a/src/proxy/tcp_proxy_server_doc.md +++ b/src/proxy/tcp_proxy_server_doc.md @@ -28,6 +28,7 @@ TCP прокси-сервер на exit node. Принимает CONNECT-зап - **Graceful CLOSE:** при получении CLOSE от клиента чтение из OS сокета паузится через `tcp_conn_pause_read()`, read_queue сливается (данные отбрасываются), push-ится close в OS сокет. - **Ретрай CLOSE/ERROR:** при неудачной отправке `send_msg()` — экспоненциальный backoff (50→100→200→…→5000 tb = 5ms→500ms) через `close_timer` + `close_retry_cb`. - **Диагностический таймер (1s):** логирует TCP буферы ОС, размеры read/write очередей, free counts entry/data pool-ов. +- **Connect-таймаут:** `connect()` к destination асинхронный; если за `connect_timeout` мс (конфиг `[tcp_proxy_server]`, по умолчанию 2000, `0` = без таймаута) соединение не установилось — срабатывает `tcp_conn_handle_error(ETIMEDOUT)` → клиенту уходит ERROR, соединение освобождается (защита от «чёрной дыры» и утечки fd). - **CLOSE_ALL:** специальный 9-байтовый пакет без валидного conn удаляет все соединения для заданного peer. ## 3. API diff --git a/src/proxy/udp_proxy.c b/src/proxy/udp_proxy.c index 8a3710c4..32d9cd5f 100644 --- a/src/proxy/udp_proxy.c +++ b/src/proxy/udp_proxy.c @@ -55,13 +55,12 @@ static void flow_read_cb(socket_t sock, void* arg) { if (!e->dgram) { queue_entry_free(e); return; } e->dgram[0] = ETCP_RT_ID_UDP_PROXY; e->dgram[1] = UDP_PROXY_SUBCMD_DATA; - memcpy(e->dgram + 2, &f->client_node_id, 8); // Ответ src = оригинальный dst_ip:dest_port - memcpy(e->dgram + 10, &f->dst_ip, 4); - memcpy(e->dgram + 14, &f->dst_port, 2); + memcpy(e->dgram + 2, &f->dst_ip, 4); + memcpy(e->dgram + 6, &f->dst_port, 2); // Ответ dst = оригинальный src_ip:src_port - memcpy(e->dgram + 16, &f->src_ip, 4); - memcpy(e->dgram + 20, &f->src_port, 2); + memcpy(e->dgram + 8, &f->src_ip, 4); + memcpy(e->dgram + 12, &f->src_port, 2); memcpy(e->dgram + UDP_PROXY_HDR_SIZE, buf, n); e->len = UDP_PROXY_HDR_SIZE + n; @@ -73,15 +72,15 @@ static void flow_read_cb(socket_t sock, void* arg) { // ==================================================================== static void exit_handle_data(struct ETCP_CONN* conn, struct ll_entry* entry) { struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_udp_ctx ? g_udp_ctx->inst : NULL); - if (!inst || !g_udp_ctx || entry->len < UDP_PROXY_HDR_SIZE + 1) goto drop; + if (!inst || !g_udp_ctx || entry->len < UDP_PROXY_RECV_HDR_SIZE + 1) goto drop; - uint64_t client_node_id; memcpy(&client_node_id, entry->dgram + 2, 8); - uint32_t src_ip; memcpy(&src_ip, entry->dgram + 10, 4); - uint16_t src_port; memcpy(&src_port, entry->dgram + 14, 2); - uint32_t dst_ip; memcpy(&dst_ip, entry->dgram + 16, 4); - uint16_t dst_port; memcpy(&dst_port, entry->dgram + 20, 2); - uint8_t* payload = entry->dgram + UDP_PROXY_HDR_SIZE; - size_t payload_len = entry->len - UDP_PROXY_HDR_SIZE; + uint64_t client_node_id; memcpy(&client_node_id, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + uint32_t src_ip; memcpy(&src_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); + uint16_t src_port; memcpy(&src_port, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 5, 2); + uint32_t dst_ip; memcpy(&dst_ip, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 7, 4); + uint16_t dst_port; memcpy(&dst_port, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 11, 2); + uint8_t* payload = entry->dgram + UDP_PROXY_RECV_HDR_SIZE; + size_t payload_len = entry->len - UDP_PROXY_RECV_HDR_SIZE; struct udp_flow* f = flow_find(g_udp_ctx->flows, client_node_id, src_ip, src_port, dst_ip, dst_port); if (!f) { @@ -115,17 +114,17 @@ drop: // Сторона клиента: принять REPLY, доставить в TUN // ==================================================================== static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (entry->len < UDP_PROXY_HDR_SIZE + 1) { queue_dgram_free(entry); queue_entry_free(entry); return; } - // svc_id уже обработан диспетчером etcp_router + if (entry->len < UDP_PROXY_RECV_HDR_SIZE + 1) { queue_dgram_free(entry); queue_entry_free(entry); return; } + const uint8_t* d = entry->dgram; uint32_t src_ip, dst_ip; uint16_t src_port, dst_port; - uint8_t* payload = entry->dgram + 1 + 1 + 8 + 4 + 2 + 4 + 2; // svc_id + subcmd + node_id + src_ip:port + dst_ip:port - size_t payload_len = entry->len - UDP_PROXY_HDR_SIZE; - - memcpy(&dst_ip, entry->dgram + 1 + 1 + 8 + 4 + 2, 4); // src в сообщении → dst в ответе - memcpy(&dst_port, entry->dgram + 1 + 1 + 8 + 4 + 2 + 4, 2); - memcpy(&src_ip, entry->dgram + 1 + 1 + 8, 4); // dst в сообщении → src в ответе - memcpy(&src_port, entry->dgram + 1 + 1 + 8 + 4, 2); + // В сообщении: src_ip:src_port = оригинальный dst, dst_ip:dst_port = оригинальный src + memcpy(&src_ip, d + ROUTER_SVC_PAYLOAD_OFF + 1, 4); + memcpy(&src_port, d + ROUTER_SVC_PAYLOAD_OFF + 5, 2); + memcpy(&dst_ip, d + ROUTER_SVC_PAYLOAD_OFF + 7, 4); + memcpy(&dst_port, d + ROUTER_SVC_PAYLOAD_OFF + 11, 2); + uint8_t* payload = d + UDP_PROXY_RECV_HDR_SIZE; + size_t payload_len = entry->len - UDP_PROXY_RECV_HDR_SIZE; struct UTUN_INSTANCE* inst = conn ? conn->instance : (g_udp_ctx ? g_udp_ctx->inst : NULL); udp_proxy_deliver_reply(inst, src_ip, src_port, dst_ip, dst_port, payload, payload_len); @@ -136,11 +135,11 @@ static void client_handle_reply(struct ETCP_CONN* conn, struct ll_entry* entry) // Единый etcp_router коллбэк // ==================================================================== void udp_proxy_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || !entry->dgram || entry->len < 2) { + if (!entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - uint8_t subcmd = entry->dgram[1]; + uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF]; if (subcmd == UDP_PROXY_SUBCMD_DATA) { if (g_udp_ctx && g_udp_ctx->is_exit) exit_handle_data(conn, entry); else client_handle_reply(conn, entry); @@ -163,11 +162,10 @@ int udp_proxy_send_to_exit(struct UTUN_INSTANCE* inst, uint64_t exit_node_id, if (!e->dgram) { queue_entry_free(e); return -1; } e->dgram[0] = ETCP_RT_ID_UDP_PROXY; e->dgram[1] = UDP_PROXY_SUBCMD_DATA; - memcpy(e->dgram + 2, &inst->node_id, 8); - memcpy(e->dgram + 10, &src_ip, 4); - memcpy(e->dgram + 14, &src_port, 2); - memcpy(e->dgram + 16, &dst_ip, 4); - memcpy(e->dgram + 20, &dst_port, 2); + memcpy(e->dgram + 2, &src_ip, 4); + memcpy(e->dgram + 6, &src_port, 2); + memcpy(e->dgram + 8, &dst_ip, 4); + memcpy(e->dgram + 12, &dst_port, 2); if (payload_len > 0) memcpy(e->dgram + UDP_PROXY_HDR_SIZE, payload, payload_len); e->len = UDP_PROXY_HDR_SIZE + payload_len; return etcp_route_send(inst, TOPO_GROUP_UTUN, exit_node_id, e, 0); diff --git a/src/proxy/udp_proxy.h b/src/proxy/udp_proxy.h index f3910266..4669e3b0 100644 --- a/src/proxy/udp_proxy.h +++ b/src/proxy/udp_proxy.h @@ -9,6 +9,7 @@ extern "C" { #include #include "../lib/socket_compat.h" +#include "../routing_layer/etcp_router.h" struct UTUN_INSTANCE; struct UASYNC; @@ -18,9 +19,12 @@ struct ETCP_CONN; // Подкоманды #define UDP_PROXY_SUBCMD_DATA 0x01 // client→exit: пробрось датаграмму | exit→client: ответ -// Заголовок сообщения (не включая байт svc_id) -// svc_id(1) + subcmd(1) + sender_node_id(8) + src_ip(4) + src_port(2) + dst_ip(4) + dst_port(2) + payload -#define UDP_PROXY_HDR_SIZE 22 +// Заголовок сообщения на отправке (включая байт svc_id; src/dst node_id добавляет роутер): +// svc_id(1) + subcmd(1) + src_ip(4) + src_port(2) + dst_ip(4) + dst_port(2) + payload +#define UDP_PROXY_HDR_SIZE 14 +// Полный заголовок на приёме (после router_deliver): +// svc_id(1) + src_node_id(8) + dst_node_id(8) + subcmd(1) + src_ip(4) + src_port(2) + dst_ip(4) + dst_port(2) +#define UDP_PROXY_RECV_HDR_SIZE (ROUTER_SVC_PAYLOAD_OFF + 13) // 30 // Один UDP поток (сторона exit узла: сопоставляет client_node_id + кортеж адресов с OS сокетом) struct udp_flow { diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index 6f610999..8a32e46c 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/src/routing_layer/conn_mgr_core.c @@ -725,12 +725,12 @@ void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id) { } void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!entry || !entry->dgram || entry->len < 3) return; + if (!entry || !entry->dgram || entry->len < ROUTER_SVC_PAYLOAD_OFF + 2) return; if (!conn || !conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "conn_mgr_router_recv: conn=%p", (void*)conn); return; } - uint8_t* d = entry->dgram + 1; size_t len = entry->len - 1; uint8_t sub = d[1]; + uint8_t* d = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF; uint8_t sub = d[1]; struct CONN_MGR* mgr = conn->instance->conn_mgr; if (mgr) { struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, conn->peer_node_id); if (e && e->state == CONN_MGR_STATE_CONNECTED) e->last_traffic_tb = get_time_tb(); } switch (sub) { diff --git a/src/routing_layer/etcp_router.c b/src/routing_layer/etcp_router.c index 22f546d2..fc0d4201 100644 --- a/src/routing_layer/etcp_router.c +++ b/src/routing_layer/etcp_router.c @@ -104,8 +104,12 @@ static struct TRANSIT_QUEUE* transit_queue_get_or_create(struct ETCP_CONN* conn, static void transit_queue_destroy(struct ETCP_CONN* conn, struct TRANSIT_QUEUE* tq) { if (conn->send_input_q) queue_waiter_cancel(conn->send_input_q, &tq->waiter); + int dropped = 0; struct ll_entry* e; - while ((e = queue_data_get(tq->q)) != NULL) free_entry(e); + while ((e = queue_data_get(tq->q)) != NULL) { free_entry(e); dropped++; } + if (dropped > 0) + DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "transit_q: destroyed with %d pending pkts src=%016llx dst=%016llx", + dropped, (unsigned long long)tq->src_node_id, (unsigned long long)tq->dst_node_id); queue_free(tq->q); queue_remove_data(conn->transit_queues, &tq->ll); queue_entry_free(&tq->ll); @@ -117,7 +121,11 @@ static void transit_queue_drain_cb(struct ll_queue* q, void* arg) { struct ll_entry* e = queue_data_get(tq->q); if (!e) { transit_queue_destroy(conn, tq); return; } int ret = etcp_send(conn, e); - if (ret != 0) free_entry(e); + if (ret != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "transit_drain: etcp_send failed ret=%d src=%016llx dst=%016llx", + ret, (unsigned long long)tq->src_node_id, (unsigned long long)tq->dst_node_id); + free_entry(e); + } if (queue_entry_count(tq->q) > 0) queue_waiter_wait(conn->send_input_q, &tq->waiter, transit_queue_drain_cb, tq); else @@ -341,11 +349,12 @@ static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) { if (!cb) return; struct ll_entry* e = queue_entry_new(0); if (!e) return; - e->dgram = u_malloc(9); + e->dgram = u_malloc(ROUTER_SVC_HDR_SIZE); if (!e->dgram) { queue_entry_free(e); return; } e->dgram[0] = rconn->svc_id; - memcpy(e->dgram + 1, &rconn->remote_node_id, 8); - e->len = 9; + memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &rconn->remote_node_id, 8); + memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8); + e->len = ROUTER_SVC_HDR_SIZE; DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_close_svc: svc_id=%u remote=%016llx", rconn->svc_id, (unsigned long long)rconn->remote_node_id); cb(loopback_conn(rconn->inst), e); @@ -762,6 +771,25 @@ static void router_idle_ack_timer_cb(void* arg) { // Reorder / assembly (аналог etcp_output_try_assembly) // ==================================================================== +// Доставка сервисной кодограммы в едином формате [svc_id][src][dst][payload]. +// Для loopback (dst == self) src и dst оба = self. +static void router_deliver_loopback(struct UTUN_INSTANCE* inst, uint8_t svc_id, struct ll_entry* entry) { + etcp_recv_fn cb = inst->router_bindings.callbacks[svc_id]; + if (!cb) { queue_dgram_free(entry); queue_entry_free(entry); return; } + size_t payload_len = entry->len - 1; + struct ll_entry* e = queue_entry_new(0); + if (!e) { queue_dgram_free(entry); queue_entry_free(entry); return; } + e->dgram = u_malloc(ROUTER_SVC_HDR_SIZE + payload_len); + if (!e->dgram) { queue_entry_free(e); queue_dgram_free(entry); queue_entry_free(entry); return; } + e->dgram[0] = svc_id; + memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &inst->node_id, 8); + memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8); + if (payload_len > 0) memcpy(e->dgram + ROUTER_SVC_PAYLOAD_OFF, entry->dgram + 1, payload_len); + e->len = ROUTER_SVC_HDR_SIZE + payload_len; + queue_dgram_free(entry); queue_entry_free(entry); + cb(loopback_conn(inst), e); +} + static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* conn, const uint8_t* payload, size_t payload_len) { struct UTUN_INSTANCE* inst = rconn->inst; @@ -770,17 +798,20 @@ static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* con DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router_deliver: no handler for svc_id=%u", rconn->svc_id); return; } - size_t entry_len = 1 + payload_len; + size_t entry_len = ROUTER_SVC_HDR_SIZE + payload_len; struct ll_entry* svc_entry = queue_entry_new(0); if (!svc_entry) { DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "router_deliver: queue_entry_new failed"); return; } svc_entry->len = entry_len; svc_entry->dgram = u_malloc(entry_len); if (!svc_entry->dgram) { queue_entry_free(svc_entry); return; } svc_entry->dgram[0] = rconn->svc_id; - if (payload_len > 0) memcpy(svc_entry->dgram + 1, payload, payload_len); + memcpy(svc_entry->dgram + ROUTER_SVC_SRC_OFF, &rconn->remote_node_id, 8); + memcpy(svc_entry->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8); + if (payload_len > 0) memcpy(svc_entry->dgram + ROUTER_SVC_PAYLOAD_OFF, payload, payload_len); - DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router_deliver: svc_id=%u len=%zu from remote=%016llx conn=%p cb=%p", - rconn->svc_id, payload_len, (unsigned long long)rconn->remote_node_id, (void*)conn, (void*)cb); + DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "router_deliver: svc_id=%u len=%zu src=%016llx dst=%016llx conn=%p cb=%p", + rconn->svc_id, payload_len, (unsigned long long)rconn->remote_node_id, + (unsigned long long)inst->node_id, (void*)conn, (void*)cb); rconn->c_pkts_rcvd++; cb(conn, svc_entry); } @@ -946,7 +977,11 @@ static void router_forward_transit(struct UTUN_INSTANCE* inst, struct ll_entry* DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "etcp_router: forwarding svc_id=%u → %016llx via %s", hdr->svc_id, (unsigned long long)hdr->dst_node_id, next->log_name); int ret = etcp_send(next, entry); - if (ret != 0) free_entry(entry); + if (ret != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCPROUTE, "etcp_router: forward etcp_send failed ret=%d svc_id=%u → %016llx", + ret, hdr->svc_id, (unsigned long long)hdr->dst_node_id); + free_entry(entry); + } return; } @@ -1172,9 +1207,7 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_ if (dst_node_id == inst->node_id) { DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "loopback svc_id=%u len=%zu", svc_id, payload_len); - if (inst->router_bindings.callbacks[svc_id]) - inst->router_bindings.callbacks[svc_id](loopback_conn(inst), entry); - else { queue_dgram_free(entry); queue_entry_free(entry); } + router_deliver_loopback(inst, svc_id, entry); return 0; } @@ -1222,9 +1255,7 @@ int etcp_route_send_encrypted(struct UTUN_INSTANCE* inst, uint64_t group_id, uin if (dst_node_id == inst->node_id) { DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "E2E send: loopback svc_id=%u len=%zu", svc_id, pl_len); - if (inst->router_bindings.callbacks[svc_id]) - inst->router_bindings.callbacks[svc_id](loopback_conn(inst), entry); - else { queue_dgram_free(entry); queue_entry_free(entry); } + router_deliver_loopback(inst, svc_id, entry); return 0; } @@ -1306,11 +1337,12 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uin if (cb) { struct ll_entry* e = queue_entry_new(0); if (e) { - e->dgram = u_malloc(9); + e->dgram = u_malloc(ROUTER_SVC_HDR_SIZE); if (e->dgram) { e->dgram[0] = svc_id; - memcpy(e->dgram + 1, &remote_node_id, 8); - e->len = 9; + memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &remote_node_id, 8); + memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8); + e->len = ROUTER_SVC_HDR_SIZE; cb(loopback_conn(inst), e); } else { queue_entry_free(e); } } diff --git a/src/routing_layer/etcp_router.h b/src/routing_layer/etcp_router.h index 56a1529c..741f7ed1 100644 --- a/src/routing_layer/etcp_router.h +++ b/src/routing_layer/etcp_router.h @@ -34,6 +34,14 @@ struct SVC_ROUTE_HDR { #define SVC_ROUTE_HDR_SIZE sizeof(struct SVC_ROUTE_HDR) // 33 #define SVC_ROUTE_MAX_BINDINGS 256 +// Единый формат доставки сервисной кодограммы (router → сервис): +// [svc_id:1][src_node_id:8][dst_node_id:8][payload...] +// src/dst — реальные end-to-end узлы из SVC_ROUTE заголовка (не промежуточные). +#define ROUTER_SVC_SRC_OFF 1 +#define ROUTER_SVC_DST_OFF 9 +#define ROUTER_SVC_PAYLOAD_OFF 17 +#define ROUTER_SVC_HDR_SIZE 17 // svc_id(1) + src_node_id(8) + dst_node_id(8) + // Биты в flags #define ROUTER_FLAG_START 0x80 #define ROUTER_FLAG_RST 0x40 diff --git a/src/routing_layer/routing.c b/src/routing_layer/routing.c index a9be50e2..6676f898 100644 --- a/src/routing_layer/routing.c +++ b/src/routing_layer/routing.c @@ -199,8 +199,19 @@ static void routing_pkt_from_etcp_cb(struct ETCP_CONN* conn, struct ll_entry* pk return; } - // Source is the remote node we received the packet from - uint64_t src_node_id = conn->peer_node_id; + // Реальный отправитель — из [svc_id][src_node_id][dst_node_id][ip] (роутер встраивает src/dst) + if (pkt->len < ROUTER_SVC_HDR_SIZE + 1) { + DEBUG_WARN(DEBUG_CATEGORY_ROUTING, "routing: short packet len=%u", pkt->len); + queue_entry_free(pkt); + queue_dgram_free(pkt); + return; + } + uint64_t src_node_id; + memcpy(&src_node_id, pkt->dgram + ROUTER_SVC_SRC_OFF, 8); + // route_pkt ожидает [cmd(1)][ip_data]: сдвигаем ip-пакет на позицию 1 (svc_id остаётся cmd-байтом) + size_t ip_len = pkt->len - ROUTER_SVC_HDR_SIZE; + memmove(pkt->dgram + 1, pkt->dgram + ROUTER_SVC_PAYLOAD_OFF, ip_len); + pkt->len = (uint16_t)(1 + ip_len); route_pkt(instance, pkt, src_node_id); } diff --git a/tests/Makefile.am b/tests/Makefile.am index a83e3a01..1b34e108 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -460,3 +460,8 @@ check-local: $(check_PROGRAMS) echo "Run 'cat $(TEST_LOG_DIR)/.log' to see details"; \ exit 1; \ fi + +# Интеграционный тест tcp_proxy_full (требует sudo, TUN, iptables). Запускается из check.sh. +.PHONY: check-proxy +check-proxy: + @cd $(srcdir)/tcp_proxy_full && bash run_test.sh diff --git a/tests/tcp_proxy_full/client.conf b/tests/tcp_proxy_full/client.conf index 5baeebf2..8c91af03 100644 --- a/tests/tcp_proxy_full/client.conf +++ b/tests/tcp_proxy_full/client.conf @@ -1,5 +1,5 @@ [global] -my_node_id=0xAAAA000000000002 +my_node_id=0x5a77374f79310824 my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01 tun_ip=10.200.20.1/24 @@ -12,14 +12,14 @@ type=public [client: c1] keepalive=1 -link=s1:127.0.0.1:15001 -peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a +link=s1:127.0.0.1:15003 +peer_public_key=23ea42d345a0efbcfaab4f47c27bf6da05e3682e27a667ba60a16645ce2c9024 [tcp_proxy_client] enabled=yes tun_name=tun_test_proxy tun_ip=10.200.30.1 -via_node=0xAAAA000000000001 +via_node=0x4239ac368c0df079 [debug] connection=info @@ -27,3 +27,5 @@ socket=info general=info traffic=info tun=info +proxy=debug +etcp_route=debug diff --git a/tests/tcp_proxy_full/exit.conf b/tests/tcp_proxy_full/exit.conf index 9b11670a..22fdd540 100644 --- a/tests/tcp_proxy_full/exit.conf +++ b/tests/tcp_proxy_full/exit.conf @@ -1,5 +1,5 @@ [global] -my_node_id=0xAAAA000000000001 +my_node_id=0x4239ac368c0df079 my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68 my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a tun_ip=10.200.10.1/24 @@ -20,3 +20,5 @@ enabled=yes connection=info socket=info traffic=info +proxy=debug +etcp_route=debug diff --git a/tests/tcp_proxy_full/intermediate.conf b/tests/tcp_proxy_full/intermediate.conf new file mode 100644 index 00000000..8d39ce64 --- /dev/null +++ b/tests/tcp_proxy_full/intermediate.conf @@ -0,0 +1,28 @@ +[global] +my_node_id=0x628471cde1291456 +my_private_key=2012fcb6f33004ee64c899a229c54619ed89f83d8fe7e003c4d2eca43c453b58 +my_public_key=23ea42d345a0efbcfaab4f47c27bf6da05e3682e27a667ba60a16645ce2c9024 +tun_ip=10.200.40.1/24 +tun_ifname=tun_test_inter +debug_level=error + +[server: s1] +addr=127.0.0.1:15003 +type=public + +[allowed_keys] +allow_all=1 + +[client: c1] +keepalive=1 +link=s1:127.0.0.1:15001 +peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a + +[debug] +connection=info +socket=info +general=info +traffic=info + etcp_route=debug + proxy=debug + diff --git a/tests/tcp_proxy_full/run_test.sh b/tests/tcp_proxy_full/run_test.sh index 2d587d36..4d38cf48 100755 --- a/tests/tcp_proxy_full/run_test.sh +++ b/tests/tcp_proxy_full/run_test.sh @@ -1,9 +1,9 @@ #!/bin/bash # run_test.sh — интеграционный тест tcp_proxy (transparent proxy) # -# Трафик: +# Топология (3 узла, трафик идёт через промежуточный узел транзитом): # tcp_client.py (SO_MARK=1, dest=10.200.100.N) → table 100 → tun_test_proxy -# → tcp_proxy(lwIP) → ETCP → exit(remote_proxy) +# → tcp_proxy(lwIP, клиент) → ETCP → intermediate (транзит) → ETCP → exit # → connect(10.200.100.N) → iptables DNAT (!mark=1) → 127.0.0.1 → echo → обратно # # Запуск: sudo ./run_test.sh @@ -41,6 +41,9 @@ setup_net() { iptables -t nat -C OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT" 2>/dev/null \ || iptables -t nat -A OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT" done + # conn_refused: .99 → порт без слушателя, exit получает ECONNREFUSED (иначе SYN уходит в tun_test_proxy и зацикливается) + iptables -t nat -C OUTPUT -d 10.200.100.99 -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$REFUSED_PORT" 2>/dev/null \ + || iptables -t nat -A OUTPUT -d 10.200.100.99 -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$REFUSED_PORT" echo " Done" } @@ -63,6 +66,7 @@ cleanup_net() { for i in $(seq 1 $N_IP); do iptables -t nat -D OUTPUT -d "10.200.100.${i}" -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$ECHO_PORT" 2>/dev/null || true done + iptables -t nat -D OUTPUT -d 10.200.100.99 -p tcp -m mark ! --mark 1 -j DNAT --to-destination 127.0.0.1:"$REFUSED_PORT" 2>/dev/null || true ip rule del fwmark 1 priority 199 table 100 2>/dev/null || true ip route del 10.200.100.0/24 dev tun_test_proxy table 100 2>/dev/null || true ip route del 10.200.100.0/24 dev tun_test_proxy 2>/dev/null || true @@ -75,18 +79,22 @@ cleanup_net() { cleanup() { echo ""; echo "=== Cleanup ===" kill $EXIT_PID 2>/dev/null || true + kill $INTER_PID 2>/dev/null || true kill $CLIENT_PID 2>/dev/null || true kill $ECHO_PID 2>/dev/null || true wait $EXIT_PID 2>/dev/null || true + wait $INTER_PID 2>/dev/null || true wait $CLIENT_PID 2>/dev/null || true wait $ECHO_PID 2>/dev/null || true sleep 0.5; cleanup_net } wait_for_etcp() { - for i in $(seq 1 30); do - grep -q "Connection established\|initialized and marked as UP (client)" "$LOG_DIR/exit_utun.log" 2>/dev/null && return 0 - grep -q "initialized and marked as UP (client)" "$LOG_DIR/client_utun.log" 2>/dev/null && return 0 + for i in $(seq 1 40); do + grep -q "Connection UP" "$LOG_DIR/exit_utun.log" 2>/dev/null \ + && grep -q "Connection UP" "$LOG_DIR/intermediate_utun.log" 2>/dev/null \ + && grep -q "Connection UP" "$LOG_DIR/client_utun.log" 2>/dev/null \ + && return 0 sleep 0.5 done return 1 @@ -202,6 +210,10 @@ echo "Starting utun exit ..." "$UTUN_BIN" -c "$SCRIPT_DIR/exit.conf" -f -l "$LOG_DIR/exit_utun.log" >"$LOG_DIR/exit_stdout.log" 2>&1 & EXIT_PID=$! +echo "Starting utun intermediate ..." +"$UTUN_BIN" -c "$SCRIPT_DIR/intermediate.conf" -f -l "$LOG_DIR/intermediate_utun.log" >"$LOG_DIR/intermediate_stdout.log" 2>&1 & +INTER_PID=$! + echo "Starting utun client ..." "$UTUN_BIN" -c "$SCRIPT_DIR/client.conf" -f -l "$LOG_DIR/client_utun.log" >"$LOG_DIR/client_stdout.log" 2>&1 & CLIENT_PID=$! @@ -216,15 +228,15 @@ echo "ETCP ready"; sleep 1 echo ""; echo "=== Running tests ===" -# 1: stress first (catch resource leaks) -run_stress && ((++PASS)) || ((++FAIL)) -sleep 2 - -# 2: basic_1mb +# 1: basic_1mb СНАЧАЛА (изолируем от стресса, чтобы понять — падает ли 1MB сам по себе) run_test basic_1mb --host 10.200.100.1 --port "$ECHO_PORT" --size 1048576 --verify --timeout 10 \ && ((++PASS)) || ((++FAIL)) sleep 1 +# 2: stress (catch resource leaks) +run_stress && ((++PASS)) || ((++FAIL)) +sleep 2 + # 3: concurrent_5 run_concurrent concurrent_5 5 204800 && ((++PASS)) || ((++FAIL)) sleep 1 diff --git a/tests/test_etcp_router.c b/tests/test_etcp_router.c index abe48b83..468768c0 100644 --- a/tests/test_etcp_router.c +++ b/tests/test_etcp_router.c @@ -248,23 +248,23 @@ static int send_one_pkt(uint32_t seq, int data_len) { // ======================== Server handler ======================== static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { - // Restart notification: conn=NULL, len=9, format [svc_id:1][node_id:8] - if (!conn && entry && entry->dgram && entry->len == 9) { + // Restart/close notification: [svc_id][src][dst] без payload + if (entry && entry->dgram && entry->len == ROUTER_SVC_HDR_SIZE) { DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "srv_handler: restart notification, resetting expected_seq from %u to 0", g_expected_seq); g_expected_seq = 0; queue_dgram_free(entry); queue_entry_free(entry); return; } - if (!entry || !entry->dgram || entry->len < 10) { + if (!entry || !entry->dgram || entry->len < ROUTER_SVC_PAYLOAD_OFF + 9) { printf("[FAIL] srv_handler: bad entry len=%zu conn=%p\n", entry ? entry->len : 0, (void*)conn); if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } g_fail = 4; g_done = -1; return; } - uint8_t subcmd = entry->dgram[1]; - uint32_t seq = 0; memcpy(&seq, entry->dgram + 2, 4); - uint32_t data_len = 0; memcpy(&data_len, entry->dgram + 6, 4); - uint8_t* payload = entry->dgram + 10; + uint8_t subcmd = entry->dgram[ROUTER_SVC_PAYLOAD_OFF]; + uint32_t seq = 0; memcpy(&seq, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 1, 4); + uint32_t data_len = 0; memcpy(&data_len, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 5, 4); + uint8_t* payload = entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 9; if (subcmd != 0x01) { printf("[FAIL] srv_handler: unexpected subcmd=%u seq=%u\n", subcmd, seq); @@ -325,9 +325,9 @@ static void srv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { static int loop_rcvd = 0, loop_ok = 0; static void loop_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; - if (entry && entry->dgram && entry->len >= 6) { + if (entry && entry->dgram && entry->len >= ROUTER_SVC_HDR_SIZE + 3) { loop_rcvd++; - if (entry->dgram[1] == 0xAA && entry->dgram[2] == 0xBB && entry->dgram[3] == 0xCC) loop_ok = 1; + if (entry->dgram[ROUTER_SVC_PAYLOAD_OFF] == 0xAA && entry->dgram[ROUTER_SVC_PAYLOAD_OFF + 1] == 0xBB && entry->dgram[ROUTER_SVC_PAYLOAD_OFF + 2] == 0xCC) loop_ok = 1; } if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } } diff --git a/tests/test_etcp_router_reconnect.c b/tests/test_etcp_router_reconnect.c index 3c4627d6..074d76bd 100644 --- a/tests/test_etcp_router_reconnect.c +++ b/tests/test_etcp_router_reconnect.c @@ -128,14 +128,14 @@ static void cleanup(void) { test_unlink(a_conf); test_unlink(b_conf); test_unlin static void recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { if (!entry || !entry->dgram) { fprintf(stderr, "RECV: null entry\n"); g_test_ok=-1; g_fail_code=10; return; } - if (entry->len == 9) { queue_dgram_free(entry); queue_entry_free(entry); return; } - if (entry->len < 7) { fprintf(stderr, "RECV: short entry len=%u\n", (unsigned)entry->len); g_test_ok=-1; g_fail_code=10; queue_dgram_free(entry); queue_entry_free(entry); return; } - uint32_t seq=0; memcpy(&seq,entry->dgram+1,4); - uint16_t sz=0; memcpy(&sz, entry->dgram+5,2); - if (entry->len != (size_t)(7+sz)) { fprintf(stderr, "FAIL: size mismatch seq=%u\n", seq); g_test_ok=-1; g_fail_code=11; queue_dgram_free(entry); queue_entry_free(entry); return; } + if (entry->len == ROUTER_SVC_HDR_SIZE) { queue_dgram_free(entry); queue_entry_free(entry); return; } + if (entry->len < ROUTER_SVC_PAYLOAD_OFF + 6) { fprintf(stderr, "RECV: short entry len=%u\n", (unsigned)entry->len); g_test_ok=-1; g_fail_code=10; queue_dgram_free(entry); queue_entry_free(entry); return; } + uint32_t seq=0; memcpy(&seq,entry->dgram+ROUTER_SVC_PAYLOAD_OFF,4); + uint16_t sz=0; memcpy(&sz, entry->dgram+ROUTER_SVC_PAYLOAD_OFF+4,2); + if (entry->len != (size_t)(ROUTER_SVC_PAYLOAD_OFF+6+sz)) { fprintf(stderr, "FAIL: size mismatch seq=%u\n", seq); g_test_ok=-1; g_fail_code=11; queue_dgram_free(entry); queue_entry_free(entry); return; } if (seq != g_expected_seq) { fprintf(stderr, "FAIL: seq broken seq=%u expected=%u\n", seq, g_expected_seq); g_test_ok=-1; g_fail_code=12; queue_dgram_free(entry); queue_entry_free(entry); return; } uint8_t exp[MAX_PAYLOAD]; gen_payload(seq, exp, sz); - if (memcmp(entry->dgram+7, exp, sz) != 0) { fprintf(stderr, "FAIL: corrupt seq=%u\n", seq); g_test_ok=-1; g_fail_code=13; queue_dgram_free(entry); queue_entry_free(entry); return; } + if (memcmp(entry->dgram+ROUTER_SVC_PAYLOAD_OFF+6, exp, sz) != 0) { fprintf(stderr, "FAIL: corrupt seq=%u\n", seq); g_test_ok=-1; g_fail_code=13; queue_dgram_free(entry); queue_entry_free(entry); return; } g_expected_seq++; g_rcvd_total++; queue_dgram_free(entry); queue_entry_free(entry); } diff --git a/tests/test_etcp_router_unit.c b/tests/test_etcp_router_unit.c index 0ab2f029..52b2b4f9 100644 --- a/tests/test_etcp_router_unit.c +++ b/tests/test_etcp_router_unit.c @@ -38,18 +38,18 @@ static int g_notify_count = 0; static void test_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; - if (!entry || !entry->dgram || entry->len < 2) { + if (!entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) { g_notify_count++; if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - // close/restart: len==9 (svc_id + node_id), conn==NULL - if (entry->len == 9) { + // close/restart: [svc_id][src][dst] без payload + if (entry->len == ROUTER_SVC_HDR_SIZE) { g_notify_count++; queue_dgram_free(entry); queue_entry_free(entry); return; } - if (entry->dgram[1] != rx.marker) { rx.errors++; } else { rx.delivered++; rx.last_seq = rx.expected_seq; rx.expected_seq++; } + if (entry->dgram[ROUTER_SVC_PAYLOAD_OFF] != rx.marker) { rx.errors++; } else { rx.delivered++; rx.last_seq = rx.expected_seq; rx.expected_seq++; } queue_dgram_free(entry); queue_entry_free(entry); } diff --git a/tests/test_icmp_proxy.c b/tests/test_icmp_proxy.c index 149019ff..9af5c62e 100644 --- a/tests/test_icmp_proxy.c +++ b/tests/test_icmp_proxy.c @@ -82,16 +82,16 @@ static void* g_to_id = NULL; static void cli_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; - if (!entry || !entry->dgram || entry->len < ICMP_PROXY_HDR_SIZE + 1) { + if (!entry || !entry->dgram || entry->len < ICMP_PROXY_RECV_HDR_SIZE + 1) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - if (entry->dgram[1] == ICMP_PROXY_SUBCMD_REPLY) { + if (entry->dgram[ROUTER_SVC_PAYLOAD_OFF] == ICMP_PROXY_SUBCMD_REPLY) { uint16_t rid, rseq; - memcpy(&rid, entry->dgram + 18, 2); - memcpy(&rseq, entry->dgram + 20, 2); + memcpy(&rid, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 9, 2); + memcpy(&rseq, entry->dgram + ROUTER_SVC_PAYLOAD_OFF + 11, 2); if (rid == test_echo_id && rseq == test_echo_seq) { - size_t payload_len = entry->len - ICMP_PROXY_HDR_SIZE; - uint8_t* payload = entry->dgram + ICMP_PROXY_HDR_SIZE; + size_t payload_len = entry->len - ICMP_PROXY_RECV_HDR_SIZE; + uint8_t* payload = entry->dgram + ICMP_PROXY_RECV_HDR_SIZE; if (payload_len == ICMP_PAYLOAD_SIZE && memcmp(payload, send_buf, ICMP_PAYLOAD_SIZE) == 0) { reply_rcvd = 1; g_done = 1; g_ok = 1; } else { diff --git a/tests/test_nat_transport.c b/tests/test_nat_transport.c index a38c1b58..938477bf 100644 --- a/tests/test_nat_transport.c +++ b/tests/test_nat_transport.c @@ -35,7 +35,7 @@ #define TEST_TIMEOUT_MS 15000 static uint64_t g_provider_node_id = 0; static uint64_t g_client_node_id = 0; -#define NAT_SVC_HDR_SIZE 9 // svc_id(1) + src_node_id(8) +#define NAT_SVC_HDR_SIZE 1 // send: svc_id(1) + ip_data (src/dst добавляет роутер) static struct UTUN_INSTANCE* inst_provider = NULL; static struct UTUN_INSTANCE* inst_client = NULL; @@ -293,14 +293,13 @@ static int test_provider_egress(void) { } DEBUG_INFO(DEBUG_CATEGORY_NAT, "Client conn peer_node_id=0x%016llx", (unsigned long long)client_conn->peer_node_id); - // Build ETCP_RT_ID_NAT packet via etcp_route_send (new format: svc_id + src_node_id + ip_data) + // Build ETCP_RT_ID_NAT packet via etcp_route_send (src/dst node_id добавляет роутер) size_t ip_len; uint8_t* raw_ip = build_udp_pkt(0x0A0000FE, 40000, 0x08080808, 53, &ip_len); size_t total = NAT_SVC_HDR_SIZE + ip_len; uint8_t* dgram = u_malloc(total); dgram[0] = ETCP_RT_ID_NAT; - memcpy(dgram + 1, &inst_client->nat_tr.self_node_id, 8); - memcpy(dgram + 9, raw_ip, ip_len); + memcpy(dgram + 1, raw_ip, ip_len); free(raw_ip); struct ll_entry* entry = queue_entry_new(0); @@ -469,8 +468,7 @@ static int test_full_roundtrip(void) { size_t total = NAT_SVC_HDR_SIZE + ip_len; uint8_t* dgram = u_malloc(total); dgram[0] = ETCP_RT_ID_NAT; - memcpy(dgram + 1, &inst_client->nat_tr.self_node_id, 8); - memcpy(dgram + 9, raw_ip, ip_len); + memcpy(dgram + 1, raw_ip, ip_len); free(raw_ip); struct ll_entry* entry = queue_entry_new(0); diff --git a/tests/test_tcp_io.c b/tests/test_tcp_io.c index 7fc04318..871d4663 100644 --- a/tests/test_tcp_io.c +++ b/tests/test_tcp_io.c @@ -307,6 +307,52 @@ static void test_error_callback(void) { TEST_PASS(); } +static void test_connect_timeout(void) { + TEST_START("Connect timeout on black hole (200ms)"); +#ifdef _WIN32 + TEST_PASS(); + return; +#else + uasync_t* ua = uasync_create(); + ASSERT_TRUE(ua != NULL, "uasync_create failed"); + reset_counters(); + + int fd = socket(AF_INET, SOCK_STREAM, 0); + ASSERT_TRUE(fd >= 0, "socket failed"); + fcntl(fd, F_SETFL, fcntl(fd, F_GETFL, 0) | O_NONBLOCK); + + struct sockaddr_in addr; + memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_port = htons(81); + addr.sin_addr.s_addr = htonl(0xC0000201); // 192.0.2.1 (TEST-NET-1, RFC 5737) — чёрная дыра + + int ret = connect(fd, (struct sockaddr*)&addr, sizeof(addr)); + ASSERT_TRUE(ret < 0 && errno == EINPROGRESS, "connect should return EINPROGRESS"); + + struct tcp_conn* tc = tcp_conn_create(ua, fd, 1500, 8192, 32, 8, 0, on_fin, on_error, NULL); + ASSERT_TRUE(tc != NULL, "tcp_conn_create failed"); + ASSERT_EQ(tc->connected, 0, "should not be connected"); + + tcp_conn_set_connect_timeout(tc, 200); + + // Крутим цикл (1 мс × 500) пока таймаут не сработает + int iterations = 0; + while (g_error_count == 0 && iterations < 500) { + uasync_poll(ua, 10); + iterations++; + } + ASSERT_EQ(g_error_count, 1, "on_error not called on connect timeout"); + ASSERT_EQ(g_last_error, ETIMEDOUT, "on_error should report ETIMEDOUT"); + ASSERT_EQ(tc->connected, 0, "connection should still be not connected"); + ASSERT_EQ(tc->error, 1, "tc->error should be set"); + + tcp_conn_destroy(tc); + uasync_destroy(ua, 0); + TEST_PASS(); +#endif +} + // ==================================================================== // FIN / Close через очередь — новые тесты // ==================================================================== @@ -552,6 +598,7 @@ int main(void) { test_high_water_pause(); test_connect_detection(); test_error_callback(); + test_connect_timeout(); test_push_fin(); test_push_close(); test_fin_data_ordering(); diff --git a/tests/test_udp_proxy.c b/tests/test_udp_proxy.c index d3b6ded7..603efea9 100644 --- a/tests/test_udp_proxy.c +++ b/tests/test_udp_proxy.c @@ -92,13 +92,12 @@ static void udp_echo_cb(socket_t sock, void* arg) { static void cli_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; - if (!entry || !entry->dgram || entry->len < UDP_PROXY_HDR_SIZE + 1) { + if (!entry || !entry->dgram || entry->len < UDP_PROXY_RECV_HDR_SIZE + 1) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } - // svc_id(1) + subcmd(1) + sender_node_id(8) + src_ip(4) + src_port(2) + dst_ip(4) + dst_port(2) + payload - size_t payload_len = entry->len - UDP_PROXY_HDR_SIZE; - if (entry->dgram[1] == UDP_PROXY_SUBCMD_DATA) { - uint8_t* payload = entry->dgram + UDP_PROXY_HDR_SIZE; + size_t payload_len = entry->len - UDP_PROXY_RECV_HDR_SIZE; + if (entry->dgram[ROUTER_SVC_PAYLOAD_OFF] == UDP_PROXY_SUBCMD_DATA) { + uint8_t* payload = entry->dgram + UDP_PROXY_RECV_HDR_SIZE; if (payload_len == PAYLOAD_SIZE && memcmp(payload, send_buf, PAYLOAD_SIZE) == 0) { reply_rcvd = 1; g_done = 1; g_ok = 1; } else { diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index c7148b37..ed4bfe3b 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -250,9 +250,10 @@ QList UtunNode::getInviteAddresses(DbManager* db) { } void UtunNode::recvCallback(struct ETCP_CONN* conn, struct ll_entry* entry) { - if (!g_currentNode || !entry) return; - uint64_t src = conn ? conn->peer_node_id : 0; - QByteArray data((const char*)entry->data, (int)(entry->len ? entry->len : 0)); + if (!g_currentNode || !entry || !entry->dgram || entry->len < ROUTER_SVC_HDR_SIZE) return; + uint64_t src; + memcpy(&src, entry->dgram + ROUTER_SVC_SRC_OFF, 8); + QByteArray data((const char*)(entry->dgram + ROUTER_SVC_PAYLOAD_OFF), (int)(entry->len - ROUTER_SVC_PAYLOAD_OFF)); QMetaObject::invokeMethod(g_currentNode, [=] { emit g_currentNode->messageReceived(src, data); }, Qt::QueuedConnection);