From ca7f6c0d012308baf1ecfcf4b762847bd19ae149 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 3 Jun 2026 20:45:42 +0300 Subject: [PATCH] feat: peer restart detection + etcp_router_conn_restart with service notification MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit etcp_router: new etcp_router_conn_restart(inst, node_id, svc_id) — closes all rconns for the peer, notifies service via cb(NULL, entry) with node_id, resets seq state. seq=0 detection in etcp_router_recv_cb: when rx_seq >= 256 (MAX_INFLIGHT) and a data packet arrives with seq=0, detect peer restart and reset rconn. tcp_proxy_client: handle TCP_PROXY_SUBCMD_RESTART — clear all client conns and server conns for the restarted peer. Threshold prevents false positives on in-window seq=0 duplicates. --- src/etcp_router.c | 33 +++++++++++++++++++++++++++++++++ src/etcp_router.h | 4 ++++ src/proxy/tcp_proxy_client.c | 19 +++++++++++++++++++ src/proxy/tcp_proxy_server.h | 1 + 4 files changed, 57 insertions(+) diff --git a/src/etcp_router.c b/src/etcp_router.c index 24ceb60a..4aef3f2e 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -368,6 +368,15 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return; } } + // Обнаружение перезапуска peer'a: seq=0 при rx_seq >= MAX_INFLIGHT + if (rconn->rx_seq >= (uint32_t)ROUTER_MAX_INFLIGHT && hdr->seq == 0 && pl_len > 0) { + DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router: peer restart detected svc_id=%u from %016llx — resetting", + hdr->svc_id, (unsigned long long)hdr->src_node_id); + etcp_router_conn_restart(inst, hdr->src_node_id, hdr->svc_id); + rconn = etcp_router_conn_get(inst, hdr->src_node_id, hdr->svc_id); + if (!rconn) { queue_dgram_free(entry); queue_entry_free(entry); return; } + } + // Seq-состояние есть — reorder rconn->last_dgram_ts = get_current_timestamp(); uint32_t seq = hdr->seq; @@ -573,6 +582,30 @@ void etcp_router_consumer_ack(struct UTUN_INSTANCE* inst, uint64_t remote_node_i router_ack_do_send(rconn); } +void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t remote_node_id, uint8_t svc_id) { + struct ETCP_ROUTER_CONN* rconn = router_conn_find(inst, remote_node_id, svc_id); + if (!rconn) return; + + DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_restart: svc_id=%u remote=%016llx — notifying service", + svc_id, (unsigned long long)remote_node_id); + + etcp_recv_fn cb = inst->router_bindings.callbacks[svc_id]; + if (cb) { + struct ll_entry* e = queue_entry_new(0); + if (e) { + e->dgram = u_malloc(10); + if (e->dgram) { + e->dgram[0] = svc_id; + e->dgram[1] = 0xFE; // RESTART subcmd + memcpy(e->dgram + 2, &remote_node_id, 8); + e->len = 10; + cb(NULL, e); + } else { queue_entry_free(e); } + } + } + router_close_and_notify(rconn); +} + // ==================================================================== // Seq-connection API // ==================================================================== diff --git a/src/etcp_router.h b/src/etcp_router.h index 1047c2e1..a88cd39c 100644 --- a/src/etcp_router.h +++ b/src/etcp_router.h @@ -105,4 +105,8 @@ void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id // Уведомить router о потреблении данных потребителем (для consumer-driven ACK) void etcp_router_consumer_ack(struct UTUN_INSTANCE* inst, uint64_t remote_node_id, uint8_t svc_id); +// Сбросить состояние роутера для конкретного peer+svc (перезапуск удалённой стороны). +// Очищает send_q/recv_q, уведомляет сервис через cb(NULL, entry) с remote_node_id. +void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t remote_node_id, uint8_t svc_id); + #endif // ETCP_ROUTER_H diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index bae5ab27..e6149e86 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -476,6 +476,25 @@ void tcp_proxy_client_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entr struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL; struct tcp_proxy_client* proxy = inst ? inst->tcp_proxy_client : NULL; + if (subcmd == TCP_PROXY_SUBCMD_RESTART) { + uint64_t peer_id = 0; + if (entry->len >= 10) memcpy(&peer_id, entry->dgram + 2, 8); + DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY RESTART from %016llx — clearing all conns for peer", (unsigned long long)peer_id); + if (proxy) { + struct tcp_proxy_client_conn *pc, *next; + for (pc = proxy->conns; pc; pc = next) { next = pc->next; tcp_proxy_client_conn_free(pc); } + proxy->conns = NULL; proxy->conn_count = 0; + } + if (inst && inst->tcp_proxy_server.enabled) { + struct tcp_proxy_server_conn *rc, *next; + for (rc = inst->tcp_proxy_server.conns; rc; rc = next) { + next = rc->next; + if (rc->peer_node_id == peer_id) tcp_proxy_server_conn_free(rc); + } + } + queue_dgram_free(entry); queue_entry_free(entry); return; + } + if (subcmd == TCP_PROXY_SUBCMD_CONNECT) { uint64_t src_node_id = conn ? conn->peer_node_id : (inst ? inst->node_id : 0); tcp_proxy_server_handle_connect(inst, entry, stream_id, src_node_id); diff --git a/src/proxy/tcp_proxy_server.h b/src/proxy/tcp_proxy_server.h index 64f2520b..416430aa 100644 --- a/src/proxy/tcp_proxy_server.h +++ b/src/proxy/tcp_proxy_server.h @@ -16,6 +16,7 @@ struct ETCP_CONN; #define TCP_PROXY_SUBCMD_DATA 0x03 #define TCP_PROXY_SUBCMD_CLOSE 0x04 #define TCP_PROXY_SUBCMD_ERROR 0x05 +#define TCP_PROXY_SUBCMD_RESTART 0xFE #define TCP_PROXY_HDR_SIZE 6 // svc_id(1)+subcmd(1)+stream_id(4) #define TCP_PROXY_CONNECT_HDR_SIZE 12 // HDR_SIZE + dest_ip(4)+dest_port(2)