* [RFC net-next 1/2] tipc: abort pending connects when the peer node is lost
2026-09-29 6:47 [RFC net-next 0/2] tipc: notify pending connects of peer node loss liushike
@ 2026-09-29 6:47 ` liushike
2026-09-29 6:47 ` [RFC net-next 2/2] selftests: net: cover pending TIPC connects on peer node loss liushike
2026-09-29 6:54 ` [RFC net-next 0/2] tipc: notify pending connects of " netdev-bot+sinfo
2 siblings, 0 replies; 4+ messages in thread
From: liushike @ 2026-09-29 6:47 UTC (permalink / raw)
To: netdev
Cc: jmaloy, tung.quang.nguyen, davem, edumazet, kuba, pabeni, horms,
shuah, tipc-discussion, linux-kselftest, linux-kernel
From: liushike <liushike@ruijie.com.cn>
A client can send a SYN and block in connect() while the listening server
has not yet called accept(). If the last usable link to that server is
lost, the client may remain asleep until its connection timeout expires.
Only established connections are in the peer node's conn_sks list, so
node_lost_contact() cannot notify the pending connection.
Register an active open before sending SYN and check node availability
under the node write lock. Let the existing node-loss notification abort
the pending handshake with EHOSTUNREACH and wake the connect waiter.
Losing one link does not abort the connect while another link survives.
Update the listening port to the accepted port in place on receipt of
ACK, preserving the node-loss registration. Reject a late ACK if node
loss already removed the registration, even if the node has recovered.
Transfer the registration when a named SYN is rerouted to another node.
Avoid registering active opens twice; remove registrations on transmit
failure and terminal rejection, but retain them across overload retries
and connection-wait timeouts, which do not cancel the handshake.
Signed-off-by: liushike <liushike@ruijie.com.cn>
---
net/tipc/node.c | 37 ++++++++++++++++++++++++++++++++++++-
net/tipc/node.h | 1 +
net/tipc/socket.c | 38 +++++++++++++++++++++++++++++++++++---
3 files changed, 72 insertions(+), 4 deletions(-)
diff --git a/net/tipc/node.c b/net/tipc/node.c
index bd91378b7540..31f2ef44ad37 100644
--- a/net/tipc/node.c
+++ b/net/tipc/node.c
@@ -713,13 +713,48 @@ int tipc_node_add_conn(struct net *net, u32 dnode, u32 port, u32 peer_port)
conn->peer_port = peer_port;
tipc_node_write_lock(node);
- list_add_tail(&conn->list, &node->conn_sks);
+ if (node_is_up(node)) {
+ list_add_tail(&conn->list, &node->conn_sks);
+ } else {
+ err = -EHOSTUNREACH;
+ kfree(conn);
+ }
tipc_node_write_unlock(node);
exit:
tipc_node_put(node);
return err;
}
+/* Replace the listening port with the accepted port without losing the
+ * node loss subscription. A missing entry means node_lost_contact() has
+ * already aborted this connection attempt, even if the node is up again.
+ */
+int tipc_node_update_conn(struct net *net, u32 dnode, u32 port, u32 peer_port)
+{
+ struct tipc_sock_conn *conn;
+ struct tipc_node *node;
+ int err = -EHOSTUNREACH;
+
+ if (in_own_node(net, dnode))
+ return 0;
+
+ node = tipc_node_find(net, dnode);
+ if (!node)
+ return err;
+
+ tipc_node_write_lock(node);
+ list_for_each_entry(conn, &node->conn_sks, list) {
+ if (conn->port != port)
+ continue;
+ conn->peer_port = peer_port;
+ err = 0;
+ break;
+ }
+ tipc_node_write_unlock(node);
+ tipc_node_put(node);
+ return err;
+}
+
void tipc_node_remove_conn(struct net *net, u32 dnode, u32 port)
{
struct tipc_node *node;
diff --git a/net/tipc/node.h b/net/tipc/node.h
index 154a5bbb0d29..48693a050aa8 100644
--- a/net/tipc/node.h
+++ b/net/tipc/node.h
@@ -107,6 +107,7 @@ void tipc_node_subscribe(struct net *net, struct list_head *subscr, u32 addr);
void tipc_node_unsubscribe(struct net *net, struct list_head *subscr, u32 addr);
void tipc_node_broadcast(struct net *net, struct sk_buff *skb, int rc_dests);
int tipc_node_add_conn(struct net *net, u32 dnode, u32 port, u32 peer_port);
+int tipc_node_update_conn(struct net *net, u32 dnode, u32 port, u32 peer_port);
void tipc_node_remove_conn(struct net *net, u32 dnode, u32 port);
int tipc_node_get_mtu(struct net *net, u32 addr, u32 sel, bool connected);
bool tipc_node_is_up(struct net *net, u32 addr);
diff --git a/net/tipc/socket.c b/net/tipc/socket.c
index d5d70eb230b5..191579636351 100644
--- a/net/tipc/socket.c
+++ b/net/tipc/socket.c
@@ -1511,6 +1511,16 @@ static int __tipc_sendmsg(struct socket *sock, struct msghdr *m, size_t dlen)
return -ENOMEM;
}
+ /* Subscribe before sending SYN so node loss also aborts pending connects. */
+ if (syn) {
+ rc = tipc_node_add_conn(net, skaddr.node, tsk->portid, skaddr.ref);
+ if (rc) {
+ __skb_queue_purge(&pkts);
+ __skb_queue_purge(&sk->sk_write_queue);
+ return rc;
+ }
+ }
+
/* Send message */
trace_tipc_sk_sendmsg(sk, skb_peek(&pkts), TIPC_DUMP_SK_SNDQ, " ");
rc = tipc_node_xmit(net, &pkts, skaddr.node, tsk->portid);
@@ -1520,6 +1530,11 @@ static int __tipc_sendmsg(struct socket *sock, struct msghdr *m, size_t dlen)
rc = 0;
}
+ if (unlikely(syn && rc)) {
+ tipc_node_remove_conn(net, skaddr.node, tsk->portid);
+ __skb_queue_purge(&sk->sk_write_queue);
+ }
+
if (unlikely(syn && !rc)) {
tipc_set_sk_state(sk, TIPC_CONNECTING);
if (dlen && timeout) {
@@ -1674,8 +1689,10 @@ static void tipc_sk_finish_conn(struct tipc_sock *tsk, u32 peer_port,
msg_set_hdr_sz(msg, SHORT_H_SIZE);
sk_reset_timer(sk, &sk->sk_timer, jiffies + CONN_PROBING_INTV);
+ /* Active opens already subscribed before sending SYN. */
+ if (sk->sk_state != TIPC_CONNECTING)
+ tipc_node_add_conn(net, peer_node, tsk->portid, peer_port);
tipc_set_sk_state(sk, TIPC_ESTABLISHED);
- tipc_node_add_conn(net, peer_node, tsk->portid, peer_port);
tsk->max_pkt = tipc_node_get_mtu(net, peer_node, tsk->portid, true);
tsk->peer_caps = tipc_node_get_capabilities(net, peer_node);
tsk_set_nagle(tsk);
@@ -2199,12 +2216,14 @@ static bool tipc_sk_filter_connect(struct tipc_sock *tsk, struct sk_buff *skb,
struct net *net = sock_net(sk);
struct tipc_msg *hdr = buf_msg(skb);
bool con_msg = msg_connected(hdr);
+ bool reject_ack = false;
u32 pport = tsk_peer_port(tsk);
u32 pnode = tsk_peer_node(tsk);
u32 oport = msg_origport(hdr);
u32 onode = msg_orignode(hdr);
int err = msg_errcode(hdr);
unsigned long delay;
+ int rc;
if (unlikely(msg_mcast(hdr)))
return false;
@@ -2216,6 +2235,18 @@ static bool tipc_sk_filter_connect(struct tipc_sock *tsk, struct sk_buff *skb,
if (likely(con_msg)) {
if (err)
break;
+ rc = tipc_node_update_conn(net, pnode, tsk->portid, oport);
+ /* A named SYN may have been rerouted to another node. */
+ if (!rc && onode != pnode) {
+ tipc_node_remove_conn(net, pnode, tsk->portid);
+ rc = tipc_node_add_conn(net, onode, tsk->portid, oport);
+ }
+ if (rc) {
+ err = TIPC_ERR_NO_NODE;
+ /* Return a late ACK to its sender, including any data. */
+ reject_ack = true;
+ break;
+ }
tipc_sk_finish_conn(tsk, oport, onode);
msg_set_importance(&tsk->phdr, msg_importance(hdr));
/* ACK+ message with data is added to receive queue */
@@ -2285,9 +2316,10 @@ static bool tipc_sk_filter_connect(struct tipc_sock *tsk, struct sk_buff *skb,
}
/* Abort connection setup attempt */
tipc_set_sk_state(sk, TIPC_DISCONNECTING);
- sk->sk_err = ECONNREFUSED;
+ tipc_node_remove_conn(net, pnode, tsk->portid);
+ sk->sk_err = err == TIPC_ERR_NO_NODE ? EHOSTUNREACH : ECONNREFUSED;
sk->sk_state_change(sk);
- return true;
+ return !reject_ack;
}
/**
--
2.34.1
^ permalink raw reply [flat|nested] 4+ messages in thread* [RFC net-next 2/2] selftests: net: cover pending TIPC connects on peer node loss
2026-09-29 6:47 [RFC net-next 0/2] tipc: notify pending connects of peer node loss liushike
2026-09-29 6:47 ` [RFC net-next 1/2] tipc: abort pending connects when the peer node is lost liushike
@ 2026-09-29 6:47 ` liushike
2026-09-29 6:54 ` [RFC net-next 0/2] tipc: notify pending connects of " netdev-bot+sinfo
2 siblings, 0 replies; 4+ messages in thread
From: liushike @ 2026-09-29 6:47 UTC (permalink / raw)
To: netdev
Cc: jmaloy, tung.quang.nguyen, davem, edumazet, kuba, pabeni, horms,
shuah, tipc-discussion, linux-kselftest, linux-kernel
From: liushike <liushike@ruijie.com.cn>
Exercise a client waiting for a listening server to accept a connection
using two network namespaces and one or two Ethernet veth bearers.
Cover normal acceptance, last-link loss, survival with one remaining
link, loss of both links, rejection, connection-wait timeout followed by
node loss, nonblocking connects, and service-name addressing. Run each
case with SOCK_STREAM and SOCK_SEQPACKET.
Wait for actual unicast links at both peers; the always-UP broadcast
link does not establish peer reachability. Require node-loss completion
within two seconds and check EHOSTUNREACH. For asynchronous failure,
check POLLHUP with poll(), matching TIPC_DISCONNECTING semantics rather
than assuming the socket becomes writable. Save topology diagnostics
when a test raises an exception.
Signed-off-by: liushike <liushike@ruijie.com.cn>
---
MAINTAINERS | 1 +
tools/testing/selftests/net/Makefile | 1 +
tools/testing/selftests/net/config | 1 +
tools/testing/selftests/net/tipc_connect.py | 278 ++++++++++++++++++++
4 files changed, 281 insertions(+)
create mode 100755 tools/testing/selftests/net/tipc_connect.py
diff --git a/MAINTAINERS b/MAINTAINERS
index df8ab9b82402..4d595e60a22d 100644
--- a/MAINTAINERS
+++ b/MAINTAINERS
@@ -27561,6 +27561,7 @@ S: Maintained
W: http://tipc.sourceforge.net/
F: include/uapi/linux/tipc*.h
F: net/tipc/
+F: tools/testing/selftests/net/tipc_connect.py
TLAN NETWORK DRIVER
M: Samuel Chessman <chessman@tux.org>
diff --git a/tools/testing/selftests/net/Makefile b/tools/testing/selftests/net/Makefile
index 3ee3378f8b26..c453ff976386 100644
--- a/tools/testing/selftests/net/Makefile
+++ b/tools/testing/selftests/net/Makefile
@@ -118,6 +118,7 @@ TEST_PROGS := \
test_vxlan_vnifilter_notify.sh \
test_vxlan_vnifiltering.sh \
tfo_passive.sh \
+ tipc_connect.py \
traceroute.sh \
txtimestamp.sh \
udpgro.sh \
diff --git a/tools/testing/selftests/net/config b/tools/testing/selftests/net/config
index 30d5fcb09a83..725c8fbc5f03 100644
--- a/tools/testing/selftests/net/config
+++ b/tools/testing/selftests/net/config
@@ -128,6 +128,7 @@ CONFIG_TCP_CONG_DCTCP=y
CONFIG_TCP_MD5SIG=y
CONFIG_TEST_BLACKHOLE_DEV=m
CONFIG_TEST_BPF=m
+CONFIG_TIPC=m
CONFIG_TLS=m
CONFIG_TRACEPOINTS=y
CONFIG_TUN=y
diff --git a/tools/testing/selftests/net/tipc_connect.py b/tools/testing/selftests/net/tipc_connect.py
new file mode 100755
index 000000000000..02685b3ce8dd
--- /dev/null
+++ b/tools/testing/selftests/net/tipc_connect.py
@@ -0,0 +1,278 @@
+#!/usr/bin/env python3
+# SPDX-License-Identifier: GPL-2.0
+
+"""Exercise pending TIPC connects across bearer loss, before accept()."""
+
+import errno
+import os
+import select
+import shutil
+import socket
+import threading
+import time
+from contextlib import ExitStack, contextmanager
+
+from lib.py import (KsftNamedVariant, KsftSkipEx, NetNS, NetNSEnter,
+ cmd, ip, ksft_eq, ksft_exit, ksft_run, ksft_true,
+ ksft_variants)
+
+
+SOCKET_TYPES = [KsftNamedVariant("stream", socket.SOCK_STREAM),
+ KsftNamedVariant("seqpacket", socket.SOCK_SEQPACKET)]
+CONNECT_TIMEOUT = 10000
+RETURN_TIMEOUT = 2
+
+
+def tipc(ns, args):
+ return cmd("tipc " + args, ns=ns).stdout
+
+
+def unicast_links_up(output):
+ """Return distinct UP links, excluding the always-UP broadcast link."""
+ links = set()
+ for line in output.splitlines():
+ name, separator, state = line.strip().rpartition(": ")
+ if (separator and state.lower() == "up" and
+ not name.startswith("broadcast-link")):
+ links.add(name)
+ return links
+
+
+def wait_for_links(client_ns, server_ns, links):
+ deadline = time.monotonic() + 10
+ while True:
+ # These isolated namespaces contain only the two test peers. Count
+ # their unicast links, not the namespace's always-UP broadcast link.
+ outputs = [tipc(ns, "link list") for ns in (client_ns, server_ns)]
+ if all(len(unicast_links_up(output)) == links for output in outputs):
+ return
+ if time.monotonic() >= deadline:
+ raise TimeoutError(
+ f"Expected {links} UP unicast links at each peer; "
+ f"client link list: {outputs[0]!r}; "
+ f"server link list: {outputs[1]!r}")
+ time.sleep(0.05)
+
+
+@contextmanager
+def topology(sock_type, links=1):
+ if not hasattr(socket, "AF_TIPC"):
+ raise KsftSkipEx("Python lacks AF_TIPC support")
+ if os.geteuid() != 0:
+ raise KsftSkipEx("root privileges required for network namespaces")
+ for tool in ("ip", "tipc"):
+ if not shutil.which(tool):
+ raise KsftSkipEx(f"{tool} is required")
+ try:
+ with socket.socket(socket.AF_TIPC, sock_type):
+ pass
+ except OSError as error:
+ if error.errno in (errno.EAFNOSUPPORT, errno.EPROTONOSUPPORT):
+ raise KsftSkipEx("TIPC support is unavailable; load tipc") from error
+ raise
+
+ with ExitStack() as stack:
+ client_ns = stack.enter_context(NetNS())
+ server_ns = stack.enter_context(NetNS())
+ for ns, address in ((client_ns, "1.1.1"), (server_ns, "1.1.2")):
+ tipc(ns, "node set netid 4711")
+ tipc(ns, f"node set address {address}")
+ for index in range(links):
+ ip(f"link add c{index} type veth peer name s{index} "
+ f"netns {server_ns}", ns=client_ns)
+ for ns, dev in ((client_ns, f"c{index}"),
+ (server_ns, f"s{index}")):
+ ip(f"link set {dev} up", ns=ns)
+ tipc(ns, f"bearer enable media eth device {dev}")
+
+ wait_for_links(client_ns, server_ns, links)
+
+ with NetNSEnter(server_ns):
+ listener = stack.enter_context(socket.socket(socket.AF_TIPC,
+ sock_type))
+ listener.listen(8)
+ with NetNSEnter(client_ns):
+ client = stack.enter_context(socket.socket(socket.AF_TIPC,
+ sock_type))
+ client.setsockopt(socket.SOL_TIPC, socket.TIPC_CONN_TIMEOUT,
+ CONNECT_TIMEOUT)
+ try:
+ yield client_ns, client, listener
+ except Exception:
+ # Collect diagnostics before ExitStack destroys the topology.
+ for label, ns in (("client", client_ns), ("server", server_ns)):
+ for args in ("link list", "node list"):
+ try:
+ output = tipc(ns, args)
+ print(f"# {label} {args}: {output!r}")
+ except Exception as error:
+ print(f"# {label} {args} failed: {error}")
+ raise
+
+
+@contextmanager
+def pending_connect(client, listener, address=None):
+ result = []
+ if address is None:
+ address = listener.getsockname()
+
+ def connect():
+ result.append(client.connect_ex(address))
+
+ worker = threading.Thread(target=connect)
+ worker.start()
+ try:
+ # Readability proves SYN is queued on the listener, without accepting it.
+ if not select.select([listener], [], [], 5)[0]:
+ raise TimeoutError(f"SYN did not reach {address!r}: {result}")
+ ksft_eq(result, [], "connect must wait for accept")
+ yield worker, result
+ finally:
+ # Bound cleanup even on the unfixed kernel or a failed assertion.
+ if worker.is_alive():
+ try:
+ client.shutdown(socket.SHUT_RDWR)
+ except OSError:
+ pass
+ worker.join(CONNECT_TIMEOUT / 1000 + 1)
+
+
+def drop_link(ns, index):
+ tipc(ns, f"bearer disable media eth device c{index}")
+
+
+def check_node_loss(client):
+ # TIPC_DISCONNECTING reports POLLHUP, not POLLOUT. Unlike select()'s
+ # write set, poll() reports hangup even when only POLLOUT is requested.
+ poller = select.poll()
+ poller.register(client, select.POLLOUT)
+ events = poller.poll(int(RETURN_TIMEOUT * 1000))
+ ksft_true(events, "node loss must notify pending connect within 2s")
+ if events:
+ ksft_eq(events[0][0], client.fileno())
+ ksft_true(events[0][1] & select.POLLHUP,
+ f"node loss must report hangup, got {events!r}")
+ ksft_eq(client.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR),
+ errno.EHOSTUNREACH)
+
+
+@ksft_variants(SOCKET_TYPES)
+def normal_accept(sock_type):
+ with topology(sock_type) as (_, client, listener):
+ with pending_connect(client, listener) as (worker, result):
+ with listener.accept()[0] as accepted:
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [0])
+ if result == [0]:
+ accepted.settimeout(RETURN_TIMEOUT)
+ client.sendall(b"connected")
+ ksft_eq(accepted.recv(32), b"connected")
+
+
+@ksft_variants(SOCKET_TYPES)
+def last_link_down(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ with pending_connect(client, listener) as (worker, result):
+ drop_link(ns, 0)
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [errno.EHOSTUNREACH],
+ "node loss must abort connect before its 10s timeout")
+
+
+@ksft_variants(SOCKET_TYPES)
+def surviving_link(sock_type):
+ with topology(sock_type, links=2) as (ns, client, listener):
+ with pending_connect(client, listener) as (worker, result):
+ drop_link(ns, 0)
+ worker.join(0.2)
+ ksft_eq(result, [], "one remaining link must preserve connect")
+ with listener.accept()[0]:
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [0])
+
+
+@ksft_variants(SOCKET_TYPES)
+def both_links_down(sock_type):
+ with topology(sock_type, links=2) as (ns, client, listener):
+ with pending_connect(client, listener) as (worker, result):
+ drop_link(ns, 0)
+ worker.join(0.2)
+ ksft_eq(result, [])
+ drop_link(ns, 1)
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [errno.EHOSTUNREACH])
+
+
+@ksft_variants(SOCKET_TYPES)
+def rejected_connect(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ with pending_connect(client, listener) as (worker, result):
+ listener.close()
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [errno.ECONNREFUSED])
+ drop_link(ns, 0)
+ ksft_eq(client.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR), 0,
+ "node loss must not change a completed rejection")
+
+
+@ksft_variants(SOCKET_TYPES)
+def connect_timeout_then_node_loss(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ client.setsockopt(socket.SOL_TIPC, socket.TIPC_CONN_TIMEOUT, 500)
+ with pending_connect(client, listener) as (worker, result):
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [errno.ETIMEDOUT])
+ # A wait timeout leaves the handshake pending, just as EINPROGRESS.
+ drop_link(ns, 0)
+ check_node_loss(client)
+
+
+@ksft_variants(SOCKET_TYPES)
+def nonblocking_connect(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ client.setblocking(False)
+ ksft_eq(client.connect_ex(listener.getsockname()), errno.EINPROGRESS)
+ if not select.select([listener], [], [], 5)[0]:
+ raise TimeoutError("SYN did not reach listener")
+ drop_link(ns, 0)
+ check_node_loss(client)
+
+
+def publish_service(ns, listener):
+ service = 18888
+ listener.bind((socket.TIPC_ADDR_NAMESEQ, service, 1, 1,
+ socket.TIPC_CLUSTER_SCOPE))
+ deadline = time.monotonic() + 10
+ while str(service) not in tipc(ns, "nametable show").split():
+ if time.monotonic() >= deadline:
+ raise TimeoutError("TIPC service publication did not arrive")
+ time.sleep(0.05)
+ return socket.TIPC_ADDR_NAME, service, 1, 0
+
+
+@ksft_variants(SOCKET_TYPES)
+def named_accept(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ address = publish_service(ns, listener)
+ with pending_connect(client, listener, address) as (worker, result):
+ with listener.accept()[0]:
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [0])
+
+
+@ksft_variants(SOCKET_TYPES)
+def named_node_loss(sock_type):
+ with topology(sock_type) as (ns, client, listener):
+ address = publish_service(ns, listener)
+ with pending_connect(client, listener, address) as (worker, result):
+ drop_link(ns, 0)
+ worker.join(RETURN_TIMEOUT)
+ ksft_eq(result, [errno.EHOSTUNREACH])
+
+
+if __name__ == "__main__":
+ ksft_run(cases=[normal_accept, last_link_down, surviving_link,
+ both_links_down, rejected_connect,
+ connect_timeout_then_node_loss, nonblocking_connect,
+ named_accept, named_node_loss])
+ ksft_exit()
--
2.34.1
^ permalink raw reply [flat|nested] 4+ messages in thread