From fbea9c3869b479cda1f38bccf311081e318029e4 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 3 Sep 2026 18:13:00 +0300 Subject: [PATCH] socks_proxy: fix client stall + non-blocking DNS via async_dns - fix HTTP_STATE_RELAY == SOCKS_STATE_CONNECTING enum collision (relay data dropped) - fix wrong waiter cancel (etcp_router_waiter_cancel -> etcp_router_cancel_send_ready) - add backpressure retry timer, restore lost queue_resume_callback, defer FIN on pending tx_buf - async_dns: /etc/hosts fast-path + hosts loading (avoid 4s DNS timeout for localhost) --- doc/TASK_tcp_proxy_throughput.md | 115 ------------ doc/_plan.txt | 35 ---- doc/db_sync_v2.md | 292 ------------------------------- doc/etcp_config.txt | 10 -- lib/async_dns.c | 99 +++++++++++ src/proxy/socks_proxy.c | 253 ++++++++++++++++++-------- src/proxy/socks_proxy.h | 22 ++- 7 files changed, 298 insertions(+), 528 deletions(-) delete mode 100644 doc/TASK_tcp_proxy_throughput.md delete mode 100644 doc/_plan.txt delete mode 100644 doc/db_sync_v2.md delete mode 100644 doc/etcp_config.txt diff --git a/doc/TASK_tcp_proxy_throughput.md b/doc/TASK_tcp_proxy_throughput.md deleted file mode 100644 index e9bd377d..00000000 --- a/doc/TASK_tcp_proxy_throughput.md +++ /dev/null @@ -1,115 +0,0 @@ -# Задача: продолжить исправление 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/_plan.txt b/doc/_plan.txt deleted file mode 100644 index d5c4a24f..00000000 --- a/doc/_plan.txt +++ /dev/null @@ -1,35 +0,0 @@ -План доработок -[ ] продумать архитектуру репликации таблицы узлов группы -[ ] добавить узлы хелперы (для проксирования трафика) - - -Категории узлов: -- узлы без нат [] -- узлы с nat [restricted] - -Узлы за nat сами подключаются к прямым - -Узлы за eim nat сами подключаются и сообщают обновление адреса. - -Прямые узлы должны сами должны объединиться в сеть следующим образом: - -Инит: база содержит RTT по узлам. -выбираем n=16 узлов с минимальным rtt. запускаем подключение к ним. -синхронизируем роутинг таблицу с ними (обмен маршрутами должен инициироваться сразу после подключения). -при потере связи узлы через nat должны найти альтернативы. - -в таблице обновляется - -как обмен завершен - -Нужно собирать и обновлять информацию о том -- к каким публичным нодам подключены nat узлы. - -- метрика доступности узла - как давно не менялся ip:port - собираем с клиентов. вместе с рейтигном (каждый клиент по узлу фолрмирует свой рейтинг кандидата в суперузлы и отправляет подписанный рейтинг узлу) -рейтинг включает: -- число разных часов фигурирует в замерах -- число разных дат фигурирует в замерах -- гистограмма rtt: 0-2ms, 2-4ms, 4-6ms, 6-10ms, 10-13ms, 14-18ms, 18-25ms, 25-35ms, 35-50ms, 50-75ms, 75-100ms, 100-130ms, 130-160ms, 160+ms -- гистограмма потерь: 0-1% 1-2% 2-4% 4-8% >8% (обновляется каждый час) - -Сделать подключение к группе diff --git a/doc/db_sync_v2.md b/doc/db_sync_v2.md deleted file mode 100644 index 7b488422..00000000 --- a/doc/db_sync_v2.md +++ /dev/null @@ -1,292 +0,0 @@ -# DB Sync Protocol v2 — Техническое задание - -## 1. Проблема - -chain_hash = SHA256(prev_chain_hash || id || ts || datahash) — хеш-цепочка, где хеш каждой записи зависит от предыдущей. - -`db_record_insert` (db_sync.c:341-463) при вставке записи в середину сортированного порядка делает каскадный пересчёт chain_hash всех последующих записей — O(N). При двустороннем обмене PUSH'ами получается O(N²). - -## 2. Решение - -**Разделить два концерна:** -- **Сравнение при синхронизации** — перейти на `datahash` (первые 8 байт SHA256(data), позиционно-независимый). Не требует каскада. -- **Целостность цепочки** — `chain_hash` остаётся в схеме БД, но используется только для однократной стартап-проверки и диагностики. - -**Синхронизация:** per-peer позиционная (`synced_pos`). При PUSH каскад не делается. При SEND_DATA батч-вставка + один каскад после батча. - -## 3. Модель данных - -### 3.1. SI_PEER - -```c -struct SI_PEER { - uint64_t node_id; - uint32_t synced_pos; // последняя подтверждённо общая позиция (0-based index) - uint8_t sync_state; // 0=not_synced, 1=syncing, 2=synced -}; -``` - -- `synced_pos` = индекс последней записи в глобальном порядке `ORDER BY timestamp, datahash`, для которой chain_hash (а после v2 — datahash) подтверждённо совпадает у обоих пиров. -- При PUSH-вставке записи на позицию P, для всех peer'ов у которых `synced_pos >= P`: `synced_pos = P - 1`. -- При SEND_DATA: `synced_pos = from + received_count - 1`. - -### 3.2. Схема БД (без изменений) - -```sql -CREATE TABLE "db_sync__" ( - timestamp INTEGER NOT NULL, - datahash INTEGER NOT NULL, - id INTEGER NOT NULL, - chain_hash BLOB NOT NULL, - author INTEGER NOT NULL, - flags INTEGER NOT NULL DEFAULT 0, - data BLOB, - author_signature BLOB, - delivered_peers INTEGER NOT NULL DEFAULT 0, - delivery_chain TEXT NOT NULL DEFAULT '', - PRIMARY KEY (timestamp, datahash) -); -``` - -## 4. Wire-формат - -### 4.1. Сообщения, которые меняются - -#### INIT_SYNC - -``` -Было: [type:1][count:4][last_chain_hash:32] = 37 байт -Стало: [type:1][count:4] = 5 байт -``` - -#### INIT_RESP - -``` -Было: [type:1][tp:4][chain_at_tp:32][sc:1][(pos:4,chain_hash:32)*sc] -Стало: [type:1][tp:4][dh_at_tp:8] [sc:1][(pos:4,dh:8)*sc] - - dh_at_tp — datahash на позиции tp (8 байт вместо 32) - sparse — datahash вместо chain_hash (8 байт на элемент вместо 32) - - При sc=0 (exact match): длина 1+4+8+1 = 14 байт (было 37) - При sc=16 sparse: длина 14 + 16*12 = 206 байт (было 37 + 16*36 = 613) -``` - -#### SYNC_DONE - -``` -Было: [type:1][count:4][last_chain_hash:32] = 37 байт -Стало: [type:1][count:4][last_dh:8] = 13 байт -``` - -### 4.2. Сообщения, которые НЕ меняются - -- **PUSH**: `[type:1][id:8][ts:8][dh:8][dlen:4][data][sig_len:1][sig]` — без изменений. `prev_ch` НЕ добавляется. -- **ACK_PUSH**: `[type:1][dh:8][ts:8]` — без изменений. -- **REFINE**: `[type:1][from:4][to:4][hc:1][(pos:4,dh:8)*hc]` — уже использует datahash, без изменений. -- **SEND_DATA**: `[type:1][from:4][count:2][(id:8,ts:8,dh:8,dlen:4,data,sig_len:1,sig)*count]` — без изменений. -- **ERROR**: `[type:1][code:1]` — без изменений. - -## 5. Новые/изменённые функции - -### 5.1. `db_cascade_from` - -```c -static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos); -``` - -Выполняет каскадный пересчёт `chain_hash` для записей с позиции `from_pos` до конца таблицы. - -Алгоритм: -1. BEGIN IMMEDIATE -2. SELECT chain_hash записи на позиции `from_pos - 1` (или zero если from_pos = 0) как prev_ch -3. SELECT id, timestamp, datahash от позиции from_pos до конца -4. Для каждой: chain_hash = SHA256(prev_ch || id || ts || dh), UPDATE в БД, prev_ch = новый chain_hash -5. COMMIT - -Выделяется из текущего кода `db_record_insert` (строки 412-451) в отдельную функцию. - -### 5.2. `db_record_insert` — новый параметр - -```c -static int db_record_insert(struct DB_SYNC_INSTANCE* si, - uint64_t id, uint64_t ts, uint64_t dh, - const char* json, size_t jlen, - const uint8_t* sig, size_t sig_len, - int do_cascade); -``` - -- `do_cascade=1` — после INSERT выполняется cascade (строки 412-451). Используется для локальных вставок. -- `do_cascade=0` — без cascade. Используется для PUSH и SEND_DATA (каскад делается отдельно). - -### 5.3. `db_verify_chain` (стартап-проверка) - -```c -static void db_verify_chain(struct DB_SYNC_INSTANCE* si); -``` - -Вызывается один раз после `db_sync_instance_add` (перед инициацией sync с пирами). - -Алгоритм: -1. Проход записей 0..N-1 в порядке `ORDER BY timestamp, datahash` -2. Для каждой: chain_hash = SHA256(prev_ch || id || ts || dh) -3. Сравнить с хранимым в БД -4. При первом расхождении на позиции P: пересчитать цепочку от P до конца через `db_cascade_from(si, P)`. Завершить. - -Сложность: O(N) SHA256, однократно при старте. - -### 5.4. `si_find_pos` (новая) - -```c -static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, - uint64_t ts, uint64_t dh); -``` - -Возвращает позицию (0-based index) записи с ключом `(ts, dh)` в глобальном порядке. Используется после вставки PUSH для корректировки `synced_pos`. - -### 5.5. `db_datahash_at` — уже существует (стр. 295) - -Возвращает `datahash` на заданной позиции. Используется для сравнения в INIT_SYNC/INIT_RESP/SYNC_DONE. - -## 6. Обработчики сообщений - -### 6.1. PUSH (`db_handle_push`) - -``` -Было: db_record_insert(si, ..., /* cascade встроен */) -Стало: db_record_insert(si, ..., /* do_cascade= */ 0) - uint32_t pos = si_find_pos(si, ts, dh); - for each peer: if peer->synced_pos >= pos: peer->synced_pos = pos - 1; -``` - -### 6.2. INIT_SYNC (`db_handle_init_sync`) - -``` -Изменения: - - Принимает [count:4] вместо [count:4][chain_hash:32] (проверка len >= 5 вместо 36) - - Сравнение: datahash вместо chain_hash (8 байт вместо 32) - - sparse: datahash вместо chain_hash (8 байт на элемент вместо 32) - - Размер sparse-элемента: 4(pos) + 8(dh) = 12 байт (было 4+32=36) - - resp буфер: 4096 → достаточно (14 + 16*12 = 206 байт при max sparse) -``` - -### 6.3. INIT_RESP (`db_handle_init_resp`) - -``` -Изменения: - - Принимает [tp:4][dh:8][sc:1][...] вместо [tp:4][ch:32][sc:1][...] - - Проверка len >= 13 вместо 37 - - Сравнение: datahash вместо chain_hash - - Пустая БД пира: zero-хеш 8 байт вместо 32 - - Sparse элементы: 12 байт вместо 36 - - При совпадении: synced_pos = tp (в дополнение к sync_state = 2) -``` - -### 6.4. SEND_DATA (`db_handle_send_data`) - -``` -Стало: - struct SI_PEER* sp = si_peer_find(si, src); - uint32_t fix_from = sp ? sp->synced_pos : 0; - - for each record: - db_record_insert(si, ..., /* do_cascade= */ 0); - - db_cascade_from(si, fix_from); - - sp->synced_pos = from + received_count - 1; - sp->sync_state = 2; -``` - -Каскад от `fix_from` (позиция, которая была synced ДО этого батча) гарантирует, что все chain_hash после этой точки пересчитаны — независимо от того, какие PUSH'и испортили их между синхронизациями. - -### 6.5. SYNC_DONE (`db_handle_sync_done`) - -``` -Изменения: - - Принимает [count:4][last_dh:8] вместо [count:4][last_chain_hash:32] (len >= 13 вместо 36) - - Сравнение последнего datahash вместо chain_hash -``` - -## 7. Механика многопировой синхронизации - -При наличии нескольких пиров с одинаковым `sync_state` (например, syncing=1) выбирается пир с минимальным `synced_pos` — у него самый старый общий префикс. После завершения его синхронизации выбор повторяется. - -Это гарантирует: двигаемся от самого старого несинхронизированного участка, не прыгая. Новые записи (в хвосте) синхронизируются последними. - -## 8. Полный сценарий: три участника A, B, C - -``` -Начальное состояние: все пусты. - -=== Шаг 1. A создаёт R1 (ts=1000) === -A: PUSH R1 → B, C. Все вставляют в конец, cascade 0. -synced_pos = 0 у всех. - -=== Шаг 2. B создаёт R2 (ts=2000), C создаёт R3 (ts=1500) === -B: [R1(1000), R2(2000)], PUSH R2 → A, C. - A: R2 в конец. C: R2 в конец. -C: [R1(1000), R2(2000)]? Нет — C создал R3(1500). - C: [R1(1000), R3(1500)] локально, потом PUSH R2 → R2.ts=2000 → в конец. - C: [R1(1000), R3(1500), R2(2000)] -C: PUSH R3 → A, B. - A: [R1(1000), R2(2000)] + R3(1500) → между R1 и R2: - A: [R1, R3, R2] ← R3 в середину, cascade ПРОПУЩЕН - B: [R1, R3, R2] ← аналогично - - synced_pos после PUSH R3 на A относительно B: - A вставил R3 на поз.1 → synced_pos_B: было 1 → стало 0 - (A считает, что совпадает с B только на позиции 0) - -=== Шаг 3. A инициирует sync с B (чемпион: min synced_pos) === -A→B: INIT_SYNC [count=3] -B→A: INIT_RESP [tp=min(3,3)-1=2, dh_at_tp=H2, sc=2, sparse:{(1,H3),(0,H1)}] -A проверяет: A@2=R2→dh=H2✓, A@1=R3→dh=H3✓, A@0=R1→dh=H1✓ -→ sync complete, synced_pos_B=2, sync_state=2 - -=== Шаг 4. Рестарт A === -db_verify_chain: - pos 0: R1.ch = SHA256(0||R1) ✓ - pos 1: R3.ch = SHA256(R1.ch||R3) ✓ (посчитан при PUSH без cascade — корректен) - pos 2: R2.ch = SHA256(R3.ch||R2) ✗ (STALE: был посчитан как SHA256(R1.ch||R2) до PUSH R3) - → расхождение на pos=2 → db_cascade_from(2): пересчёт R2.ch = SHA256(R3.ch||R2) ✓ - -=== Шаг 5. A, B, C активно обмениваются === -Каждый PUSH: вставка без cascade. -Периодические sync (по min synced_pos): SEND_DATA + один cascade от synced_pos. -synced_pos растёт. -``` - -## 9. Что НЕ трогать - -- `db_sync_insert_signed` — локальная вставка. ts всегда монотонный, запись в конец, `do_cascade=1`, каскад всегда 0 строк. -- TTL cleanup (`db_sync_instance_ttl_cb`) — удаление старых записей. Нужен cascade после удаления (уже существующая логика, не меняется). -- ACK_PUSH, ERROR, REFINE — без изменений. -- Схема БД — без изменений. -- `etcp_bind`/эпилог — без изменений. - -## 10. Сложность операций (итого) - -| Операция | Каскад | Сложность | -|----------|--------|-----------| -| Локальная вставка | 0 строк (конец) | O(1) | -| PUSH (приём) | нет | O(1) | -| SEND_DATA (батч N записей) | 1 каскад после батча | O(N + tail) | -| Стартап-проверка | 1 раз | O(total) | -| Sync-протокол (INIT_SYNC/INIT_RESP/REFINE) | нет cascade | O(log N) сообщений | - -## 11. Порядок реализации - -1. **Выделить `db_cascade_from(si, from_pos)`** из тела `db_record_insert`. -2. **Добавить параметр `do_cascade`** в `db_record_insert`. При `0` — только INSERT, cascade вызывается отдельно. -3. **Добавить `synced_pos`** в `struct SI_PEER`. Инициализировать в 0. -4. **Добавить `si_find_pos(si, ts, dh)`** — поиск позиции записи в глобальном порядке. -5. **Изменить wire-формат** INIT_SYNC, INIT_RESP, SYNC_DONE (chain_hash → datahash, 32→8 байт). -6. **Изменить обработчики**: - - `db_handle_init_sync`: новые размеры, datahash вместо chain_hash - - `db_handle_init_resp`: новые размеры, datahash вместо chain_hash, обновлять synced_pos - - `db_handle_send_data`: `do_cascade=0`, `db_cascade_from(fix_from)`, обновить synced_pos - - `db_handle_push`: `do_cascade=0`, скорректировать все synced_pos - - `db_handle_sync_done`: обновить synced_pos -7. **Добавить `db_verify_chain(si)`** — вызвать при старте после `db_sync_instance_add`. -8. **Многопировая синхронизация**: в `db_sync_peer_check_cb` выбирать пира с min synced_pos вместо первого попавшегося. diff --git a/doc/etcp_config.txt b/doc/etcp_config.txt deleted file mode 100644 index 3a06b151..00000000 --- a/doc/etcp_config.txt +++ /dev/null @@ -1,10 +0,0 @@ -В конфиге явно прописываются сокеты с ip и портом. -Через эти сокеты идёт весь трафик, случайные порты клиент не использует - только эти сокеты. -Соответственно, возможен только один линк между двумя пирами по одному каналу. - -виртуальные сокеты - плохо т.к. используются ip/port для поиск подключения - - - -./configure --with-openssl - openssl -./configure --without-openssl - tinycrypt diff --git a/lib/async_dns.c b/lib/async_dns.c index 1eef6bce..e86255d8 100644 --- a/lib/async_dns.c +++ b/lib/async_dns.c @@ -20,6 +20,7 @@ #include #include #include +#include #ifdef _WIN32 #include @@ -37,6 +38,9 @@ struct adns_query { void* timer; /* handle таймаута (NULL когда нет) */ int finished; char* name; + int fast; /* 1 = готовый /etc/hosts-ответ, без libdns */ + struct adns_result fast_result; + void* call_soon_id; /* отложенный вызов для fast-пути */ }; static void adns_drive(struct adns_query* q); @@ -44,6 +48,8 @@ static void adns_finish(struct adns_query* q, int status, struct sockaddr_in* ad static void adns_read_cb(socket_t sock, void* arg); static void adns_write_cb(socket_t sock, void* arg); static void adns_timer_cb(void* arg); +static void adns_fast_ready_cb(void* arg); +static int adns_hosts_lookup(const char* name, struct sockaddr_in* out); /* ─── system DNS servers discovery ─── */ @@ -124,6 +130,49 @@ int adns_system_servers(struct sockaddr_in* out, int max) { } #endif +/* ─── синхронный /etc/hosts-лукап (быстрый путь, без libdns) ─── */ + +// Ищет name в /etc/hosts. Возвращает 1 если найден IPv4-адрес (out заполнен). +// Блокирующее чтение небольшого файла — только для локальных имён (localhost и т.п.). +static int adns_hosts_lookup(const char* name, struct sockaddr_in* out) { +#ifndef _WIN32 + if (!name || !out) return 0; + FILE* f = fopen("/etc/hosts", "r"); + if (!f) return 0; + char line[512]; + int found = 0; + while (!found && fgets(line, sizeof(line), f)) { + char ip[64]; int n = 0; + const char* p = line; + while (*p && !isspace(*p) && n < 63) ip[n++] = *p++; + ip[n] = '\0'; + if (!ip[0] || ip[0] == '#') continue; + struct in_addr a; + if (inet_pton(AF_INET, ip, &a) != 1) continue; + while (*p && isspace(*p)) p++; + while (*p && *p != '\n') { + if (*p == '#') break; + char h[256]; int hn = 0; + while (*p && !isspace(*p) && *p != '#' && hn < 255) h[hn++] = *p++; + h[hn] = '\0'; + if (h[0] && strcasecmp(h, name) == 0) { + memset(out, 0, sizeof(*out)); + out->sin_family = AF_INET; + out->sin_addr = a; + found = 1; + break; + } + while (*p && isspace(*p)) p++; + } + } + fclose(f); + return found; +#else + (void)name; (void)out; + return 0; +#endif +} + /* ─── resolver setup ─── */ static struct dns_resolver* adns_open_resolver(const struct sockaddr_in* servers, int count, @@ -155,6 +204,11 @@ static struct dns_resolver* adns_open_resolver(const struct sockaddr_in* servers dns_resconf_close(rc); return NULL; } +#ifndef _WIN32 + // Локальные имена из /etc/hosts (localhost и т.п.). Не фатально, если файла нет. + if (dns_hosts_loadpath(hosts, "/etc/hosts") != 0) + DEBUG_DEBUG(ADNS_DEBUG_CAT, "adns: /etc/hosts not loaded"); +#endif struct dns_hints* hints = dns_hints_local(rc, &error); if (!hints) { @@ -326,6 +380,23 @@ static void adns_finish(struct adns_query* q, int status, struct sockaddr_in* ad u_free(name); } +// Отложенный вызов коллбэка для fast-пути (/etc/hosts): сохраняет async-контракт +// (adns_resolve возвращает handle, коллбэк вызывается в потоке event loop). +static void adns_fast_ready_cb(void* arg) { + struct adns_query* q = (struct adns_query*)arg; + q->call_soon_id = NULL; + if (q->finished) return; // отменён до срабатывания + q->finished = 1; + adns_done_cb cb = q->cb; + void* a = q->arg; + struct adns_result res = q->fast_result; + char* nm = q->name; + q->name = NULL; + u_free(q); + if (cb) cb(&res, a); + u_free(nm); +} + /* ─── public API ─── */ struct adns_query* adns_resolve(struct UASYNC* ua, const char* name, @@ -334,6 +405,28 @@ struct adns_query* adns_resolve(struct UASYNC* ua, const char* name, if (!ua || !name || !name[0] || !cb) return NULL; if (strlen(name) > 255) { DEBUG_ERROR(ADNS_DEBUG_CAT, "adns: name too long"); return NULL; } + // Быстрый путь: имя из /etc/hosts (localhost и т.п.) — без libdns и сети. + // libdns по умолчанию идёт в DNS первым (lookup="bf"), поэтому локальные имена + // ждали бы таймаут DNS-запроса (2×2s); обрабатываем их сразу. + { + struct sockaddr_in ha; + if (adns_hosts_lookup(name, &ha)) { + struct adns_query* q = u_calloc(1, sizeof(*q)); + if (!q) { DEBUG_ERROR(ADNS_DEBUG_CAT, "adns: calloc failed"); return NULL; } + q->ua = ua; q->cb = cb; q->arg = arg; + q->name = u_strdup(name); + if (!q->name) { u_free(q); DEBUG_ERROR(ADNS_DEBUG_CAT, "adns: strdup failed"); return NULL; } + q->fast = 1; + q->fast_result.status = ADNS_OK; + q->fast_result.count = 1; + q->fast_result.addrs[0] = ha; + q->call_soon_id = uasync_call_soon(ua, q, adns_fast_ready_cb); + if (!q->call_soon_id) { u_free(q->name); u_free(q); return NULL; } + DEBUG_DEBUG(ADNS_DEBUG_CAT, "adns: '%s' resolved from /etc/hosts (fast)", name); + return q; + } + } + struct sockaddr_in servers[3]; int count; if (opts && opts->server.sin_family != 0) { @@ -380,6 +473,12 @@ struct adns_query* adns_resolve(struct UASYNC* ua, const char* name, void adns_cancel(struct adns_query* q) { if (!q || q->finished) return; q->finished = 1; + if (q->fast) { + if (q->call_soon_id) { uasync_call_soon_cancel(q->ua, q->call_soon_id); q->call_soon_id = NULL; } + u_free(q->name); + u_free(q); + return; + } adns_teardown(q); u_free(q->name); u_free(q); diff --git a/src/proxy/socks_proxy.c b/src/proxy/socks_proxy.c index 01f170e7..736d6b31 100644 --- a/src/proxy/socks_proxy.c +++ b/src/proxy/socks_proxy.c @@ -12,6 +12,7 @@ #include "../lib/memory_pool.h" #include "../lib/tcp_io.h" #include "../lib/mem.h" +#include "../lib/async_dns.h" #include #include #include @@ -19,7 +20,6 @@ #include #include #include -#include #include #endif @@ -29,9 +29,18 @@ static void on_fin_cb(struct tcp_conn* tc, void* arg); static void on_error_cb(struct tcp_conn* tc, int err, void* arg); static void on_closed_cb(struct tcp_conn* tc, void* arg); static void tx_waiter_cb(struct ll_queue* q, void* arg); +static void tx_retry_timer_cb(void* arg); +static void socks_issue_dns(struct socks_proxy_conn* c, const char* host, uint8_t kind); +static void socks_finalize(struct socks_proxy_conn* c); +static void socks_dns_done_cb(const struct adns_result* res, void* arg); +static void socks_dns_error_and_close(struct socks_proxy_conn* c); +static void socks_bp_register(struct socks_proxy_conn* c); +static void socks_maybe_relay_fin(struct socks_proxy_conn* c); static int 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); +enum { DNSK_SOCKS = 0, DNSK_HTTP_CONNECT = 1, DNSK_HTTP_PROXY = 2 }; + struct listen_ctx { socket_t listen_sock; void* socket_id; @@ -76,8 +85,8 @@ static int send_connect(struct socks_proxy_conn* c) { return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_CONNECT, c->stream_id, buf, 6, 1); } -static int send_data(struct socks_proxy_conn* c, const uint8_t* data, uint16_t len) { - return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_DATA, c->stream_id, data, len, 0); +static int send_data(struct socks_proxy_conn* c, const uint8_t* data, uint16_t len, int force) { + return send_msg(c->inst, TOPO_GROUP_UTUN, c->via_node_id, TCP_PROXY_SUBCMD_DATA, c->stream_id, data, len, force); } static void send_close(struct socks_proxy_conn* c) { @@ -149,35 +158,26 @@ static void process_socks_request(struct socks_proxy_conn* c) { if (atyp == 1) { memcpy(c->dest_ip, c->buf + 4, 4); memcpy(&c->dest_port, c->buf + 8, 2); + c->buf_len = 0; + socks_finalize(c); + return; } else if (atyp == 4) { // Извлекаем первые 4 байта IPv6 в dest_ip (для упрощения: IPv6 mapped IPv4 или реальный IPv6) memcpy(c->dest_ip, c->buf + 12, 4); memcpy(&c->dest_port, c->buf + 20, 2); // TODO: proper IPv6 support DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: IPv6 unsupported, using last 4 bytes of addr"); - } else { // atyp == 3 (domain) - uint8_t dlen = c->buf[4]; - char domain[256]; memcpy(domain, c->buf + 5, dlen); domain[dlen] = '\0'; - memcpy(&c->dest_port, c->buf + 5 + dlen, 2); - // resolve domain - struct hostent* he = gethostbyname(domain); - if (!he || he->h_addrtype != AF_INET) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS failed for %s", domain); - goto error; - } - memcpy(c->dest_ip, he->h_addr, 4); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: resolved %s → %d.%d.%d.%d:%d", - domain, c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], ntohs(c->dest_port)); + c->buf_len = 0; + socks_finalize(c); + return; } + // atyp == 3 (domain) + uint8_t dlen = c->buf[4]; + char domain[256]; memcpy(domain, c->buf + 5, dlen); domain[dlen] = '\0'; + memcpy(&c->dest_port, c->buf + 5 + dlen, 2); c->buf_len = 0; - c->state = SOCKS_STATE_RELAY; // relay immediately, no waiting for first DATA - uint8_t reply[] = { 0x05, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; - write_to_client(c, reply, 10); - if (send_connect(c) < 0) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: send_connect failed sid=%08x", c->stream_id); - tcp_conn_push_close(c->tc); - } + socks_issue_dns(c, domain, DNSK_SOCKS); return; error: { @@ -213,21 +213,10 @@ static void process_http_request(struct socks_proxy_conn* c) { *colon = '\0'; int port = atoi(colon + 1); if (port <= 0 || port > 65535) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: bad port %d", port); goto error; } - struct hostent* he = gethostbyname(host_port); - if (!he || he->h_addrtype != AF_INET) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS failed for %s", host_port); - uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return; - } - memcpy(c->dest_ip, he->h_addr, 4); c->dest_port = htons((uint16_t)port); - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP CONNECT %s → %d.%d.%d.%d:%d", - host_port, c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], port); - + c->dest_port = htons((uint16_t)port); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP CONNECT %s:%d sid=%08x", host_port, port, c->stream_id); c->buf_len = 0; - c->state = HTTP_STATE_RELAY; - uint8_t resp[] = "HTTP/1.1 200 Connection Established\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); - if (send_connect(c) < 0) { tcp_conn_push_close(c->tc); } + socks_issue_dns(c, host_port, DNSK_HTTP_CONNECT); return; } @@ -276,23 +265,9 @@ static void process_http_request(struct socks_proxy_conn* c) { if (host[0] == '\0') { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: empty host in %s", url); goto error; } if (!path_start || *path_start == '\0') path_start = "/"; - struct hostent* he = gethostbyname(host); - if (!he || he->h_addrtype != AF_INET) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS failed for %s", host); - uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return; - } - memcpy(c->dest_ip, he->h_addr, 4); c->dest_port = htons((uint16_t)port); - - DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP %s %s:%d%s → %d.%d.%d.%d:%d", - method, host, port, path_start, - c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], port); - - if (send_connect(c) < 0) { - DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: send_connect failed for %s", host); - uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n"; - write_to_client(c, resp, (uint16_t)strlen((char*)resp)); tcp_conn_push_close(c->tc); return; - } + c->dest_port = htons((uint16_t)port); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: HTTP %s %s:%d%s sid=%08x", + method, host, port, path_start, c->stream_id); const char* version = strstr(line, " HTTP/"); char version_str[16] = " HTTP/1.1"; @@ -321,16 +296,10 @@ static void process_http_request(struct socks_proxy_conn* c) { memcpy(pkt, hdr_buf, hdr_chunk); if (extra) memcpy(pkt + hdr_chunk, c->buf + headers_len, extra); - int ret = send_data(c, pkt, total); - if (ret == 0) { - u_free(pkt); - } else { - c->tx_buf = pkt; c->tx_len = total; - etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); - } - + c->dns_pending = pkt; + c->dns_pending_len = total; c->buf_len = 0; - c->state = HTTP_STATE_RELAY; + socks_issue_dns(c, host, DNSK_HTTP_PROXY); return; error: { @@ -339,6 +308,88 @@ error: { } } +// ==================================================================== +// DNS-резолвинг (неблокирующий) + финализация рукопожатия +// ==================================================================== +static void socks_finalize(struct socks_proxy_conn* c) { + switch (c->dns_kind) { + case DNSK_SOCKS: { + uint8_t reply[] = { 0x05, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; + write_to_client(c, reply, 10); + break; + } + case DNSK_HTTP_CONNECT: { + uint8_t resp[] = "HTTP/1.1 200 Connection Established\r\n\r\n"; + write_to_client(c, resp, (uint16_t)strlen((char*)resp)); + break; + } + case DNSK_HTTP_PROXY: + default: + break; // reply нет — сразу CONNECT + данные + } + + if (send_connect(c) < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: send_connect failed sid=%08x kind=%d", c->stream_id, c->dns_kind); + tcp_conn_push_close(c->tc); + return; + } + + if (c->dns_kind == DNSK_HTTP_PROXY && c->dns_pending) { + uint8_t* pkt = c->dns_pending; + uint16_t len = c->dns_pending_len; + c->dns_pending = NULL; c->dns_pending_len = 0; + int ret = send_data(c, pkt, len, 0); + if (ret == 0) { + u_free(pkt); + } else { + c->tx_buf = pkt; c->tx_len = len; + socks_bp_register(c); + } + } + + c->state = c->is_http ? HTTP_STATE_RELAY : SOCKS_STATE_RELAY; +} + +static void socks_dns_error_and_close(struct socks_proxy_conn* c) { + if (c->is_http) { + uint8_t resp[] = "HTTP/1.1 502 Bad Gateway\r\n\r\n"; + write_to_client(c, resp, (uint16_t)strlen((char*)resp)); + } else { + uint8_t err[] = { 0x05, 0x04, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00 }; + write_to_client(c, err, 10); + } + tcp_conn_push_close(c->tc); +} + +static void socks_issue_dns(struct socks_proxy_conn* c, const char* host, uint8_t kind) { + c->dns_kind = kind; + if (inet_pton(AF_INET, host, &c->dest_ip) == 1) { + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: literal IP %s sid=%08x", host, c->stream_id); + socks_finalize(c); + return; + } + c->state = c->is_http ? HTTP_STATE_CONNECTING : SOCKS_STATE_CONNECTING; + c->dns_q = adns_resolve(c->ua, host, NULL, socks_dns_done_cb, c); + if (!c->dns_q) { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS resolve start failed for %s sid=%08x", host, c->stream_id); + socks_dns_error_and_close(c); + } +} + +static void socks_dns_done_cb(const struct adns_result* res, void* arg) { + struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; + c->dns_q = NULL; + if (res->status == ADNS_OK && res->count > 0) { + memcpy(c->dest_ip, &res->addrs[0].sin_addr, 4); + DEBUG_INFO(DEBUG_CATEGORY_PROXY, "socks_proxy: resolved → %d.%d.%d.%d:%d sid=%08x", + c->dest_ip[0], c->dest_ip[1], c->dest_ip[2], c->dest_ip[3], ntohs(c->dest_port), c->stream_id); + socks_finalize(c); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: DNS failed status=%d sid=%08x", res->status, c->stream_id); + socks_dns_error_and_close(c); + } +} + // ==================================================================== // tcp_io read callback — парсинг рукопожатия + релей данных // ==================================================================== @@ -374,13 +425,33 @@ static void on_read_cb(struct ll_queue* q, void* arg) { return; } + if (c->state == SOCKS_STATE_CONNECTING || c->state == HTTP_STATE_CONNECTING) { + // DNS в процессе: копим входящий body (HTTP proxy) / дропаем (CONNECT) + if (c->dns_kind == DNSK_HTTP_PROXY) { + size_t nl = (size_t)c->dns_pending_len + e->len; + if (nl <= UINT16_MAX) { + uint8_t* nb = u_realloc(c->dns_pending, nl); + if (nb) { memcpy(nb + c->dns_pending_len, e->dgram, e->len); c->dns_pending = nb; c->dns_pending_len = (uint16_t)nl; } + else DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: dns_pending realloc failed sid=%08x", c->stream_id); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: dns_pending overflow sid=%08x", c->stream_id); + } + } else { + DEBUG_WARN(DEBUG_CATEGORY_PROXY, "socks_proxy: data during DNS (kind=%d) sid=%08x — drop len=%u", c->dns_kind, c->stream_id, e->len); + } + memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); + queue_resume_callback(q); + return; + } + // RELAY = релей данных в ETCP if (c->tx_buf) { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: RELAY with pending tx_buf sid=%08x — drop new len=%u", c->stream_id, e->len); memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); + queue_resume_callback(q); return; } - int ret = send_data(c, e->dgram, e->len); + int ret = send_data(c, e->dgram, e->len, 0); DEBUG_TRACE(DEBUG_CATEGORY_PROXY, "SOCKS PROXY SEND sid=%08x len=%u ret=%d is_http=%d", c->stream_id, e->len, ret, c->is_http); if (ret == 0) { memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); @@ -390,7 +461,7 @@ static void on_read_cb(struct ll_queue* q, void* arg) { if (c->tx_buf) { memcpy(c->tx_buf, e->dgram, e->len); c->tx_len = e->len; } else { DEBUG_ERROR(DEBUG_CATEGORY_PROXY, "socks_proxy: tx_buf malloc=%u failed sid=%08x — drop", e->len, c->stream_id); } memory_pool_free(c->tc->data_pool, e->dgram); queue_entry_free(e); - etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); + socks_bp_register(c); DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY BP: sid=%08x tx_buf=%u waiter_reg", c->stream_id, c->tx_len); } } @@ -398,14 +469,47 @@ static void on_read_cb(struct ll_queue* q, void* arg) { static void tx_waiter_cb(struct ll_queue* q, void* arg) { (void)q; struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; - if (!c->tx_buf || c->rem_closed || c->close_sent) return; - int ret = send_data(c, c->tx_buf, c->tx_len); + if (c->rem_closed || c->close_sent) return; + if (!c->tx_buf) { if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); return; } + int ret = send_data(c, c->tx_buf, c->tx_len, 0); DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY WAKE: sid=%08x tx_buf=%u ret=%d is_http=%d", c->stream_id, c->tx_len, ret, c->is_http); if (ret == 0) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; - queue_resume_callback(c->tc->read_queue); + if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); + socks_maybe_relay_fin(c); + } else { + socks_bp_register(c); + } +} + +static void tx_retry_timer_cb(void* arg) { + struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; + c->tx_retry_timer = NULL; + if (c->rem_closed || c->close_sent) return; + if (!c->tx_buf) return; + int ret = send_data(c, c->tx_buf, c->tx_len, 1); + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "SOCKS PROXY RETRY: sid=%08x len=%u force=1 ret=%d", c->stream_id, c->tx_len, ret); + if (ret == 0) { + u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; + if (c->tc && c->tc->read_queue) queue_resume_callback(c->tc->read_queue); + socks_maybe_relay_fin(c); } else { - etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); + c->tx_retry_timer = uasync_set_timeout(c->ua, 5000, c, tx_retry_timer_cb, "socks_retry"); + } +} + +static void socks_bp_register(struct socks_proxy_conn* c) { + etcp_router_on_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter, tx_waiter_cb, c); + if (!c->tx_retry_timer) + c->tx_retry_timer = uasync_set_timeout(c->ua, 5000, c, tx_retry_timer_cb, "socks_retry"); +} + +static void socks_maybe_relay_fin(struct socks_proxy_conn* c) { + if (c->fin_deferred && !c->tx_buf) { + c->fin_deferred = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: relay deferred FIN sid=%08x", c->stream_id); + send_fin(c); + if (c->fin_remote && !c->close_sent && !c->close_pending) send_close(c); } } @@ -415,6 +519,12 @@ static void tx_waiter_cb(struct ll_queue* q, void* arg) { static void on_fin_cb(struct tcp_conn* tc, void* arg) { struct socks_proxy_conn* c = (struct socks_proxy_conn*)arg; if (c->rem_closed || c->close_sent) return; + if (c->tx_buf) { + // отложенный FIN: дождёмся сброса tx_buf (backpressure) + c->fin_deferred = 1; + DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: local FIN deferred (tx_buf pending) sid=%08x", c->stream_id); + return; + } if (!tc->fin_local && !tc->write_buf && !tc->write_queue->head) { DEBUG_DEBUG(DEBUG_CATEGORY_PROXY, "socks_proxy: local FIN → relay FIN sid=%08x", c->stream_id); send_fin(c); @@ -560,8 +670,11 @@ void socks_proxy_conn_free(struct socks_proxy_conn* c) { } if (head) { struct socks_proxy_conn** prev = head; while (*prev) { if (*prev == c) { *prev = c->next; if (count) (*count)--; break; } prev = &(*prev)->next; } } if (c->tx_buf) { u_free(c->tx_buf); c->tx_buf = NULL; c->tx_len = 0; } + if (c->dns_pending) { u_free(c->dns_pending); c->dns_pending = NULL; c->dns_pending_len = 0; } + if (c->dns_q) { adns_cancel(c->dns_q); c->dns_q = NULL; } + if (c->tx_retry_timer) { uasync_cancel_timeout(c->ua, c->tx_retry_timer); c->tx_retry_timer = NULL; } if (c->tc) { tcp_conn_destroy(c->tc); c->tc = NULL; } - if (c->inst) etcp_router_waiter_cancel(c->inst, TOPO_GROUP_UTUN, c->via_node_id, &c->tx_waiter); + if (c->inst) etcp_router_cancel_send_ready(c->inst, TOPO_GROUP_UTUN, c->via_node_id, ETCP_RT_ID_TCP_PROXY_SERVER, &c->tx_waiter); u_free(c); } diff --git a/src/proxy/socks_proxy.h b/src/proxy/socks_proxy.h index 5b31882e..871c0660 100644 --- a/src/proxy/socks_proxy.h +++ b/src/proxy/socks_proxy.h @@ -16,18 +16,22 @@ struct UTUN_INSTANCE; struct ll_entry; struct ll_queue; struct tcp_conn; +struct adns_query; enum { SOCKS_STATE_GREETING = 0, // ждём SOCKS5 greeting (ver + methods) - SOCKS_STATE_REQUEST, // ждём SOCKS5 CONNECT request - SOCKS_STATE_CONNECTING, // отправили ETCP CONNECT, ждём ответ от exit - SOCKS_STATE_RELAY // релей данных + SOCKS_STATE_REQUEST, // ждём SOCKS5 CONNECT request (== HTTP_STATE_REQUEST) + SOCKS_STATE_CONNECTING, // DNS в процессе (== HTTP_STATE_CONNECTING) + SOCKS_STATE_RELAY // релей данных (== HTTP_STATE_RELAY) }; +// HTTP-состояния совпадают по значениям с SOCKS-состояниями (общая state-машина +// в on_read_cb), поэтому явно привязаны к SOCKS_STATE_* — иначе возникает коллизия +// SOCKS_STATE_CONNECTING(2) == HTTP_STATE_RELAY(2) и RELAY-данные дропаются. enum { - HTTP_STATE_REQUEST = 0, // ждём "CONNECT host:port HTTP/1.1\r\n" - HTTP_STATE_CONNECTING, // отправили ETCP CONNECT, ждём ответ от exit - HTTP_STATE_RELAY // релей данных + HTTP_STATE_REQUEST = SOCKS_STATE_REQUEST, // ждём первую строку HTTP-запроса + HTTP_STATE_CONNECTING = SOCKS_STATE_CONNECTING, // DNS в процессе + HTTP_STATE_RELAY = SOCKS_STATE_RELAY // релей данных }; struct socks_proxy_conn { @@ -55,6 +59,12 @@ struct socks_proxy_conn { uint32_t bytes_to_client; // байт записано клиенту (curl) через write_queue uint32_t bytes_from_exit; // байт получено от exit node через ETCP DATA struct queue_waiter_handle tx_waiter; + void* tx_retry_timer; // safety-net таймер backpressure (force=1) + uint8_t fin_deferred; // FIN отложен до сброса tx_buf + struct adns_query* dns_q; // активный DNS-запрос (NULL когда нет) + uint8_t dns_kind; // DNSK_SOCKS / DNSK_HTTP_CONNECT / DNSK_HTTP_PROXY + uint8_t* dns_pending; // HTTP proxy: накопленный запрос (headers+body) до DNS + uint16_t dns_pending_len; }; struct listen_ctx;