Browse Source

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
topo_upd
evgeny 2 months ago
parent
commit
a73cc1beea
  1. 2
      src/transport_layer/auto_socket.c
  2. 13
      src/transport_layer/etcp_connections.c
  3. 208
      tests/test_auto_socket_dynamic.c

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

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

208
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_<ifname>_<v4/v6>_<udp/tcp> */
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();

Loading…
Cancel
Save