Re: [PATCH net-next v4 5/6] selftests: net: add multithread server support to iou-zcrx
David Wei <[email protected]> Tue, 4 Aug 2026 09:45:05 -0700
| Newsgroups | gmane.linux.network |
|---|---|
| Message-ID | <[email protected]> |
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.