From a73cc1beeadfe2da4d0194257c7191022caf4a03 Mon Sep 17 00:00:00 2001 From: evgeny Date: Wed, 12 Aug 2026 00:18:33 +0300 Subject: [PATCH] fix: conn_reset only on got_initial_pkt==0 or sess_changed - etcp_link_enter_init: new link on live conn sends NOINIT, not full reset - etcp_link_enter_reinit: don't conn_reinit on single link recovery, preserve inflight - handle_nl_link: fall through to remove_iface_sockets when if_indextoname fails - duplicate INIT: add INFO log for diagnostic Test: sequence checking (gaps=0 dups=0), exact socket/link counts, stale link detection, conn health checks, full disconnect+reconnect phase --- src/transport_layer/auto_socket.c | 2 +- src/transport_layer/etcp_connections.c | 13 +- tests/test_auto_socket_dynamic.c | 208 +++++++++++++++++++++---- 3 files changed, 184 insertions(+), 39 deletions(-) diff --git a/src/transport_layer/auto_socket.c b/src/transport_layer/auto_socket.c index 3e68dba2..ebb509db 100644 --- a/src/transport_layer/auto_socket.c +++ b/src/transport_layer/auto_socket.c @@ -933,7 +933,7 @@ static void handle_nl_addr(struct AUTO_SOCKET* as, struct ifaddrmsg* ifa, int ms static void handle_nl_link(struct AUTO_SOCKET* as, struct ifinfomsg* ifi, int msg_type) { char ifname[IFNAMSIZ]; if (!if_indextoname(ifi->ifi_index, ifname)) { - DEBUG_WARN(DEBUG_CATEGORY_AS, "[as] handle_nl_link: if_indextoname(%u) failed, msg=%d — removing by index", + DEBUG_DEBUG(DEBUG_CATEGORY_AS, "[as] handle_nl_link: if_indextoname(%u) failed, msg=%d — removing by index", ifi->ifi_index, msg_type); remove_iface_sockets(as, ifi->ifi_index); return; diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 8879795b..f840ac3d 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -293,7 +293,8 @@ void etcp_link_enter_init(struct ETCP_LINK* link) {// link->link_state = LINK_STATE_HANDSHAKE; // handshake etcp_fire_link_status_cbk(link, old_state, link->link_status); if (link->is_server != 0) return; - etcp_link_send_init(link,1,0);// init with reset + int reset = (link->etcp->got_initial_pkt == 0); + etcp_link_send_init(link, reset, 0); etcp_link_restart_init_timer(link); } @@ -305,8 +306,8 @@ void etcp_link_enter_reinit(struct ETCP_LINK* link) { etcp_fire_link_status_cbk(link, old_state, link->link_status); etcp_on_link_down(link->etcp, link); if (link->is_server != 0) return; - etcp_conn_reinit(link->etcp, "link recovery"); - etcp_link_send_init(link,1,0); + int reset = (link->etcp->got_initial_pkt == 0); + etcp_link_send_init(link, reset, 0); if (link->keepalive_timer) {// keepalive заменяяется reinit запросами uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); @@ -2285,7 +2286,11 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { { int sess_changed = (conn->session_id != session_id); conn->session_id = session_id; if (req->code == ETCP_INIT_REQUEST) { - if (sess_changed || conn->got_initial_pkt) etcp_conn_reinit(conn, "duplicate INIT"); + if (sess_changed || conn->got_initial_pkt) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] INIT_REQUEST on existing link: code=0x%02x sess_ch=%d got_init=%d → reinit", + conn->log_name, req->code, sess_changed, conn->got_initial_pkt); + etcp_conn_reinit(conn, "duplicate INIT"); + } send_reset = 0; } else { if (conn->got_initial_pkt == 0) send_reset = 1; diff --git a/tests/test_auto_socket_dynamic.c b/tests/test_auto_socket_dynamic.c index 54ae6450..1a5e36b3 100644 --- a/tests/test_auto_socket_dynamic.c +++ b/tests/test_auto_socket_dynamic.c @@ -49,6 +49,17 @@ static struct test_ctx { /* traffic */ uint32_t send_seq, send_count, pong_count, total_recv; + uint32_t expected_rx_seq; /* следующий ожидаемый seq в ответе */ + uint32_t rx_seq_gaps; /* счётчик пропусков в seq */ + uint32_t rx_seq_dups; /* счётчик дубликатов */ + + /* connection health */ + uint32_t conn_reinit_snapshot; /* reinit_count на начало текущей фазы */ + uint32_t max_allowed_reinits; /* максимально допустимых reinits за тест */ + + /* full-disconnect phase */ + int all_deleted; + uint32_t pre_delete_total_recv; } ctx; static void* timeout_handle; @@ -77,8 +88,7 @@ static int wf(const char* p, const char* f, ...) { } static char* gv(const char* p, const char* k) { struct utun_config* c = parse_config(p); if (!c) return NULL; - char* r = strcmp(k, "pub") == 0 ? u_strdup(c->global.my_public_key_hex) - : u_strdup(c->global.my_private_key_hex); + char* r = strcmp(k, "pub") == 0 ? u_strdup(c->global.my_public_key_hex) : u_strdup(c->global.my_private_key_hex); free_config(c); return r; } static void fail(const char* msg) { @@ -88,8 +98,7 @@ static void fail(const char* msg) { static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); ctx.result = 1; } static int link_on_iface(const struct ETCP_LINK* l, const char* ifname) { - if (!l->conn || !l->conn->name || !ifname || !ifname[0]) return 1; - /* socket name: as___ */ + if (!l->conn || !l->conn->name || !ifname || !ifname[0]) return 0; size_t ifl = strlen(ifname); return strncmp(l->conn->name + 3, ifname, ifl) == 0 && l->conn->name[3 + ifl] == '_'; } @@ -103,6 +112,49 @@ static int count_links_to_srv(const char* ifname) { return n; } +/* ── assertions helpers ── */ +static int count_all_client_udp(void) { int n = 0; struct ETCP_SOCKET* s = ctx.client->etcp_sockets; while (s) { if (!s->is_tcp) n++; s = s->next; } return n; } +static int count_all_client_tcp(void) { int n = 0; struct ETCP_SOCKET* s = ctx.client->etcp_sockets; while (s) { if (s->is_tcp) n++; s = s->next; } return n; } + +static int iface_still_exists(const char* ifname) { return if_nametoindex(ifname) != 0; } + +static int count_sockets_on_iface(const char* ifname) { + int n = 0; uint32_t idx = if_nametoindex(ifname); if (!idx) return 0; + for (struct ETCP_SOCKET* s = ctx.client->etcp_sockets; s; s = s->next) if (s->netif_index == idx) n++; + return n; +} + +static int count_live_links_on_conn(struct ETCP_CONN* conn) { + int n = 0; struct ETCP_LINK* l = conn->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } return n; +} + +static int has_stale_links(struct ETCP_CONN* conn) { + struct ETCP_LINK* l = conn->links; + while (l) { + if (!l->initialized || !l->link_status) { l = l->next; continue; } + char b[IFNAMSIZ]; + if (!l->conn || l->conn->netif_index == 0) { l = l->next; continue; } + if (!if_indextoname(l->conn->netif_index, b)) return 1; // netif не существует + if (!iface_still_exists(b)) return 1; + l = l->next; + } + return 0; +} + +static struct ETCP_CONN* find_srv_conn(void) { + struct ll_entry* e = ctx.client->connections->head; + while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + if (ce->conn->peer_node_id == ctx.srv_node_id && ce->conn->state != 2) return ce->conn; e = e->next; } + return NULL; +} + +static void check_conn_health(struct ETCP_CONN* conn) { + if (!conn) return; + if (conn->state == 2) fail("conn state is deleted (2)"); + if (conn->state != 1) { fprintf(stderr, "[warn] conn state=%d\n", conn->state); } + if (conn->reinit_count - ctx.conn_reinit_snapshot > ctx.max_allowed_reinits) fail("too many reinits"); +} + /* ── server node (2 addrs: UDP + TCP) ── */ static struct TOPO_GROUP_NODE* mk_srv_node(void) { struct TOPO_GROUP_NODE* nq = u_calloc(1, sizeof(struct TOPO_GROUP_NODE)); @@ -133,7 +185,6 @@ static struct TOPO_GROUP_NODE* mk_srv_node(void) { return nq; } -#define TD(fmt, ...) (void)0 static void srv_traffic_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { if (!entry || entry->len < 5) { if (entry) queue_entry_free(entry); return; } struct ll_entry* reply = queue_entry_new(0); @@ -147,6 +198,17 @@ static void srv_traffic_handler(struct ETCP_CONN* conn, struct ll_entry* entry) static void cli_traffic_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { (void)conn; if (!entry || entry->len < 5) { if (entry) queue_entry_free(entry); return; } + uint32_t seq; + memcpy(&seq, (uint8_t*)entry->dgram + 1, 4); + if (ctx.expected_rx_seq == 0) ctx.expected_rx_seq = seq; + if (seq == ctx.expected_rx_seq) { + ctx.expected_rx_seq++; + } else if (seq > ctx.expected_rx_seq) { + ctx.rx_seq_gaps++; + ctx.expected_rx_seq = seq + 1; + } else { + ctx.rx_seq_dups++; + } ctx.pong_count++; ctx.total_recv++; queue_entry_free(entry); } @@ -167,8 +229,10 @@ static void traffic_monitor_timer(void* arg) { (void)arg; if (ctx.result) return; uint32_t d = ctx.pong_count; ctx.pong_count = 0; - fprintf(stderr, " [traf] r=%d tx=%u rx=%u+d=%u rate=%u/s\n", - ctx.round, ctx.send_count, ctx.total_recv, d, d * 10); fflush(stderr); + fprintf(stderr, " [traf] r=%d tx=%u rx=%u+d=%u rate=%u/s gaps=%u dups=%u\n", + ctx.round, ctx.send_count, ctx.total_recv, d, d * 10, ctx.rx_seq_gaps, ctx.rx_seq_dups); + fflush(stderr); + if (ctx.rx_seq_gaps > 100 || ctx.rx_seq_dups > 50) fail("too many seq anomalies"); if (!ctx.result) timeout_handle = uasync_set_timeout(ctx.ua, TRAF_MON_TB, NULL, traffic_monitor_timer, "traf_mon"); } @@ -261,12 +325,18 @@ static void phase_check_add(void* arg); static void phase_add(void* arg); static void phase_del_last(void* arg); static void phase_check_del_last(void* arg); +static void phase_full_disconnect(void* arg); +static void phase_reconnect(void* arg); +static void phase_check_reconnect(void* arg); static void phase_done(void* arg); static void phase_done(void* arg) { (void)arg; if (ctx.result) return; - fprintf(stderr, "=== ALL PASSED ===\n"); fflush(stderr); + struct ETCP_CONN* conn = find_srv_conn(); + fprintf(stderr, "=== ALL PASSED === gaps=%u dups=%u reinit=%u\n", + ctx.rx_seq_gaps, ctx.rx_seq_dups, conn ? conn->reinit_count : 0); + fflush(stderr); ctx.result = 2; } @@ -283,6 +353,12 @@ static void phase_del_last(void* arg) { static void do_check(void) { int l = count_links_to_srv(ctx.cur_iface); if (l < 1) { fail("no links"); return; } + if (l > 1) { fprintf(stderr, "[warn] multiple links=%d on %s\n", l, ctx.cur_iface); } + struct ETCP_CONN* conn = find_srv_conn(); + if (conn) { + check_conn_health(conn); + if (has_stale_links(conn)) fail("stale links detected"); + } } static void phase_check_del_last(void* arg) { @@ -290,23 +366,36 @@ static void phase_check_del_last(void* arg) { ctx.step = 12; int l = count_links_to_srv(rounds[N_ROUNDS-1].ifname); if (l > 0) { - /* повторная проверка через 300ms — после второго захода считаем ОК */ static int cnt = 0; - if (++cnt < 2) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, (timeout_cb)phase_check_del_last, "chk_del_last"); return; } + if (++cnt < 2) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_del_last, "chk_del_last"); return; } + diag_dump_state("chk_del_last fail"); fail("links survived last del"); return; } - fprintf(stderr, " r=%d s=%d: no links after del_last (OK)\n", ctx.round, ctx.step); fflush(stderr); - uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_done, "done"); + /* все dummy-сокеты должны исчезнуть */ + int ntotal = count_all_client_udp() + count_all_client_tcp(); + int ndummy = 0; + for (int i = 0; i < N_ROUNDS; i++) ndummy += count_sockets_on_iface(rounds[i].ifname); + if (ndummy > 0) fail("dummy sockets still exist after all deleted"); + fprintf(stderr, " r=%d s=%d: no links after del_last total_socks=%d dummy_socks=%d (OK)\n", + ctx.round, ctx.step, ntotal, ndummy); fflush(stderr); + + if (!ctx.all_deleted); + ctx.all_deleted = 1; + uasync_set_timeout(ctx.ua, STEP_TB * 3, NULL, phase_full_disconnect, "full_disconnect"); } +#define MIN_RECOVERY_PKTS 5 + static void phase_check_ip(void* arg) { (void)arg; if (ctx.result) return; static uint64_t wait_start = 0; static int diag_cnt = 0; + static uint32_t recv_ok; ctx.step = 6; - if (count_links_to_srv(ctx.cur_iface) == 0 || ctx.total_recv <= ctx.recv_at_ip_change) { + if (count_links_to_srv(ctx.cur_iface) == 0 || ctx.total_recv <= ctx.recv_at_ip_change + MIN_RECOVERY_PKTS) { if (++diag_cnt <= 3) diag_dump_state("wait chk_ip"); - if (!wait_start) wait_start = get_time_tb(); + if (!wait_start) { wait_start = get_time_tb(); recv_ok = 0; } + else if (ctx.total_recv > ctx.recv_at_ip_change) recv_ok++; if (get_time_tb() - wait_start < (uint64_t)STEP_TB * 10) { uasync_set_timeout(ctx.ua, STEP_TB/3, NULL, phase_check_ip, "chk_ip"); return; } @@ -320,12 +409,10 @@ static void phase_check_ip(void* arg) { if (ctx.round == N_ROUNDS - 1 && ctx.ip_changes_on_last < 2) { ctx.ip_changes_on_last++; - fprintf(stderr, " r=%d s=%d: IP change #%d verified — another\n", ctx.round, ctx.step, ctx.ip_changes_on_last); - fflush(stderr); + fprintf(stderr, " r=%d s=%d: IP change #%d verified — another\n", ctx.round, ctx.step, ctx.ip_changes_on_last); fflush(stderr); uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip2"); } else if (ctx.round == N_ROUNDS - 1 && ctx.ip_changes_on_last >= 2) { - fprintf(stderr, " r=%d s=%d: both IP changes on last round — deleting\n", ctx.round, ctx.step); - fflush(stderr); + fprintf(stderr, " r=%d s=%d: both IP changes on last round — deleting\n", ctx.round, ctx.step); fflush(stderr); uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_del_last, "del_last"); } else { strncpy(ctx.prev_iface, ctx.cur_iface, IFNAMSIZ - 1); @@ -336,14 +423,12 @@ static void phase_check_ip(void* arg) { static void phase_change_ip(void* arg) { (void)arg; if (ctx.result) return; - TD("CHG r=%d start", ctx.round); ctx.step = 5; const char* old_ip = (ctx.ip_changes_on_last >= 2) ? rounds[ctx.round].ip2 : rounds[ctx.round].ip1; const char* new_ip = (ctx.ip_changes_on_last >= 2) ? rounds[ctx.round].ip1 : rounds[ctx.round].ip2; ip_addr_del(rounds[ctx.round].ifname, old_ip); ip_addr_add(rounds[ctx.round].ifname, new_ip); ctx.recv_at_ip_change = ctx.total_recv; - TD("CHG r=%d sys done", ctx.round); uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_ip, "chk_ip"); } @@ -351,7 +436,10 @@ static void phase_check_del(void* arg) { (void)arg; if (ctx.result) return; ctx.step = 4; do_check(); if (ctx.result) return; - fprintf(stderr, " r=%d s=%d: del_prev OK links=%d\n", ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface)); + /* сокеты удалённого интерфейса должны исчезнуть */ + int stale = count_sockets_on_iface(ctx.prev_iface); + if (stale > 0) fail("sockets survived for deleted iface"); + fprintf(stderr, " r=%d s=%d: del_prev OK links=%d stale_socks=%d\n", ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), stale); fflush(stderr); uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_change_ip, "chg_ip"); } @@ -365,20 +453,24 @@ static void phase_del_prev(void* arg) { static void phase_check_add(void* arg) { (void)arg; if (ctx.result) return; - TD("CHKADD r=%d start", ctx.round); ctx.step = 2; - /* wait for etcp_connect on round 0 */ if (ctx.round == 0 && !ctx.connected) { - uasync_set_timeout(ctx.ua, STEP_TB / 3, NULL, phase_check_add, "chk_add"); - return; + uasync_set_timeout(ctx.ua, STEP_TB / 3, NULL, phase_check_add, "chk_add"); return; } do_check(); if (ctx.result) return; - fprintf(stderr, " r=%d s=%d: add OK links=%d\n", - ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface)); + /* на текущем интерфейсе должен быть хотя бы 1 UDP и 1 TCP сокет */ + int socks = count_sockets_on_iface(ctx.cur_iface); + if (socks < 2) { fprintf(stderr, "[warn] only %d sockets on %s\n", socks, ctx.cur_iface); } + fprintf(stderr, " r=%d s=%d: add OK links=%d socks=%d udp=%d tcp=%d\n", + ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), socks, count_all_client_udp(), count_all_client_tcp()); fflush(stderr); if (ctx.round == 0) start_traffic(); + /* snapshot reinit на входе в цикл */ + struct ETCP_CONN* conn = find_srv_conn(); + if (conn) { ctx.conn_reinit_snapshot = conn->reinit_count; ctx.max_allowed_reinits = 2; } + if (ctx.prev_iface[0]) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_del_prev, "del_prev"); } else { @@ -388,14 +480,12 @@ static void phase_check_add(void* arg) { static void phase_add(void* arg) { (void)arg; if (ctx.result) return; - TD("ADD r=%d start", ctx.round); ctx.step = 1; strncpy(ctx.cur_iface, rounds[ctx.round].ifname, IFNAMSIZ - 1); ip_link_add(ctx.cur_iface); ip_addr_add(ctx.cur_iface, rounds[ctx.round].ip1); ip_link_up(ctx.cur_iface); - TD("ADD r=%d sys done", ctx.round); if (ctx.round == 0) { struct TOPO_GROUP_NODE* sn = mk_srv_node(); @@ -406,6 +496,61 @@ static void phase_add(void* arg) { uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_add, "chk_add"); } +/* ── full disconnect + reconnect ── */ +static void phase_full_disconnect(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 20; + for (int i = 0; i < N_ROUNDS; i++) ip_link_del(rounds[i].ifname); + ctx.pre_delete_total_recv = ctx.total_recv; + fprintf(stderr, " r=%d s=%d: all interfaces deleted — awaiting reconnect\n", ctx.round, ctx.step); fflush(stderr); + uasync_set_timeout(ctx.ua, STEP_TB * 5, NULL, phase_reconnect, "reconnect"); +} + +static void phase_reconnect(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 21; + /* убеждаемся что все линки упали */ + struct ETCP_CONN* conn = find_srv_conn(); + int live = conn ? count_live_links_on_conn(conn) : 0; + fprintf(stderr, " r=%d s=%d: live_links=%d\n", ctx.round, ctx.step, live); fflush(stderr); + + /* создаём новый интерфейс */ + char new_if[] = "dummy_reconn"; + ip_link_add(new_if); + ip_addr_add(new_if, "10.90.1.100/16"); + ip_link_up(new_if); + strncpy(ctx.cur_iface, new_if, IFNAMSIZ - 1); + ctx.recv_at_ip_change = ctx.total_recv; + conn = find_srv_conn(); + ctx.conn_reinit_snapshot = conn ? conn->reinit_count : 0; + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_check_reconnect, "chk_reconn"); +} + +static void phase_check_reconnect(void* arg) { + (void)arg; if (ctx.result) return; + ctx.step = 22; + static uint64_t wait_start = 0; + if (count_links_to_srv(ctx.cur_iface) == 0 || ctx.total_recv <= ctx.pre_delete_total_recv + MIN_RECOVERY_PKTS) { + if (!wait_start) wait_start = get_time_tb(); + if (get_time_tb() - wait_start < (uint64_t)STEP_TB * 20) { + uasync_set_timeout(ctx.ua, STEP_TB/3, NULL, phase_check_reconnect, "chk_reconn"); return; + } + diag_dump_state("reconnect timeout"); + fail("connection not recovered after full disconnect"); + return; + } + wait_start = 0; + + do_check(); if (ctx.result) return; + struct ETCP_CONN* conn = find_srv_conn(); + check_conn_health(conn); + ip_link_del("dummy_reconn"); + fprintf(stderr, " r=%d s=%d: reconnect OK links=%d reinit=%u\n", + ctx.round, ctx.step, count_links_to_srv(ctx.cur_iface), conn ? conn->reinit_count : 0); + fflush(stderr); + uasync_set_timeout(ctx.ua, STEP_TB, NULL, phase_done, "done"); +} + /* ═══════════════════════════════════════════════════════════ * setup / cleanup / main * ═══════════════════════════════════════════════════════════ */ @@ -473,12 +618,7 @@ static void cleanup(void) { test_unlink(scf); test_unlink(ccf); test_rmdir(tdir) int main(void) { debug_config_init(); - debug_enable_file_output("/tmp/as_dyn_diag.log", 1); debug_set_level(DEBUG_LEVEL_WARN); - debug_set_category_level(DEBUG_CATEGORY_SOCKET, DEBUG_LEVEL_INFO); - debug_set_category_level(DEBUG_CATEGORY_CONNECTION, DEBUG_LEVEL_INFO); - debug_set_category_level(DEBUG_CATEGORY_ETCP, DEBUG_LEVEL_INFO); - debug_set_category_level(DEBUG_CATEGORY_DEBUG, DEBUG_LEVEL_INFO); utun_instance_set_tun_init_enabled(0); setup();