mirror of https://lore.kernel.org/lkml/
 help / color / mirror / Atom feed
From: Jia Jia <physicalmtea@gmail.com>
To: stefanha@redhat.com, sgarzare@redhat.com, netdev@vger.kernel.org,
	virtualization@lists.linux.dev, kvm@vger.kernel.org
Cc: mst@redhat.com, jasowangio@gmail.com, eperezma@redhat.com,
	xuanzhuo@linux.alibaba.com, davem@davemloft.net,
	edumazet@kernel.org, kuba@kernel.org, pabeni@redhat.com,
	horms@kernel.org, linux-kernel@vger.kernel.org,
	bpf@vger.kernel.org, Jia Jia <physicalmtea@gmail.com>
Subject: [PATCH net-next v2 5/5] vsock/virtio: defer RX readable notifications until batch unlock
Date: Sat, 10 Oct 2026 22:22:47 +0800	[thread overview]
Message-ID: <20261010142247.99223-6-physicalmtea@gmail.com> (raw)
In-Reply-To: <20261010142247.99223-1-physicalmtea@gmail.com>

An RX batch can make a socket readable while the worker still owns its
socket lock. Waking a reader at that point can make it immediately wait
for the same lock.

For the callback installed by sock_init_data(), record when the existing
low-watermark condition is met and issue one notification after the batch
releases the lock. Save the initial sock_def_readable() callback when the
AF_VSOCK socket is created because virtio_transport_common can be built as
a module and cannot refer to that unexported symbol directly.

When sk_data_ready or sk_write_space no longer matches that initial
callback, pass queued skbs to the current callback one skb at a time. The
handoff stops if the callback is the default again or the queue does not
shrink. After 64 successful callbacks, queue another bounded handoff if
data remains. vsock_bpf_update_proto() schedules the same handoff so an
attach or detach outside the worker still drains skbs already queued.

Always deliver a coalesced write-space event to the saved default callback
before sampling the callback state. A current replacement is invoked by
the handoff as well, so a callback change cannot hide the native writer
wakeup or pending sockmap work.

Sampling at batch finish takes sk_callback_lock once. The per-packet fast
path still uses READ_ONCE() and does not take that lock. Control packets,
state transitions, transport changes and non-batched sockets keep their
existing behavior.

Signed-off-by: Jia Jia <physicalmtea@gmail.com>
---
 include/linux/virtio_vsock.h            |   1 +
 include/net/af_vsock.h                  |   8 +-
 net/vmw_vsock/af_vsock.c                | 103 ++++++++++++++++++++++++
 net/vmw_vsock/virtio_transport_common.c |  73 ++++++++++++++---
 net/vmw_vsock/vsock_bpf.c               |   2 +
 5 files changed, 174 insertions(+), 13 deletions(-)

diff --git a/include/linux/virtio_vsock.h b/include/linux/virtio_vsock.h
index 3a58120d9078..21a3460b940f 100644
--- a/include/linux/virtio_vsock.h
+++ b/include/linux/virtio_vsock.h
@@ -291,6 +291,7 @@ struct virtio_transport_rx_batch {
 	struct sockaddr_vm src;
 	struct sockaddr_vm dst;
 	bool write_space_pending;
+	bool data_ready_pending;
 };
 
 void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
diff --git a/include/net/af_vsock.h b/include/net/af_vsock.h
index 9d8ae62209ed..86d6e6ba3d9d 100644
--- a/include/net/af_vsock.h
+++ b/include/net/af_vsock.h
@@ -63,8 +63,11 @@ struct vsock_sock {
 	u32 peer_shutdown;
 	bool sent_request;
 	bool ignore_connecting_rst;
-	/* Initial callback, used to identify replacements. */
+	/* Initial callbacks, used to identify replacements. */
+	void (*default_data_ready)(struct sock *sk);
 	void (*default_write_space)(struct sock *sk);
+	struct delayed_work rx_cb_work;
+	bool rx_cb_retried;
 
 	/* Protected by lock_sock(sk) */
 	u64 buffer_size;
@@ -80,6 +83,9 @@ s64 vsock_stream_has_data(struct vsock_sock *vsk);
 s64 vsock_stream_has_space(struct vsock_sock *vsk);
 struct sock *vsock_create_connected(struct sock *parent);
 void vsock_data_ready(struct sock *sk);
+bool vsock_rx_cb_is_native(struct sock *sk);
+void vsock_rx_callback_kick(struct sock *sk);
+void vsock_rx_callback_handoff(struct sock *sk);
 
 /**** TRANSPORT ****/
 
diff --git a/net/vmw_vsock/af_vsock.c b/net/vmw_vsock/af_vsock.c
index fe7f49da6c4b..df6faee0bd37 100644
--- a/net/vmw_vsock/af_vsock.c
+++ b/net/vmw_vsock/af_vsock.c
@@ -931,6 +931,107 @@ static int __vsock_bind(struct sock *sk, struct sockaddr_vm *addr)
 	return retval;
 }
 
