Browse Source

feat: peer restart detection + etcp_router_conn_restart with service notification

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.
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
ca7f6c0d01
  1. 33
      src/etcp_router.c
  2. 4
      src/etcp_router.h
  3. 19
      src/proxy/tcp_proxy_client.c
  4. 1
      src/proxy/tcp_proxy_server.h

33
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
// ====================================================================

4
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

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

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

Loading…
Cancel
Save