Browse Source
- tcp_proxy/remote_proxy → tcp_proxy_client/tcp_proxy_server - Все proxy файлы перенесены в src/proxy/ - Конфиг: [tcp_proxy]→[tcp_proxy_client], [remote_proxy]→[tcp_proxy_server] - Убран subcmd CONNECTED (0x02) — клиент шлёт DATA сразу после CONNECT - Backpressure: etcp_router_waiter_register/cancel через ll_queue threshold waiter - Клиент: recv_cb send fail → tx_buf+waiter, не tcp_recved() (окно lwIP=0) - Сервер: read_cb send fail → pause_buf+waiter, пауза read_id - CLOSE/ERROR отправляется только после drain всех очередей - Fail: conn не найден/closed → send ERROR вместо CLOSE/silent drop - PROTOCOL.md: полная архитектура с диаграммами состояний и потоковetcp-inflight-fix
41 changed files with 1656 additions and 656 deletions
@ -0,0 +1,477 @@
|
||||
## TCP Proxy Protocol — архитектура (v2, без CONNECTED) |
||||
|
||||
### 1. Протокольный формат (общий для client и server) |
||||
|
||||
``` |
||||
┌──────┬──────┬───────────────────────┬──────────────────────┐ |
||||
│svc_id│subcmd│ stream_id (4) │ data ... │ |
||||
│ 1B │ 1B │ │ │ |
||||
└──────┴──────┴───────────────────────┴──────────────────────┘ |
||||
HDR = 6 байт |
||||
|
||||
subcmd: |
||||
0x01 CONNECT client→server : dest_ip(4)+dest_port(2) |
||||
0x03 DATA ↔ bidirectional: payload |
||||
0x04 CLOSE ↔ bidirectional: (нет данных) |
||||
0x05 ERROR ↔ bidirectional: (нет данных) |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 2. Диаграмма состояний `tcp_proxy_client_conn` |
||||
|
||||
``` |
||||
┌──────────────┐ |
||||
│ ALLOC + │ accept_cb: u_calloc, stream_id++ |
||||
│ CONNECT │ send CONNECT → exit |
||||
│ отправлен │ |
||||
└──────┬───────┘ |
||||
│ |
||||
┌───────────┴───────────┐ |
||||
│ send успешен? │ |
||||
└───────────┬───────────┘ |
||||
Y │ N |
||||
┌────────┘ └──────┐ |
||||
▼ ▼ |
||||
┌────────────┐ ┌──────────────┐ |
||||
│ ACTIVE │ │ FREE │ |
||||
│ traffic │ └──────────────┘ |
||||
└────┬──┬────┘ |
||||
│ │ |
||||
│ └─────────────────────────────────┐ |
||||
│ │ |
||||
▼ lwIP FIN (p=NULL) ▼ ERROR от lwIP/exit |
||||
┌──────────────┐ ┌──────────────┐ |
||||
│ TUN_CLOSED │ │ ERROR │ |
||||
│ tun_closed=1 │ │ error=1 │ |
||||
│ ↓ send_close │ │ send ERROR→ │ |
||||
│ после drain │ └──────┬───────┘ |
||||
│ tx_buf+lwIP │ │ |
||||
└──────┬───────┘ poll_cb: tcp_abort+free |
||||
│ |
||||
│ recv CLOSE от exit |
||||
▼ |
||||
┌───────────────────┐ |
||||
│ TUN+REM CLOSED │ |
||||
│ tun=1 rem=1 │ |
||||
│ ждём flush буферов│ |
||||
└────────┬──────────┘ |
||||
│ poll_cb: pending=0, unsent=NULL, unacked=NULL |
||||
│ (или sndbuf==0) |
||||
▼ |
||||
┌──────────────┐ |
||||
│ CLEANUP │ |
||||
│ free(pc) │ |
||||
└──────────────┘ |
||||
|
||||
|
||||
ДОП. СОСТОЯНИЯ: |
||||
┌─────────────────┐ send_data в ETCP вернул -1 (normalizer полон) |
||||
│ TX_BACKPRESSURE│ не tcp_recved() → окно lwIP=0 → recv_cb остановлен |
||||
│ tx_buf != NULL │ etcp_router_waiter_register() → ждём освобождения |
||||
│ tx_waiter pend │ waiter_cb: retry send_data → tcp_recved → окно открыто |
||||
└─────────────────┘ |
||||
|
||||
┌─────────────────┐ send_close/send_error вернул -1 |
||||
│ CLOSE_PENDING │ → retry в poll_cb |
||||
│ close_pending=1│ |
||||
└─────────────────┘ |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 3. Диаграмма состояний `tcp_proxy_server_conn` |
||||
|
||||
``` |
||||
┌───────────────────┐ |
||||
│ CREATING │ recv CONNECT |
||||
│ socket() создан │ → неблокирующий connect() |
||||
│ sock = неблок. │ |
||||
└────────┬──────────┘ |
||||
│ |
||||
┌────────────┴────────────┐ |
||||
│ connect() результат? │ |
||||
└────┬──────────┬──────────┘ |
||||
EINPROGRESS│ │ошибка сразу |
||||
▼ ▼ |
||||
┌──────────┐ ┌──────────────┐ |
||||
│CONNECTING│ │ ERROR │ send ERROR→client |
||||
│ждём write│ │ FREE │ |
||||
└────┬─────┘ └──────────────┘ |
||||
│ write_cb: SO_ERROR==0 |
||||
▼ |
||||
┌─────────────┐ |
||||
│ CONNECTED │ connected=1 |
||||
│ flush pend │ → flush pending_buf в сокет |
||||
└──────┬──────┘ |
||||
│ |
||||
┌──────────────┼──────────────┐ |
||||
│ │ │ |
||||
client→DATA │ │ ERROR (sock/pipe) |
||||
(sock_send) │ ▼ |
||||
или pending_buf │ ┌──────────────┐ |
||||
│ │ │ ERROR │ send ERROR→client |
||||
│ │ │ FREE │ |
||||
│ sock EOF (n=0) │ │ |
||||
│ │ └──────────────┘ |
||||
│ ▼ |
||||
│ ┌──────────────────┐ |
||||
│ │ SOCK_CLOSED │ sock_closed=1 |
||||
│ │ out_buf/pause? │ если очереди пусты → send_close |
||||
│ │ ждём drain │ иначе ждём retry_cb/waiter_cb |
||||
│ └────────┬─────────┘ |
||||
│ │ очереди опустошены → send_close |
||||
│ │ если cli_closed → FREE |
||||
│ │ |
||||
client→CLOSE │ |
||||
│ │ |
||||
▼ │ |
||||
┌──────────────┐ │ |
||||
│ CLI_CLOSED │ │ cli_closed=1 |
||||
│ out_buf pend? │◄────────┘ если out_buf пуст → shutdown(SHUT_WR) |
||||
│ ждём flush │ иначе ждём sock_retry_cb → shutdown |
||||
└──────┬───────┘ |
||||
│ |
||||
└──────────────┐ |
||||
▼ |
||||
┌──────────────────┐ |
||||
│ cli && sock │ оба закрыты + очереди пусты |
||||
│ очереди пусты │ |
||||
└────────┬─────────┘ |
||||
▼ |
||||
┌──────────────┐ |
||||
│ FREE │ conn_free(rc) |
||||
└──────────────┘ |
||||
|
||||
|
||||
ДОП. СОСТОЯНИЯ: |
||||
┌─────────────────┐ send_msg в ETCP вернул -1 |
||||
│ PAUSED │ буферизуем в pause_buf, убираем read_id (пауза чтения) |
||||
│ pause_buf>0 │ etcp_router_waiter_register() → ждём освобождения |
||||
│ pause_waiter │ waiter_cb: retry send_msg → resume read_id |
||||
└─────────────────┘ |
||||
|
||||
┌─────────────────┐ CLOSE/ERROR не доставлен через ETCP |
||||
│ CLOSE_PENDING │ close_timer с backoff 50..5000 tb → ретрай |
||||
│ close_pending=1│ |
||||
└─────────────────┘ |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 4. DATA FLOW |
||||
|
||||
``` |
||||
══════════════════════════════════════════════════════════════════════ |
||||
APP → lwIP → recv_cb → ETCP → handle → socket → DEST |
||||
APP ← lwIP ← tcp_write ← ETCP ← read_cb ← socket ← DEST |
||||
══════════════════════════════════════════════════════════════════════ |
||||
|
||||
CLIENT (узел A) SERVER (узел B) |
||||
|
||||
APP (через TUN) DEST (реальный) |
||||
│ │ |
||||
▼ TCP SYN/ACK (lwIP) ▼ TCP SYN/ACK (OS) |
||||
┌──────────┐ ┌──────────┐ |
||||
│ lwIP │ │ socket │ |
||||
│ TCP │ │ (неблок) │ |
||||
└────┬─────┘ └────┬─────┘ |
||||
│ recv_cb │ read_cb |
||||
│ │ |
||||
┌────▼────────────┐ ┌────▼───────────────┐ |
||||
│ tcp_proxy │ ETCP │ tcp_proxy │ |
||||
│ _client │◄═══════════════════════►│ _server │ |
||||
│ │ │ │ |
||||
│ ┌────────────┐ │ │ ┌──────────────┐ │ |
||||
│ │ recv_cb │ │──CONNECT+DATA──► │ │handle_connect│ │ |
||||
│ │(lwIP→ETCP) │ │ │ │handle_data │ │ |
||||
│ │ │ │ backpressure: │ └──────┬───────┘ │ |
||||
│ │ send fail: │ │ waiter на normalizer │ │ sock_send │ |
||||
│ │ tx_buf+ │ │ │ ┌──────▼───────┐ │ |
||||
│ │ waiter │ │ │ │ out_buf │ │ |
||||
│ │ окно lwIP=0│ │ │ │ pending_buf │ │ |
||||
│ └────────────┘ │ │ └─────────────┘ │ |
||||
│ │ │ │ send() │ |
||||
│ ┌────────────┐ │ │ ┌──────▼───────┐ │ |
||||
│ │handle_data │ │◄────DATA──────── │ │ read_cb │ │ |
||||
│ │handle_close│ │ │ │(sock→ETCP) │ │ |
||||
│ └─────┬──────┘ │ │ │ │ │ |
||||
│ │ │ │ │ send fail: │ │ |
||||
│ ┌─────▼──────┐ │ │ │ pause_buf+ │ │ |
||||
│ │ to_lwip │ │ │ │ waiter+ │ │ |
||||
│ │ (очередь) │ │ │ │ пауза read_id│ │ |
||||
│ └─────┬──────┘ │ │ └─────────────┘ │ |
||||
│ │tcp_write│ │ │ |
||||
│ ▼ │ └───────────────────┘ |
||||
│ lwIP→TUN→APP │ |
||||
└─────────────────┘ |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 5. BACKPRESSURE (lwIP→ETCP и socket→ETCP) |
||||
|
||||
``` |
||||
══════════════════════════════════════════════════════════════════════ |
||||
|
||||
Используется встроенный механизм ll_queue threshold waiter: |
||||
etcp_router_waiter_register(inst, peer_node_id, &waiter) |
||||
→ внутри: queue_waiter_wait(normalizer->input_queue, &waiter) |
||||
etcp_router_waiter_cancel(inst, peer_node_id, &waiter) |
||||
|
||||
CLIENT (lwIP → ETCP): |
||||
recv_cb(buf, len): |
||||
ret = send_data(buf, len) |
||||
if ret == 0: tcp_recved(len) |
||||
if ret < 0: |
||||
tx_buf = copy(buf), tx_len = len |
||||
НЕ tcp_recved() → окно lwIP = 0, recv_cb больше не вызывается |
||||
etcp_router_waiter_register(&tx_waiter) |
||||
|
||||
tx_waiter callback: |
||||
ret = send_data(tx_buf, tx_len) |
||||
if ret == 0: |
||||
tcp_recved(tx_len) → окно открыто, lwIP возобновляет recv_cb |
||||
free(tx_buf) |
||||
if tun_closed → send_close() |
||||
// если ret < 0 → waiter остаётся, normalizer вызовет снова |
||||
|
||||
SERVER (socket → ETCP): |
||||
read_cb(buf, n): |
||||
ret = send_msg(DATA, buf, n) |
||||
if ret == 0: ok |
||||
if ret < 0: |
||||
pause_buf = copy(buf), pause_len = n |
||||
uasync_remove_socket_t(sock) → пауза чтения с сокета |
||||
etcp_router_waiter_register(&pause_waiter) |
||||
|
||||
pause_waiter callback: |
||||
ret = send_msg(DATA, pause_buf, pause_len) |
||||
if ret == 0: |
||||
free(pause_buf) |
||||
uasync_add_socket_t(sock) → возобновляем чтение |
||||
if sock_closed → send_close() |
||||
// если ret < 0 → waiter остаётся |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 6. Последовательность: HANDSHAKE (упрощённый, без CONNECTED) |
||||
|
||||
``` |
||||
CLIENT (lwIP) CLIENT PROXY ETCP SERVER PROXY SOCKET |
||||
───────────────────────────────────────────────────────────────────────────────────────── |
||||
|
||||
tcp_proxy_client_accept_cb() CONNECT |
||||
├─ u_calloc(conn) ────────────────────────────→ tcp_proxy_server_handle_connect() |
||||
├─ stream_id++ ├─ socket(AF_INET,SOCK_STREAM) |
||||
├─ send_connect() ├─ socket_set_nonblocking() |
||||
│ [dest_ip(4)+dest_port(2)] ├─ uasync_add_socket_t(read/write/error) |
||||
├─ conn→conns list ├─ connect(неблок.) → EINPROGRESS |
||||
│ └─ rc→conns list |
||||
▼ |
||||
ACTIVE — данные идут СРАЗУ (без ожидания CONNECTED) |
||||
│ |
||||
│ DATA ...connect завершён... |
||||
├──────────────────────────────────────────────────→ tcp_proxy_server_sock_write_cb() |
||||
│ ├─ connected=1 |
||||
│ DATA └─ flush pending_buf в сокет |
||||
├──────────────────────────────────────────────────→ |
||||
│ |
||||
│ DATA ←─ tcp_proxy_server_sock_read_cb() |
||||
│◄────────────────────────────────────────────── (socket recv → ETCP) |
||||
│ |
||||
└──→ трафик идёт в обе стороны |
||||
|
||||
При ошибке connect: |
||||
...connect fail... → tcp_proxy_server_send_error() → client получает ERROR, error=1, cleanup |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 7. Последовательность: CLOSE (drain очередей → close в конце) |
||||
|
||||
``` |
||||
CLIENT (lwIP) CLIENT PROXY ETCP SERVER PROXY SOCKET |
||||
───────────────────────────────────────────────────────────────────────────────────────── |
||||
|
||||
══ Вариант А: lwIP закрывает первым ══ |
||||
|
||||
tcp_proxy_client_recv_cb(p=NULL) |
||||
├─ tun_closed=1 |
||||
│ если tx_buf пуст → send_close CLOSE |
||||
│ если tx_buf pend → waiter досылает ──────────→ tcp_proxy_server_handle_close() |
||||
│ потом send_close ├─ cli_closed=1 |
||||
└─ │ если out_buf пуст → shutdown(SHUT_WR) |
||||
│ иначе ждём sock_retry_cb |
||||
│ |
||||
...out_buf дослан... |
||||
sock_retry_cb: |
||||
shutdown(SHUT_WR) |
||||
│ |
||||
...сокет дочитывает остаток... |
||||
│ |
||||
tcp_proxy_server_sock_read_cb(n=0) |
||||
├─ sock_closed=1 |
||||
│ если pause_buf пуст → send_close |
||||
│ иначе ждём waiter_cb |
||||
│ |
||||
CLOSE │ |
||||
tcp_proxy_client_handle_close() ◄────────────────────────┘ |
||||
└─ rem_closed=1 |
||||
│ |
||||
poll_cb: tun=1, rem=1, to_lwip пуст, unsent=NULL, unacked=NULL |
||||
├─ tcp_close(pcb) |
||||
└─ conn_free() |
||||
|
||||
|
||||
══ Вариант Б: серверный сокет закрывается первым ══ |
||||
|
||||
tcp_proxy_server_sock_read_cb(n=0) |
||||
├─ sock_closed=1 |
||||
│ если out_buf+pause_buf пусты → send_close |
||||
│ иначе ждём retry/waiter |
||||
CLOSE │ |
||||
tcp_proxy_client_handle_close() ◄────────────────────────┘ |
||||
└─ rem_closed=1 |
||||
│ |
||||
▼ |
||||
...приложение закрывает TCP... |
||||
│ |
||||
tcp_proxy_client_recv_cb(p=NULL) |
||||
├─ tun_closed=1 |
||||
│ если tx_buf пуст → send_close |
||||
└─ |
||||
|
||||
poll_cb: оба закрыты, буферы пусты → cleanup |
||||
|
||||
|
||||
══ Вариант В: Ошибка / аномалия ══ |
||||
|
||||
ERROR |
||||
─────────────────────── (или ←) ────────── |
||||
error=1 |
||||
│ |
||||
▼ |
||||
poll_cb: tcp_abort(pcb) → немедленная очистка |
||||
или: rc→error=1 → conn_free |
||||
|
||||
|
||||
══ Вариант Г: conn не найден / closed / error ══ |
||||
|
||||
Любой пришедший пакет (DATA/CLOSE/ERROR) для conn в закрытом состоянии: |
||||
→ tcp_proxy_*_send_msg(ERROR, stream_id) |
||||
→ удалённая сторона получает ERROR и немедленно закрывается |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 8. Диспетчеризация в `tcp_proxy_client_etcp_recv_cb()` |
||||
|
||||
``` |
||||
┌────────────────────────────────────────────────────────────────────────┐ |
||||
│ tcp_proxy_client_etcp_recv_cb() │ |
||||
│ (единая точка входа для обоих направлений) │ |
||||
├────────────────────────────────────────────────────────────────────────┤ |
||||
│ │ |
||||
│ subcmd == CONNECT? │ |
||||
│ └── tcp_proxy_server_handle_connect() ←── всегда сервер │ |
||||
│ │ |
||||
│ stream_id в server.conns? (если server.enabled) │ |
||||
│ ├── DATA → tcp_proxy_server_handle_data() │ |
||||
│ ├── CLOSE → tcp_proxy_server_handle_close() │ |
||||
│ └── ERROR → rc->error=1; conn_free() │ |
||||
│ │ |
||||
│ stream_id в client.conns? (если client proxy != NULL) │ |
||||
│ ├── DATA → tcp_proxy_client_handle_data() │ |
||||
│ ├── CLOSE → tcp_proxy_client_handle_close() │ |
||||
│ └── ERROR → tcp_proxy_client_handle_error() │ |
||||
│ │ |
||||
│ Приоритет: │ |
||||
│ 1. CONNECT всегда интерпретируется сервером │ |
||||
│ 2. Если server.enabled И stream_id найден в server.conns → server │ |
||||
│ 3. Иначе если client proxy != NULL → client │ |
||||
│ 4. Иначе drop │ |
||||
└────────────────────────────────────────────────────────────────────────┘ |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 9. INIT |
||||
|
||||
``` |
||||
CLIENT SERVER |
||||
───────────────────────────────────────────────────────── |
||||
|
||||
utun_instance.c |
||||
│ |
||||
├─ tcp_proxy_client_create() tcp_proxy_server_init() |
||||
│ ├─ tun_init_nat() (TUN iface) ├─ memset(ctx,0) |
||||
│ ├─ lwip_tcp_init() (TCP стек) ├─ ctx->enabled = config ... |
||||
│ ├─ etcp_router_bind( ├─ etcp_router_bind( |
||||
│ │ ETCP_ID_TCP_PROXY, │ ETCP_ID_TCP_PROXY, |
||||
│ │ tcp_proxy_client_etcp_recv_cb) │ tcp_proxy_client_etcp_recv_cb) |
||||
│ ├─ udp_proxy_init() ├─ udp_proxy_init() |
||||
│ └─ icmp_proxy_init() └─ icmp_proxy_init() (если client не активен) |
||||
│ |
||||
└─ ОБА слушают ETCP_ID_TCP_PROXY через один обработчик |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 10. Поля структур |
||||
|
||||
``` |
||||
tcp_proxy_client_conn: |
||||
stream_id — уникальный ID потока |
||||
tun_closed — lwIP/FIN (сторона TUN) |
||||
rem_closed — CLOSE от exit |
||||
error — немедленная очистка |
||||
close_sent — CLOSE/ERROR отправлен |
||||
close_pending — отправка не удалась, ретрай |
||||
to_lwip — очередь DATA от exit → lwIP (feed_from_transport) |
||||
pcb — lwIP tcp_pcb |
||||
tx_buf — буфер при backpressure lwIP→ETCP (send_data fail) |
||||
tx_len — длина tx_buf |
||||
tx_waiter — ll_queue waiter на normalizer.input_queue |
||||
|
||||
tcp_proxy_server_conn: |
||||
stream_id — уникальный ID потока |
||||
cli_closed — клиент прислал CLOSE |
||||
sock_closed — сокет получил EOF |
||||
error — ошибка |
||||
connected — сокет подключился (connect завершён) |
||||
close_pending — CLOSE/ERROR не доставлен, ретрай |
||||
pending_buf — буфер данных от клиента до завершения connect |
||||
out_buf — буфер отправки в сокет (неблокирующий send, EAGAIN retry) |
||||
out_off/out_len — offset/размер out_buf |
||||
out_timer — таймер повтора send в сокет |
||||
out_backoff — backoff для retry send в сокет |
||||
pause_buf — буфер при backpressure socket→ETCP (send_msg fail) |
||||
pause_len — длина pause_buf |
||||
pause_waiter — ll_queue waiter на normalizer.input_queue |
||||
close_timer — таймер повтора CLOSE/ERROR |
||||
close_backoff — backoff для retry CLOSE/ERROR |
||||
sock — OS сокет к адресату |
||||
read_id — handle uasync для событий сокета |
||||
``` |
||||
|
||||
--- |
||||
|
||||
### 11. etcp_router waiter API |
||||
|
||||
``` |
||||
// etcp_router.h |
||||
void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, |
||||
uint64_t peer_node_id, struct queue_waiter_handle* h); |
||||
void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, |
||||
uint64_t peer_node_id, struct queue_waiter_handle* h); |
||||
|
||||
// Внутри: |
||||
// conn = route_bgp_find_conn_for_node(peer_node_id) |
||||
// queue_waiter_wait(conn->normalizer->input, h) |
||||
// queue_waiter_cancel(conn->normalizer->input, h) |
||||
// |
||||
// normalizer.input_queue уже имеет threshold=0 (ждёт полного освобождения) |
||||
// через queue_set_threshold(input_queue, 0, 0) в etcp.c |
||||
``` |
||||
Loading…
Reference in new issue