Browse Source

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)
v2
evgeny 4 weeks ago
parent
commit
fbea9c3869
  1. 115
      doc/TASK_tcp_proxy_throughput.md
  2. 35
      doc/_plan.txt
  3. 292
      doc/db_sync_v2.md
  4. 10
      doc/etcp_config.txt
  5. 99
      lib/async_dns.c
  6. 249
      src/proxy/socks_proxy.c
  7. 22
      src/proxy/socks_proxy.h

115
doc/TASK_tcp_proxy_throughput.md

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

35
doc/_plan.txt

@ -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% (обновляется каждый час)
Сделать подключение к группе

292
doc/db_sync_v2.md

@ -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_<name>_<id>" (
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 вместо первого попавшегося.

10
doc/etcp_config.txt

@ -1,10 +0,0 @@
В конфиге явно прописываются сокеты с ip и портом.
Через эти сокеты идёт весь трафик, случайные порты клиент не использует - только эти сокеты.
Соответственно, возможен только один линк между двумя пирами по одному каналу.
виртуальные сокеты - плохо т.к. используются ip/port для поиск подключения
./configure --with-openssl - openssl
./configure --without-openssl - tinycrypt

99
lib/async_dns.c

@ -20,6 +20,7 @@
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <ctype.h>
#ifdef _WIN32
#include <iphlpapi.h>
@ -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);

249
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 <stdlib.h>
#include <string.h>
#include <errno.h>
@ -19,7 +20,6 @@
#include <unistd.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <netdb.h>
#include <arpa/inet.h>
#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)
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);
// 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;
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 {
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);
}

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

Loading…
Cancel
Save