Browse Source

tcp_proxy: fix throughput (tx_queue backpressure) + connect timeout + integration test

- 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)
v2
evgeny 4 weeks ago
parent
commit
f47ad7d87f
  1. 3
      AGENTS.md
  2. 11
      check.sh
  3. 115
      doc/TASK_tcp_proxy_throughput.md
  4. 66
      doc/lwip_nuances.md
  5. 27
      lib/tcp_io.c
  6. 9
      lib/tcp_io.h
  7. 9
      src/config_parser.c
  8. 1
      src/config_parser.h
  9. 12
      src/media_delivery/media_delivery.c
  10. 20
      src/nat_transport.c
  11. 59
      src/proxy/icmp_proxy.c
  12. 10
      src/proxy/icmp_proxy.h
  13. 244
      src/proxy/tcp_proxy_client.c
  14. 10
      src/proxy/tcp_proxy_client.h
  15. 55
      src/proxy/tcp_proxy_server.c
  16. 4
      src/proxy/tcp_proxy_server.h
  17. 1
      src/proxy/tcp_proxy_server_doc.md
  18. 56
      src/proxy/udp_proxy.c
  19. 10
      src/proxy/udp_proxy.h
  20. 4
      src/routing_layer/conn_mgr_core.c
  21. 70
      src/routing_layer/etcp_router.c
  22. 8
      src/routing_layer/etcp_router.h
  23. 15
      src/routing_layer/routing.c
  24. 5
      tests/Makefile.am
  25. 10
      tests/tcp_proxy_full/client.conf
  26. 4
      tests/tcp_proxy_full/exit.conf
  27. 28
      tests/tcp_proxy_full/intermediate.conf
  28. 32
      tests/tcp_proxy_full/run_test.sh
  29. 18
      tests/test_etcp_router.c
  30. 12
      tests/test_etcp_router_reconnect.c
  31. 8
      tests/test_etcp_router_unit.c
  32. 12
      tests/test_icmp_proxy.c
  33. 10
      tests/test_nat_transport.c
  34. 47
      tests/test_tcp_io.c
  35. 9
      tests/test_udp_proxy.c
  36. 7
      tools/chatgui/transport/utun_node.cpp

3
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

11
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

115
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
```

66
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`.

27
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");
}
}

9
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
}

9
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;

1
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;

12
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);

20
src/nat_transport.c

@ -12,7 +12,7 @@
#include "../lib/ll_queue.h"
#include <string.h>
#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);
}
}

59
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);

10
src/proxy/icmp_proxy.h

@ -9,6 +9,7 @@ extern "C" {
#include <stdint.h>
#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 {

244
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;

10
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;

55
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;
}

4
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;

1
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

56
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);

10
src/proxy/udp_proxy.h

@ -9,6 +9,7 @@ extern "C" {
#include <stdint.h>
#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 {

4
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) {

70
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); }
}

8
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

15
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);
}

5
tests/Makefile.am

@ -460,3 +460,8 @@ check-local: $(check_PROGRAMS)
echo "Run 'cat $(TEST_LOG_DIR)/<test>.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

10
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

4
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

28
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

32
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

18
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); }
}

12
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);
}

8
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);
}

12
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 {

10
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);

47
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();

9
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 {

7
tools/chatgui/transport/utun_node.cpp

@ -250,9 +250,10 @@ QList<NodeAddr> 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);

Loading…
Cancel
Save