+#define VSOCK_RX_CB_HANDOFF_MAX 64
+
+static void vsock_rx_callback_queue(struct sock *sk, unsigned long delay)
+{
+	struct vsock_sock *vsk = vsock_sk(sk);
+
+	sock_hold(sk);
+	if (!schedule_delayed_work(&vsk->rx_cb_work, delay))
+		sock_put(sk);
+}
+
+void vsock_rx_callback_kick(struct sock *sk)
+{
+	vsock_rx_callback_queue(sk, 0);
+}
+
+/* Caller holds sk_callback_lock. */
+bool vsock_rx_cb_is_native(struct sock *sk)
+{
+	struct vsock_sock *vsk = vsock_sk(sk);
+
+	return READ_ONCE(sk->sk_prot) == sk->sk_prot_creator &&
+	       READ_ONCE(sk->sk_data_ready) == vsk->default_data_ready &&
+	       READ_ONCE(sk->sk_write_space) == vsk->default_write_space;
+}
+EXPORT_SYMBOL_GPL(vsock_rx_cb_is_native);
+
+/* Caller holds lock_sock(). */
+void vsock_rx_callback_handoff(struct sock *sk)
+{
+	struct vsock_sock *vsk = vsock_sk(sk);
+	void (*write_space)(struct sock *sk);
+	void (*data_ready)(struct sock *sk);
+	bool wrote = false;
+	s64 before, after;
+	int i;
+
+	if (sock_flag(sk, SOCK_DEAD) || !vsk->transport)
+		return;
+
+	read_lock_bh(&sk->sk_callback_lock);
+	data_ready = READ_ONCE(sk->sk_data_ready);
+	if (READ_ONCE(sk->sk_prot) != sk->sk_prot_creator &&
+	    data_ready == vsk->default_data_ready) {
+		read_unlock_bh(&sk->sk_callback_lock);
+		/*
+		 * start_verdict publishes sk_data_ready after the proto swap.
+		 * Retry once before treating this as a stable non-RX setup.
+		 */
+		if (!vsk->rx_cb_retried) {
+			vsk->rx_cb_retried = true;
+			vsock_rx_callback_queue(sk, 1);
+			return;
+		}
+		vsk->rx_cb_retried = false;
+		data_ready(sk);
+		return;
+	}
+	read_unlock_bh(&sk->sk_callback_lock);
+	vsk->rx_cb_retried = false;
+
+	for (i = 0; i < VSOCK_RX_CB_HANDOFF_MAX; i++) {
+		read_lock_bh(&sk->sk_callback_lock);
+		data_ready = READ_ONCE(sk->sk_data_ready);
+		write_space = READ_ONCE(sk->sk_write_space);
+		read_unlock_bh(&sk->sk_callback_lock);
+
+		if (!wrote && write_space != vsk->default_write_space) {
+			write_space(sk);
+			wrote = true;
+		}
+		if (data_ready == vsk->default_data_ready) {
+			data_ready(sk);
+			break;
+		}
+		before = vsock_stream_has_data(vsk);
+		if (before <= 0)
+			break;
+		data_ready(sk);
+		after = vsock_stream_has_data(vsk);
+		if (after >= before)
+			break;
+		if (i == VSOCK_RX_CB_HANDOFF_MAX - 1 && after > 0)
+			vsock_rx_callback_queue(sk, 0);
+	}
+}
+EXPORT_SYMBOL_GPL(vsock_rx_callback_handoff);
+
+static void vsock_rx_callback_work(struct work_struct *work)
+{
+	struct vsock_sock *vsk = container_of(work, struct vsock_sock,
+					      rx_cb_work.work);
+	struct sock *sk = sk_vsock(vsk);
+
+	lock_sock(sk);
+	if (!sock_flag(sk, SOCK_DEAD))
+		vsock_rx_callback_handoff(sk);
+	release_sock(sk);
+	sock_put(sk);
+}
+
 static void vsock_connect_timeout(struct work_struct *work);
 
 static struct sock *__vsock_create(struct net *net,
@@ -958,6 +1059,7 @@ static struct sock *__vsock_create(struct net *net,
 		sk->sk_type = type;
 
 	vsk = vsock_sk(sk);
+	vsk->default_data_ready = sk->sk_data_ready;
 	vsk->default_write_space = sk->sk_write_space;
 	vsock_addr_init(&vsk->local_addr, VMADDR_CID_ANY, VMADDR_PORT_ANY);
 	vsock_addr_init(&vsk->remote_addr, VMADDR_CID_ANY, VMADDR_PORT_ANY);
@@ -976,6 +1078,7 @@ static struct sock *__vsock_create(struct net *net,
 	WRITE_ONCE(vsk->peer_shutdown, 0);
 	INIT_DELAYED_WORK(&vsk->connect_work, vsock_connect_timeout);
 	INIT_DELAYED_WORK(&vsk->pending_work, vsock_pending_work);
+	INIT_DELAYED_WORK(&vsk->rx_cb_work, vsock_rx_callback_work);
 
 	psk = parent ? vsock_sk(parent) : NULL;
 	if (parent) {
diff --git a/net/vmw_vsock/virtio_transport_common.c b/net/vmw_vsock/virtio_transport_common.c
index f0398ae2200d..5499cd1d8a25 100644
--- a/net/vmw_vsock/virtio_transport_common.c
+++ b/net/vmw_vsock/virtio_transport_common.c
@@ -1583,10 +1583,12 @@ out:
 
 static int
 virtio_transport_recv_connected(struct sock *sk,
-				struct sk_buff *skb)
+				struct sk_buff *skb,
+				bool *data_ready_pending)
 {
 	struct virtio_vsock_hdr *hdr = virtio_vsock_hdr(skb);
 	struct vsock_sock *vsk = vsock_sk(sk);
+	void (*data_ready)(struct sock *sk);
 	int err = 0;
 
 	switch (le16_to_cpu(hdr->op)) {
@@ -1602,7 +1604,23 @@ virtio_transport_recv_connected(struct sock *sk,
 			vsock_remove_sock(vsk);
 			break;
 		}
-		vsock_data_ready(sk);
+		if (!data_ready_pending) {
+			vsock_data_ready(sk);
+		} else {
+			data_ready = READ_ONCE(sk->sk_data_ready);
+			if (data_ready == vsk->default_data_ready) {
+				if (!*data_ready_pending &&
+				    (vsock_stream_has_data(vsk) >= sk->sk_rcvlowat ||
+				     sock_flag(sk, SOCK_DONE)))
+					*data_ready_pending = true;
+			} else {
+				/* Use the callback seen for this packet. */
+				*data_ready_pending = false;
+				if (vsock_stream_has_data(vsk) >= sk->sk_rcvlowat ||
+				    sock_flag(sk, SOCK_DONE))
+					data_ready(sk);
+			}
+		}
 		return err;
 	case VIRTIO_VSOCK_OP_CREDIT_REQUEST:
 		virtio_transport_send_credit_update(vsk);
@@ -1836,6 +1854,7 @@ struct virtio_transport_rx_pkt_ctx {
 	const struct sockaddr_vm *dst;
 	bool *batchable;
 	struct virtio_transport_rx_batch *batch;
+	bool defer_data_ready;
 };
 
 static bool
@@ -1909,7 +1928,9 @@ virtio_transport_recv_pkt_locked(struct virtio_transport *t,
 		kfree_skb(skb);
 		break;
 	case TCP_ESTABLISHED:
-		virtio_transport_recv_connected(sk, skb);
+		virtio_transport_recv_connected(sk, skb,
+						ctx->defer_data_ready && ctx->batch ?
+						&ctx->batch->data_ready_pending : NULL);
 		break;
 	case TCP_CLOSING:
 		virtio_transport_recv_disconnecting(sk, skb);
@@ -1978,33 +1999,45 @@ EXPORT_SYMBOL_GPL(virtio_transport_recv_pkt);
  * Finish the RX batch.
  * For a non-empty batch, the caller must hold the socket lock acquired with
  * lock_sock(). This function releases the lock and the batch's lookup
- * reference.
- * An empty batch is a no-op.
+ * reference. An empty batch is a no-op.
  */
 void virtio_transport_rx_batch_finish(struct virtio_transport_rx_batch *batch)
 {
 	bool write_space_pending = batch->write_space_pending;
-	void (*write_space)(struct sock *sk);
+	bool data_ready_pending = batch->data_ready_pending;
 	struct sock *sk = batch->sk;
 	struct vsock_sock *vsk;
+	bool native;
 
 	batch->sk = NULL;
 	batch->pkts = 0;
 	batch->bytes = 0;
 	batch->net = NULL;
 	batch->write_space_pending = false;
+	batch->data_ready_pending = false;
 
 	if (!sk)
 		return;
 
-	if (write_space_pending) {
-		vsk = vsock_sk(sk);
+	vsk = vsock_sk(sk);
+	if (write_space_pending)
 		vsk->default_write_space(sk);
-		write_space = READ_ONCE(sk->sk_write_space);
-		if (write_space != vsk->default_write_space)
-			write_space(sk);
+
+	read_lock_bh(&sk->sk_callback_lock);
+	native = vsock_rx_cb_is_native(sk);
+	read_unlock_bh(&sk->sk_callback_lock);
+
+	if (native) {
+		release_sock(sk);
+		if (data_ready_pending)
+			vsk->default_data_ready(sk);
+		sock_put(sk);
+		return;
 	}
 
+	if (write_space_pending || data_ready_pending)
+		vsock_rx_callback_handoff(sk);
+
 	release_sock(sk);
 	sock_put(sk);
 }
@@ -2015,9 +2048,9 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 				     struct virtio_transport_rx_batch *batch)
 {
 	struct virtio_vsock_hdr *hdr = virtio_vsock_hdr(skb);
+	bool batchable, defer_data_ready, start_batch;
 	struct virtio_transport_rx_pkt_ctx ctx;
 	struct sockaddr_vm src, dst;
-	bool batchable, start_batch;
 	struct sock *sk;
 	bool free_pkt;
 
@@ -2038,6 +2071,17 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 		    vsock_addr_equals_addr(&batch->dst, &dst) &&
 		    virtio_transport_recv_pkt_batchable(t, batch->sk)) {
 			sk = batch->sk;
+			defer_data_ready = READ_ONCE(sk->sk_data_ready) ==
+					   vsock_sk(sk)->default_data_ready;
+			if (unlikely(batch->data_ready_pending && !defer_data_ready)) {
+				/*
+				 * BPF map updates can replace callbacks while lock_sock()
+				 * remains held. Finish the old notification first.
+				 */
+				virtio_transport_rx_batch_finish(batch);
+				goto lookup;
+			}
+
 			if (!skb_set_owner_sk_safe(skb, sk)) {
 				WARN_ONCE(1, "receiving vsock socket has sk_refcnt == 0\n");
 				virtio_transport_rx_batch_finish(batch);
@@ -2051,6 +2095,7 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 				.dst = &dst,
 				.batchable = &batchable,
 				.batch = batch,
+				.defer_data_ready = defer_data_ready,
 			};
 			free_pkt = virtio_transport_recv_pkt_locked(t, skb, sk, &ctx);
 
@@ -2064,6 +2109,7 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 		virtio_transport_rx_batch_finish(batch);
 	}
 
+lookup:
 	sk = virtio_transport_recv_pkt_find_socket(skb, &src, &dst, net);
 	if (!sk) {
 		virtio_transport_rx_batch_finish(batch);
@@ -2089,6 +2135,8 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 
 	if (start_batch)
 		batch->sk = sk;
+	defer_data_ready = start_batch &&
+		READ_ONCE(sk->sk_data_ready) == vsock_sk(sk)->default_data_ready;
 
 	ctx = (struct virtio_transport_rx_pkt_ctx) {
 		.net = net,
@@ -2096,6 +2144,7 @@ void virtio_transport_recv_pkt_batch(struct virtio_transport *t,
 		.dst = &dst,
 		.batchable = start_batch ? &batchable : NULL,
 		.batch = start_batch ? batch : NULL,
+		.defer_data_ready = defer_data_ready,
 	};
 	free_pkt = virtio_transport_recv_pkt_locked(t, skb, sk, &ctx);
 	if (start_batch && batchable) {
diff --git a/net/vmw_vsock/vsock_bpf.c b/net/vmw_vsock/vsock_bpf.c
index 9049d2648646..7770a6c51dbe 100644
--- a/net/vmw_vsock/vsock_bpf.c
+++ b/net/vmw_vsock/vsock_bpf.c
@@ -154,6 +154,7 @@ int vsock_bpf_update_proto(struct sock *sk, struct sk_psock *psock, bool restore
 	if (restore) {
 		sk->sk_write_space = psock->saved_write_space;
 		sock_replace_proto(sk, psock->sk_proto);
+		vsock_rx_callback_kick(sk);
 		return 0;
 	}
 
@@ -166,6 +167,7 @@ int vsock_bpf_update_proto(struct sock *sk, struct sk_psock *psock, bool restore
 
 	vsock_bpf_check_needs_rebuild(psock->sk_proto);
 	sock_replace_proto(sk, &vsock_bpf_prot);
+	vsock_rx_callback_kick(sk);
 	return 0;
 }
 
-- 
2.34.1


  parent reply	other threads:[~2026-10-10 14:23 UTC|newest]

Thread overview: 9+ messages / expand[flat|nested]  mbox.gz  Atom feed  top
2026-10-10 14:22 [PATCH net-next v2 0/5] vsock/virtio: reduce RX per-packet socket overhead Jia Jia
2026-10-10 14:22 ` [PATCH net-next v2 1/5] vsock/virtio: split socket lookup from locked RX processing Jia Jia
2026-10-11 14:25   ` netdev-bot+sashiko
2026-10-10 14:22 ` [PATCH net-next v2 2/5] vsock/virtio: amortize RX socket locking for stream packets Jia Jia
2026-10-11 14:25   ` netdev-bot+sashiko
2026-10-10 14:22 ` [PATCH net-next v2 3/5] vsock/virtio: reuse same-flow socket lookup in RX batches Jia Jia
2026-10-10 14:22 ` [PATCH net-next v2 4/5] vsock/virtio: coalesce RX write-space notifications in lock batches Jia Jia
2026-10-10 14:22 ` Jia Jia [this message]
2026-10-11 14:25   ` [PATCH net-next v2 5/5] vsock/virtio: defer RX readable notifications until batch unlock netdev-bot+sashiko

Reply instructions:

You may reply publicly to this message via plain-text email
using any one of the following methods:

* Save the following mbox file, import it into your mail client,
  and reply-to-all from there: mbox

  Avoid top-posting and favor interleaved quoting:
  https://en.wikipedia.org/wiki/Posting_style#Interleaved_style

* Reply using the --to, --cc, and --in-reply-to
  switches of git-send-email(1):

  git send-email \
    --in-reply-to=20261010142247.99223-6-physicalmtea@gmail.com \
    --to=physicalmtea@gmail.com \
    --cc=bpf@vger.kernel.org \
    --cc=davem@davemloft.net \
    --cc=edumazet@kernel.org \
    --cc=eperezma@redhat.com \
    --cc=horms@kernel.org \
    --cc=jasowangio@gmail.com \
    --cc=kuba@kernel.org \
    --cc=kvm@vger.kernel.org \
    --cc=linux-kernel@vger.kernel.org \
    --cc=mst@redhat.com \
    --cc=netdev@vger.kernel.org \
    --cc=pabeni@redhat.com \
    --cc=sgarzare@redhat.com \
    --cc=stefanha@redhat.com \
    --cc=virtualization@lists.linux.dev \
    --cc=xuanzhuo@linux.alibaba.com \
    /path/to/YOUR_REPLY

  https://kernel.org/pub/software/scm/git/docs/git-send-email.html

* If your mail client supports setting the In-Reply-To header
  via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox

all inboxes | Powered by JetHome®