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, ¶ms); > > + 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