Re: [PATCH net-next v4 5/6] selftests: net: add multithread server support to iou-zcrx

Juanlu Herrero <[email protected]> Tue, 4 Aug 2026 14:54:39 -0500
Newsgroups gmane.linux.network
Message-ID <anJDFTKOcncNQUeL@jlhe0197-mac>
On Tue, Aug 04, 2026 at 09:45:05AM -0600, David Wei wrote:
> On 2026-07-29 15:18, Juanlu Herrero wrote:
> > Run the iou-zcrx server as N worker threads, each owning one receive
> > queue with its own io_uring and zero-copy receive (zcrx) ifq, so the
> > test can exercise multi-queue zero-copy receive.
> > 
> > The main thread owns the listening socket, accepts connections, and
> > dispatches each to the worker owning the queue it landed on by matching
> > SO_INCOMING_NAPI_ID against per-queue NAPI IDs.
> > 
> > Assisted-by: Claude:claude-opus-4-8
> > Signed-off-by: Juanlu Herrero <[email protected]>
> > ---
> >   .../testing/selftests/drivers/net/hw/Makefile |   5 +-
> >   .../selftests/drivers/net/hw/iou-zcrx.c       | 290 ++++++++++++++----
> >   2 files changed, 226 insertions(+), 69 deletions(-)
> > 
> > diff --git a/tools/testing/selftests/drivers/net/hw/iou-zcrx.c b/tools/testing/selftests/drivers/net/hw/iou-zcrx.c
> > index 7bc61f3b70ca6..16259129df46d 100644
> > --- a/tools/testing/selftests/drivers/net/hw/iou-zcrx.c
> > +++ b/tools/testing/selftests/drivers/net/hw/iou-zcrx.c
> [...]
> > @@ -322,28 +310,124 @@ static void server_loop(struct thread_ctx *ctx)
> >   	struct io_uring_cqe *cqe;
> >   	unsigned int count = 0;
> >   	unsigned int head;
> > -	int i, ret;
> >   	io_uring_submit_and_wait(&ctx->ring, 1);
> >   	io_uring_for_each_cqe(&ctx->ring, head, cqe) {
> > -		if (cqe->user_data == 1)
> > -			process_accept(ctx, cqe);
> > -		else if (cqe->user_data == 2)
> > -			process_recvzc(ctx, cqe);
> > -		else
> > -			error(1, 0, "unknown cqe");
> > +		process_recvzc(ctx, cqe, cqe->user_data);
> >   		count++;
> >   	}
> >   	io_uring_cq_advance(&ctx->ring, count);
> >   }
> > -static void run_server(void)
> > +static void *server_worker(void *arg)
> >   {
> > -	struct thread_ctx ctx = {};
> > -	unsigned int flags = 0;
> > -	int fd, enable, ret;
> > +	struct thread_ctx *ctx = arg;
> > +	struct io_uring_params params = { };
> 
> Reverse christmas tree.
> 
> >   	uint64_t tstop;
> > +	int i;
> > +
> > +	params.flags |= IORING_SETUP_COOP_TASKRUN;
> > +	params.flags |= IORING_SETUP_SINGLE_ISSUER;
> > +	params.flags |= IORING_SETUP_DEFER_TASKRUN;
> > +	params.flags |= IORING_SETUP_SUBMIT_ALL;
> > +	params.flags |= IORING_SETUP_CQE32;
> > +	params.flags |= IORING_SETUP_CQSIZE;
> > +	params.cq_entries = AREA_SIZE / page_size;
> > +
> > +	io_uring_queue_init_params(512, &ctx->ring, &params);
> > +	setup_zcrx(ctx);
> > +
> > +	if (cfg_dry_run)
> > +		return NULL;
> > +
> > +	{
> > +		uint64_t val = 1;
> > +
> > +		if (write(ctx->ready_fd, &val, sizeof(val)) != sizeof(val))
> > +			error(1, errno, "write(ready_fd)");
> > +		if (read(ctx->start_fd, &val, sizeof(val)) != sizeof(val))
> > +			error(1, errno, "read(start_fd)");
> 
> A pthread_barrier_t is more idiomatic and you can get rid of the epollfd
> in main thread entirely.
> 
> > +	}
> > +
> > +	for (i = 0; i < ctx->nr_conns; i++) {
> > +		if (cfg_oneshot)
> > +			add_recvzc_oneshot(ctx, i, page_size);
> > +		else
> > +			add_recvzc(ctx, i);
> > +	}
> > +
> > +	tstop = gettimeofday_ms() + 5000;
> > +	while (ctx->nr_conns > 0 && gettimeofday_ms() < tstop)
> > +		server_loop(ctx);
> > +
> > +	if (ctx->nr_conns != 0)
> > +		error(1, 0, "test failed: %d connections incomplete",
> > +		      ctx->nr_conns);
> > +
> > +	return NULL;
> > +}
> > +
> > +static int query_napi_id(unsigned int ifindex, int queue_id)
> > +{
> > +	struct netdev_queue_get_req *req;
> > +	struct netdev_queue_get_rsp *rsp;
> > +	struct ynl_error yerr;
> > +	struct ynl_sock *ys;
> > +	int napi_id;
> > +
> > +	ys = ynl_sock_create(&ynl_netdev_family, &yerr);
> > +	if (!ys)
> > +		error(1, 0, "ynl_sock_create: %s", yerr.msg);
> > +
> > +	req = netdev_queue_get_req_alloc();
> > +	netdev_queue_get_req_set_ifindex(req, ifindex);
> > +	netdev_queue_get_req_set_type(req, NETDEV_QUEUE_TYPE_RX);
> > +	netdev_queue_get_req_set_id(req, queue_id);
> > +
> > +	rsp = netdev_queue_get(ys, req);
> > +	if (!rsp)
> > +		error(1, 0, "netdev_queue_get(q=%d): %s", queue_id,
> > +		      ys->err.msg);
> > +	if (!rsp->_present.napi_id)
> > +		error(1, 0, "netdev_queue_get(q=%d): napi_id not present",
> > +		      queue_id);
> > +
> > +	napi_id = rsp->napi_id;
> > +
> > +	netdev_queue_get_req_free(req);
> > +	netdev_queue_get_rsp_free(rsp);
> > +	ynl_sock_destroy(ys);
> > +
> > +	return napi_id;
> > +}
> > +
> > +static int find_thread_by_napi(struct thread_ctx *ctxs, int napi_id)
> > +{
> > +	int i;
> > +
> > +	for (i = 0; i < cfg_num_threads; i++) {
> > +		if (ctxs[i].napi_id == napi_id)
> > +			return i;
> > +	}
> > +	return -1;
> > +}
> > +
> > +static void run_server(void)
> > +{
> > +	struct thread_ctx *ctxs;
> > +	struct epoll_event ev, out_ev;
> > +	pthread_t *threads;
> > +	unsigned int ifindex;
> > +	int conns_per_thread, total_conns, accepted = 0, connfd;
> > +	int fd, ret, enable, i;
> > +	int epfd, ready = 0;
> > +	uint64_t val;
> 
> Reverse christmas tree.
> 
> > +
> > +	ctxs = calloc(cfg_num_threads, sizeof(*ctxs));
> > +	threads = calloc(cfg_num_threads, sizeof(*threads));
> > +	if (!ctxs || !threads)
> > +		error(1, 0, "calloc()");
> >   	fd = socket(AF_INET6, SOCK_STREAM, 0);
> >   	if (fd == -1)
> > @@ -358,29 +442,105 @@ static void run_server(void)
> >   	if (ret < 0)
> >   		error(1, 0, "bind()");
> > -	flags |= IORING_SETUP_COOP_TASKRUN;
> > -	flags |= IORING_SETUP_SINGLE_ISSUER;
> > -	flags |= IORING_SETUP_DEFER_TASKRUN;
> > -	flags |= IORING_SETUP_SUBMIT_ALL;
> > -	flags |= IORING_SETUP_CQE32;
> > +	for (i = 0; i < cfg_num_threads; i++) {
> > +		ctxs[i].queue_id = cfg_queue_id + i;
> > +		ctxs[i].ready_fd = eventfd(0, 0);
> > +		ctxs[i].start_fd = eventfd(0, 0);
> > +	}
> > -	io_uring_queue_init(512, &ctx.ring, flags);
> > +	for (i = 0; i < cfg_num_threads; i++) {
> > +		ret = pthread_create(&threads[i], NULL,
> > +				     server_worker, &ctxs[i]);
> > +		if (ret)
> > +			error(1, ret, "pthread_create()");
> > +	}
> > -	setup_zcrx(&ctx);
> >   	if (cfg_dry_run)
> > -		return;
> > +		goto join;
> >   	if (listen(fd, 1024) < 0)
> >   		error(1, 0, "listen()");
> > -	add_accept(&ctx, fd);
> > +	epfd = epoll_create1(0);
> > +	if (epfd < 0)
> > +		error(1, errno, "epoll_create1()");
> > -	tstop = gettimeofday_ms() + 5000;
> > -	while (!ctx.stop && gettimeofday_ms() < tstop)
> > -		server_loop(&ctx);
> > +	for (i = 0; i < cfg_num_threads; i++) {
> > +		ev.events = EPOLLIN;
> > +		ev.data.fd = ctxs[i].ready_fd;
> > +		if (epoll_ctl(epfd, EPOLL_CTL_ADD,
> > +			      ctxs[i].ready_fd, &ev) < 0)
> > +			error(1, errno, "epoll_ctl()");
> > +	}
> > +
> > +	while (ready < cfg_num_threads) {
> > +		if (epoll_wait(epfd, &out_ev, 1, -1) < 0)
> > +			error(1, errno, "epoll_wait()");
> > +		if (read(out_ev.data.fd, &val, sizeof(val)) != sizeof(val))
> > +			error(1, errno, "read(ready_fd)");
> > +		ready++;
> > +	}
> > +
> > +	close(epfd);
> 
> All this can go if using a barrier.
> 
> > +
> > +	if (cfg_num_threads > 1) {
> > +		ifindex = if_nametoindex(cfg_ifname);
> > +		if (!ifindex)
> > +			error(1, 0, "bad interface name: %s", cfg_ifname);
> > +		for (i = 0; i < cfg_num_threads; i++)
> > +			ctxs[i].napi_id = query_napi_id(ifindex,
> > +							ctxs[i].queue_id);
> > +	}
> > +
> > +	conns_per_thread = cfg_num_threads > 1 ? CONNS_PER_THREAD : 1;
> > +	total_conns = conns_per_thread * cfg_num_threads;
> > +
> > +	while (accepted < total_conns) {
> > +		int idx;
> > +
> > +		connfd = accept(fd, NULL, NULL);
> > +		if (connfd < 0)
> > +			error(1, errno, "accept()");
> > +
> > +		if (cfg_num_threads > 1) {
> > +			int napi_id;
> > +			socklen_t len = sizeof(napi_id);
> > +
> > +			ret = getsockopt(connfd, SOL_SOCKET,
> > +					 SO_INCOMING_NAPI_ID,
> > +					 &napi_id, &len);
> > +			if (ret < 0)
> > +				error(1, errno,
> > +				      "getsockopt(SO_INCOMING_NAPI_ID)");
> > +
> > +			idx = find_thread_by_napi(ctxs, napi_id);
> 
> Make the helper take in a connfd and return the index directly.
> 
> > +			if (idx < 0)
> > +				error(1, 0, "unknown NAPI ID: %d",
> > +				      napi_id);
> > +		} else {
> > +			idx = 0;
> 
> Don't need the else, init idx to 0.

I will make these fixes in v5 of this patchset. Thanks for the review.

Best,
Juanlu