Linux block layer
 help / color / mirror / Atom feed
From: Ming Lei <tom.leiming@gmail.com>
To: linux-block@vger.kernel.org
Cc: Ming Lei <tom.leiming@gmail.com>, Jens Axboe <axboe@kernel.dk>,
	Caleb Sander Mateos <csander@purestorage.com>,
	Josef Bacik <josef@toxicpanda.com>
Subject: [PATCH 8/8] selftests: ublk: add test for going live over canceled io commands
Date: Thu,  1 Oct 2026 07:54:22 -0500	[thread overview]
Message-ID: <20261001125422.1364260-9-tom.leiming@gmail.com> (raw)
In-Reply-To: <20261001125422.1364260-1-tom.leiming@gmail.com>

Add ublk_cancel_ready, a small liburing program which sends control
commands with kublk's helpers from ctrl.c, and test_generic_18.sh,
which runs its modes:

  mode               what it does                     expected
  stop_start         STOP_DEV before START_DEV        START_DEV -EBUSY
  partial_fetch      task A fetches tags 0-2 and      START_DEV -ENODEV, or
                     dies, then task B fetches tag 3  live and reads -EIO
  recovery           in user recovery, queue 0's      END_USER_RECOVERY
                     task dies before queue 1 is      -ENODEV, or live and
                     ready                            reads -EIO
  stop_restart       STOP_DEV on a new device, then   START_DEV 0,
                     a server starts it               reads succeed
  stop_attached      STOP_DEV after open, before      FETCH ABORT and
                     FETCH, then a new server         START_DEV -EBUSY,
                                                      then a new one works
  stop_live_restart  STOP_DEV on a live device, then  FETCH and START_DEV
                     a new server starts it           0, reads succeed
  race_start         STOP_DEV and START_DEV at the    START_DEV 0 and
                     same time, 50 times              reads complete,
                                                      or -EBUSY
  race_fetch         STOP_DEV while a server opens    START_DEV 0 and
                     and fetches, then START_DEV,     reads succeed,
                     50 times                         or -EBUSY
  race_async_fetch   IOSQE_ASYNC FETCH while STOP_DEV no crash (KASAN
                     cancels, then close the ring     finds the bug)

If the device goes live, the test reads it, so an unfixed kernel shows
the oops in ublk_queue_cmd(). In stop_live_restart, an unfixed kernel
doesn't reset the FETCH round in ublk_ch_release() once the disk is
gone, so the new server's FETCH gets -EBUSY.

A dying task is a child process which fetches and then calls exec().
exec() cancels the task's uring_cmds before it returns, so the cancel
is done once the child is reaped, with no sleep. kublk can't be used:
its failure injection kills the whole process, which closes
/dev/ublkcN and starts a new round.

Link: https://lore.kernel.org/linux-block/20260928-b4-ublk-cancel-stop-v1-0-4a4360232a46@toxicpanda.com/
Signed-off-by: Ming Lei <tom.leiming@gmail.com>
Assisted-by: LLM
---
 tools/testing/selftests/ublk/.gitignore       |   1 +
 tools/testing/selftests/ublk/Makefile         |   6 +-
 .../testing/selftests/ublk/test_generic_18.sh |  37 +
 .../selftests/ublk/ublk_cancel_ready.c        | 984 ++++++++++++++++++
 4 files changed, 1026 insertions(+), 2 deletions(-)
 create mode 100755 tools/testing/selftests/ublk/test_generic_18.sh
 create mode 100644 tools/testing/selftests/ublk/ublk_cancel_ready.c

diff --git a/tools/testing/selftests/ublk/.gitignore b/tools/testing/selftests/ublk/.gitignore
index e17bd28f27e0..d8e93fef7fcd 100644
--- a/tools/testing/selftests/ublk/.gitignore
+++ b/tools/testing/selftests/ublk/.gitignore
@@ -3,3 +3,4 @@
 /tools
 kublk
 metadata_size
+ublk_cancel_ready
diff --git a/tools/testing/selftests/ublk/Makefile b/tools/testing/selftests/ublk/Makefile
index 37883e9d50ec..fc64b8f02833 100644
--- a/tools/testing/selftests/ublk/Makefile
+++ b/tools/testing/selftests/ublk/Makefile
@@ -19,6 +19,7 @@ TEST_PROGS += test_generic_12.sh
 TEST_PROGS += test_generic_13.sh
 TEST_PROGS += test_generic_16.sh
 TEST_PROGS += test_generic_17.sh
+TEST_PROGS += test_generic_18.sh
 
 TEST_PROGS += test_batch_01.sh
 TEST_PROGS += test_batch_02.sh
@@ -76,13 +77,14 @@ TEST_FILES := settings
 TEST_FILES += test_common.sh
 TEST_FILES += trace
 
-TEST_GEN_PROGS_EXTENDED = kublk metadata_size
-STANDALONE_UTILS := metadata_size.c
+TEST_GEN_PROGS_EXTENDED = kublk metadata_size ublk_cancel_ready
+STANDALONE_UTILS := metadata_size.c ublk_cancel_ready.c
 
 LOCAL_HDRS += $(wildcard *.h)
 include ../lib.mk
 
 $(OUTPUT)/kublk: $(filter-out $(STANDALONE_UTILS),$(wildcard *.c))
+$(OUTPUT)/ublk_cancel_ready: ublk_cancel_ready.c ctrl.c
 
 check:
 	shellcheck -x -f gcc *.sh
diff --git a/tools/testing/selftests/ublk/test_generic_18.sh b/tools/testing/selftests/ublk/test_generic_18.sh
new file mode 100755
index 000000000000..223944dba149
--- /dev/null
+++ b/tools/testing/selftests/ublk/test_generic_18.sh
@@ -0,0 +1,37 @@
+#!/bin/bash
+# SPDX-License-Identifier: GPL-2.0
+
+. "$(cd "$(dirname "$0")" && pwd)"/test_common.sh
+
+ERR_CODE=0
+CANCEL_PROG="$(_ublk_test_top_dir)/ublk_cancel_ready"
+
+_prep_test "generic" "start device over canceled io commands"
+
+# the modes are described in ublk_cancel_ready.c
+for mode in stop_start partial_fetch recovery stop_restart stop_attached \
+		stop_live_restart race_start race_fetch race_async_fetch; do
+	dmesg_before=$(dmesg | wc -l)
+	timeout 60 "$CANCEL_PROG" "$mode" > "$UBLK_TMP" 2>&1
+	res=$?
+	msg=""
+
+	if dmesg | tail -n +"$((dmesg_before + 1))" | \
+			grep -q -e "BUG:" -e "Oops" -e "WARNING:"; then
+		msg="$mode: kernel oops/warning"
+		ERR_CODE=255
+	elif [ "$res" -eq "$UBLK_SKIP_CODE" ]; then
+		[ "$ERR_CODE" -eq 0 ] && ERR_CODE=$UBLK_SKIP_CODE
+	elif [ "$res" -ne 0 ]; then
+		msg="$mode: failed ($res)"
+		ERR_CODE=255
+	fi
+	# the output once: on failure, or always when not quiet
+	[ -n "$msg" ] && echo "$msg"
+	if [ -n "$msg" ] || [ "$UBLK_TEST_QUIET" -eq 0 ]; then
+		cat "$UBLK_TMP"
+	fi
+done
+
+_cleanup_test
+_show_result $TID $ERR_CODE
diff --git a/tools/testing/selftests/ublk/ublk_cancel_ready.c b/tools/testing/selftests/ublk/ublk_cancel_ready.c
new file mode 100644
index 000000000000..0ed0d01fa74a
--- /dev/null
+++ b/tools/testing/selftests/ublk/ublk_cancel_ready.c
@@ -0,0 +1,984 @@
+// SPDX-License-Identifier: GPL-2.0
+/*
+ * Bring a ublk device live while some of its fetched io commands are
+ * canceled.
+ *
+ * A cancel completes a fetched command and clears io->cmd, but the io
+ * still counts as ready. Only ubq->canceling keeps ublk_queue_rq() away
+ * from the NULL io->cmd. Each mode below loses ->canceling in another way:
+ *
+ * stop_start:    fetch every tag, STOP_DEV before START_DEV, START_DEV.
+ *                STOP_DEV cancels the commands of the attached server,
+ *                so START_DEV must get -EBUSY until that server is gone.
+ *                An unfixed kernel starts the device over them.
+ *
+ * partial_fetch: task A fetches tags 0..depth-2 and dies, so its
+ *                commands are canceled. Task B fetches the last tag, and
+ *                an unfixed kernel clears ->canceling when the queue gets
+ *                ready. /dev/ublkcN stays open all the time, so
+ *                ublk_ch_release() never resets the queue.
+ *
+ * recovery:      a UBLK_F_USER_RECOVERY device with two queues loses its
+ *                server. During recovery task Q0 fetches queue 0, which
+ *                clears q0->canceling, then dies. An unfixed kernel still
+ *                has ub->canceling set because queue 1 is not ready, so
+ *                ublk_start_cancel() does not mark queue 0 again. Reads
+ *                are issued on a CPU mapped to queue 0.
+ *
+ * In partial_fetch and recovery, START_DEV / END_USER_RECOVERY may refuse
+ * with -ENODEV, or bring the device live with its queue still canceling:
+ * then every read has to complete, with -EIO. An unfixed kernel oopses
+ * in ublk_queue_cmd() on a NULL io->cmd.
+ *
+ * stop_restart:  control mode. STOP_DEV on a new device, before any
+ *                server opened it, takes no command. A server started
+ *                afterwards has to start the device and serve I/O.
+ *
+ * stop_attached: STOP_DEV on a device whose server opened it but fetched
+ *                nothing stops that server too: FETCH gets ABORT and
+ *                START_DEV -EBUSY, until a new server opens the device.
+ *
+ * stop_live_restart: STOP_DEV on a live device; once its server is gone,
+ *                a new server has to fetch and start it again. An unfixed
+ *                ublk_ch_release() skips the reset once the disk is gone,
+ *                so the new FETCH gets -EBUSY.
+ *
+ * race_start:    STOP_DEV and START_DEV at the same time on a device whose
+ *                server fetched every tag. START_DEV either wins, and
+ *                reads complete (they fail once STOP_DEV removes the
+ *                disk), or gets -EBUSY. An unfixed kernel cancels the
+ *                commands of the live disk.
+ *
+ * race_fetch:    STOP_DEV while a server opens the device and fetches,
+ *                then START_DEV. If STOP_DEV came before the open, it must
+ *                not take any command, so START_DEV works and every read
+ *                succeeds; otherwise START_DEV gets -EBUSY. An unfixed
+ *                kernel takes commands fetched after its unlock and goes
+ *                live over them.
+ *
+ * race_async_fetch: FETCH with IOSQE_ASYNC while STOP_DEV cancels,
+ *                then close the ring. An unfixed FETCH marks its command
+ *                cancelable only after publishing it and dropping
+ *                ub->mutex; a cancel in between completes a command which
+ *                is then put on io_uring's cancelable list. Closing the
+ *                ring walks that list: KASAN reports a use after free.
+ *                Needs a KASAN kernel to see the bug; otherwise it must
+ *                just not crash.
+ *
+ * A dying task is a child process which fetches and then calls exec().
+ * exec() cancels the task's uring_cmds before it returns, so the cancel
+ * is done once the child is reaped. A thread exit does not cancel them,
+ * and closing the ring cancels them later, from io_ring_exit_work().
+ */
+#include <sched.h>
+
+#include "kublk.h"
+#include "../kselftest.h"
+
+#define NR_QUEUES	2
+#define DEPTH		4
+#define BUF_SIZE	(64 << 10)
+#define DEV_SECTORS	(64 << 11)	/* 64MB */
+#define NR_READS	(DEPTH * 2)
+#define SERVE_DELAY_US	200000
+#define RACE_LOOPS	50
+#define ASYNC_LOOPS	300
+
+/* one control handle per thread, the race modes send commands from two */
+static __thread struct ublk_dev *ctrl_dev;
+static int cdev_fd = -1;
+static int dev_id = -1;
+static int nr_queues;
+static void *bufs[NR_QUEUES][DEPTH];
+static const struct ublksrv_io_desc *iods[NR_QUEUES];
+static size_t iods_len;
+
+static struct io_uring_sqe *get_sqe(struct io_uring *ring)
+{
+	struct io_uring_sqe *sqe = io_uring_get_sqe(ring);
+
+	if (!sqe) {
+		fprintf(stderr, "out of sqes\n");
+		exit(KSFT_FAIL);
+	}
+	return sqe;
+}
+
+static struct ublk_dev *ctrl(void)
+{
+	if (!ctrl_dev) {
+		ctrl_dev = ublk_ctrl_init();
+		if (!ctrl_dev) {
+			fprintf(stderr, "ublk_ctrl_init failed\n");
+			exit(KSFT_FAIL);
+		}
+	}
+	ctrl_dev->dev_info.dev_id = dev_id;
+	return ctrl_dev;
+}
+
+static void ctrl_put(void)
+{
+	if (ctrl_dev)
+		ublk_ctrl_deinit(ctrl_dev);
+	ctrl_dev = NULL;
+}
+
+static int dev_state(void)
+{
+	struct ublk_dev *dev = ctrl();
+	int ret = ublk_ctrl_get_info(dev);
+
+	return ret ? ret : dev->dev_info.state;
+}
+
+static int open_cdev(void)
+{
+	size_t max_len = UBLK_MAX_QUEUE_DEPTH * sizeof(struct ublksrv_io_desc);
+	int pg = getpagesize();
+	char path[64];
+
+	snprintf(path, sizeof(path), "/dev/ublkc%d", dev_id);
+	for (int i = 0; i < 100 && cdev_fd < 0; i++) {
+		cdev_fd = open(path, O_RDWR);
+		if (cdev_fd < 0)
+			usleep(50000);
+	}
+	if (cdev_fd < 0)
+		return -errno;
+
+	/* queue q's descriptors start at q * the size for the max depth */
+	max_len = (max_len + pg - 1) & ~(size_t)(pg - 1);
+	iods_len = (DEPTH * sizeof(struct ublksrv_io_desc) + pg - 1) &
+		~(size_t)(pg - 1);
+	for (int q = 0; q < nr_queues; q++) {
+		void *p = mmap(NULL, iods_len, PROT_READ,
+			       MAP_SHARED | MAP_POPULATE, cdev_fd,
+			       UBLKSRV_CMD_BUF_OFFSET + q * max_len);
+
+		if (p == MAP_FAILED)
+			return -errno;
+		iods[q] = p;
+	}
+	return 0;
+}
+
+/* the last reference to /dev/ublkcN runs ublk_ch_release() */
+static void close_cdev(void)
+{
+	for (int q = 0; q < nr_queues; q++) {
+		if (iods[q])
+			munmap((void *)iods[q], iods_len);
+		iods[q] = NULL;
+	}
+	if (cdev_fd >= 0)
+		close(cdev_fd);
+	cdev_fd = -1;
+}
+
+/* ADD_DEV and SET_PARAMS, without opening /dev/ublkcN */
+static int add_dev_noopen(int queues, __u64 flags)
+{
+	struct ublk_dev *dev = ctrl();
+	struct ublksrv_ctrl_dev_info info = {
+		.nr_hw_queues	= queues,
+		.queue_depth	= DEPTH,
+		.max_io_buf_bytes = BUF_SIZE,
+		.dev_id		= -1,
+		.flags		= UBLK_F_NO_AUTO_PART_SCAN | flags,
+	};
+	struct ublk_params p = {
+		.types	= UBLK_PARAM_TYPE_BASIC,
+		.basic	= {
+			.logical_bs_shift	= 9,
+			.physical_bs_shift	= 12,
+			.io_opt_shift		= 12,
+			.io_min_shift		= 9,
+			.max_sectors		= BUF_SIZE >> 9,
+			.dev_sectors		= DEV_SECTORS,
+		},
+	};
+	int ret;
+
+	nr_queues = queues;
+	dev->dev_info = info;
+	ret = ublk_ctrl_add_dev(dev);
+	if (ret)
+		return ret;
+	dev_id = dev->dev_info.dev_id;
+
+	ret = ublk_ctrl_set_params(ctrl(), &p);
+	if (ret)
+		return ret;
+
+	for (int q = 0; q < queues; q++)
+		for (int i = 0; i < DEPTH; i++)
+			if (!bufs[q][i] &&
+			    posix_memalign(&bufs[q][i], getpagesize(), BUF_SIZE))
+				return -ENOMEM;
+	return 0;
+}
+
+static int add_dev(int queues, __u64 flags)
+{
+	return add_dev_noopen(queues, flags) ?: open_cdev();
+}
+
+static void cleanup(void)
+{
+	if (dev_id < 0)
+		return;
+	ublk_ctrl_stop_dev(ctrl());
+	close_cdev();
+	ublk_ctrl_del_dev(ctrl());
+	dev_id = -1;
+}
+
+/* race_async_fetch sets IOSQE_ASYNC on the io commands */
+static int io_cmd_sqe_flags;
+
+static void queue_io_cmd(struct io_uring *ring, __u32 op, int q, int tag,
+			 int res)
+{
+	struct io_uring_sqe *sqe = get_sqe(ring);
+	struct ublksrv_io_cmd *cmd = (struct ublksrv_io_cmd *)sqe->cmd;
+
+	memset(sqe, 0, sizeof(*sqe));
+	sqe->fd = cdev_fd;
+	sqe->opcode = IORING_OP_URING_CMD;
+	sqe->flags = io_cmd_sqe_flags;
+	ublk_set_sqe_cmd_op(sqe, op);
+	cmd->q_id = q;
+	cmd->tag = tag;
+	cmd->result = res;
+	cmd->addr = (__u64)(uintptr_t)bufs[q][tag];
+	io_uring_sqe_set_data64(sqe, (q << 16) | tag);
+}
+
+struct async_arg {
+	pthread_barrier_t go, stopped;
+};
+
+/*
+ * FETCH every tag with IOSQE_ASYNC, racing STOP_DEV, then close the ring once
+ * STOP_DEV returned: that walks io_uring's list of cancelable commands.
+ */
+static void *race_async_fn(void *data)
+{
+	struct async_arg *a = data;
+	struct io_uring ring;
+	int ok = !io_uring_queue_init(DEPTH, &ring, 0);
+
+	if (ok)
+		for (int tag = 0; tag < DEPTH; tag++)
+			queue_io_cmd(&ring, UBLK_U_IO_FETCH_REQ, 0, tag, 0);
+	pthread_barrier_wait(&a->go);
+	if (ok)
+		io_uring_submit(&ring);
+	pthread_barrier_wait(&a->stopped);
+	if (ok)
+		io_uring_queue_exit(&ring);
+	return NULL;
+}
+
+/* reap @nr completions of canceled fetch commands, return how many */
+static int reap_aborts(struct io_uring *ring, int nr)
+{
+	struct __kernel_timespec ts = { .tv_sec = 2 };
+	struct io_uring_cqe *cqe;
+	int aborted = 0;
+
+	while (nr--) {
+		if (io_uring_wait_cqe_timeout(ring, &cqe, &ts))
+			break;
+		if (cqe->res == UBLK_IO_RES_ABORT)
+			aborted++;
+		io_uring_cqe_seen(ring, cqe);
+	}
+	return aborted;
+}
+
+/*
+ * A server thread fetches @nr_tags tags of queue @q, then completes each
+ * request after @delay_us, until its commands are aborted.
+ */
+struct server {
+	int q, first_tag, nr_tags, delay_us;
+	int failed;	/* result which ended the serve loop, 0 if none */
+	pthread_t thread;
+	pthread_barrier_t fetched;
+	/* race_fetch: wait for @go and @race_delay_us, open, fetch one by one */
+	pthread_barrier_t go;
+	int race, race_delay_us;
+};
+
+static void *server_fn(void *data)
+{
+	struct server *s = data;
+	struct io_uring ring;
+	struct io_uring_cqe *cqe;
+
+	if (io_uring_queue_init(DEPTH, &ring, 0)) {
+		pthread_barrier_wait(&s->fetched);
+		return NULL;
+	}
+	if (s->race) {
+		pthread_barrier_wait(&s->go);
+		usleep(s->race_delay_us);
+		if (open_cdev()) {
+			pthread_barrier_wait(&s->fetched);
+			io_uring_queue_exit(&ring);
+			return NULL;
+		}
+	}
+	for (int tag = s->first_tag; tag < s->first_tag + s->nr_tags; tag++) {
+		queue_io_cmd(&ring, UBLK_U_IO_FETCH_REQ, s->q, tag, 0);
+		if (s->race)
+			io_uring_submit(&ring);
+	}
+	io_uring_submit(&ring);
+	pthread_barrier_wait(&s->fetched);
+
+	while (!io_uring_wait_cqe(&ring, &cqe)) {
+		int tag = cqe->user_data & 0xffff;
+		const struct ublksrv_io_desc *iod = &iods[s->q][tag];
+		int res = cqe->res;
+
+		io_uring_cqe_seen(&ring, cqe);
+		if (res != UBLK_IO_RES_OK) {
+			__atomic_store_n(&s->failed, res, __ATOMIC_RELEASE);
+			break;
+		}
+		usleep(s->delay_us);
+		res = ublksrv_get_op(iod) <= UBLK_IO_OP_WRITE ?
+			iod->nr_sectors << 9 : 0;
+		queue_io_cmd(&ring, UBLK_U_IO_COMMIT_AND_FETCH_REQ, s->q, tag,
+			     res);
+		io_uring_submit(&ring);
+	}
+	io_uring_queue_exit(&ring);
+	return NULL;
+}
+
+/* returns once the fetch commands are issued; with @race, call race_go() */
+static void server_start(struct server *s)
+{
+	pthread_barrier_init(&s->fetched, NULL, 2);
+	pthread_barrier_init(&s->go, NULL, 2);
+	pthread_create(&s->thread, NULL, server_fn, s);
+	if (!s->race)
+		pthread_barrier_wait(&s->fetched);
+}
+
+/* wait up to 5s for the serve loop of @s to end, return its result */
+static int server_wait_failed(struct server *s)
+{
+	int res = 0;
+
+	for (int i = 0; i < 500 && !res; i++) {
+		res = __atomic_load_n(&s->failed, __ATOMIC_ACQUIRE);
+		if (!res)
+			usleep(10000);
+	}
+	return res;
+}
+
+/*
+ * A task fetches @nr_tags tags of queue @q and dies. Returns once its
+ * commands are canceled, see the top of this file.
+ */
+static int fetch_and_die(int q, int first_tag, int nr_tags)
+{
+	int status;
+	pid_t pid = fork();
+
+	if (pid < 0)
+		return -1;
+	if (!pid) {
+		struct io_uring ring;
+
+		if (io_uring_queue_init(DEPTH, &ring, 0))
+			_exit(1);
+		for (int tag = first_tag; tag < first_tag + nr_tags; tag++)
+			queue_io_cmd(&ring, UBLK_U_IO_FETCH_REQ, q, tag, 0);
+		if (io_uring_submit(&ring) != nr_tags)
+			_exit(1);
+		execlp("true", "true", NULL);
+		_exit(1);
+	}
+	if (waitpid(pid, &status, 0) != pid || !WIFEXITED(status) ||
+	    WEXITSTATUS(status))
+		return -1;
+	printf("q%d: task died with %d fetch cmds in flight\n", q, nr_tags);
+	return 0;
+}
+
+struct reads {
+	struct io_uring ring;
+	void *buf;
+	int fd, done, ok, eio, other;
+};
+
+/*
+ * Issue NR_READS reads at once, from @cpu if it is not negative: blk-mq
+ * maps the submitting CPU to the hw queue. A server holding its live
+ * tags for a while makes the other reads take the canceled tags.
+ */
+static int open_tries = 100;	/* 50ms each */
+
+static int reads_submit(struct reads *r, int cpu)
+{
+	char path[64];
+	cpu_set_t set, old;
+
+	snprintf(path, sizeof(path), "/dev/ublkb%d", dev_id);
+	r->fd = -1;
+	for (int i = 0; i < open_tries && r->fd < 0; i++) {
+		r->fd = open(path, O_RDONLY | O_DIRECT);
+		if (r->fd < 0)
+			usleep(50000);
+	}
+	if (r->fd < 0) {
+		fprintf(stderr, "open %s: %s\n", path, strerror(errno));
+		return -1;
+	}
+	if (posix_memalign(&r->buf, 4096, NR_READS * 4096) ||
+	    io_uring_queue_init(NR_READS, &r->ring, 0))
+		return -1;
+
+	for (int i = 0; i < NR_READS; i++)
+		io_uring_prep_read(get_sqe(&r->ring), r->fd,
+				   r->buf + i * 4096, 4096, i * 4096);
+
+	if (cpu >= 0) {
+		sched_getaffinity(0, sizeof(old), &old);
+		CPU_ZERO(&set);
+		CPU_SET(cpu, &set);
+		sched_setaffinity(0, sizeof(set), &set);
+	}
+	io_uring_submit(&r->ring);
+	if (cpu >= 0)
+		sched_setaffinity(0, sizeof(old), &old);
+	return 0;
+}
+
+static void reads_reap(struct reads *r, int timeout_s)
+{
+	struct __kernel_timespec ts = { .tv_sec = timeout_s };
+	struct io_uring_cqe *cqe;
+
+	while (r->done < NR_READS &&
+	       !io_uring_wait_cqe_timeout(&r->ring, &cqe, &ts)) {
+		if (cqe->res == 4096)
+			r->ok++;
+		else if (cqe->res == -EIO)
+			r->eio++;
+		else
+			r->other++;
+		r->done++;
+		io_uring_cqe_seen(&r->ring, cqe);
+	}
+}
+
+static void reads_put(struct reads *r)
+{
+	io_uring_queue_exit(&r->ring);
+	close(r->fd);
+	free(r->buf);
+}
+
+static int reads_result(struct reads *r)
+{
+	printf("reads: %d ok, %d -EIO, %d other, %d not completed\n",
+	       r->ok, r->eio, r->other, NR_READS - r->done);
+	reads_put(r);
+	return r->other || r->done < NR_READS ? KSFT_FAIL : KSFT_PASS;
+}
+
+static int start_dev(void)
+{
+	int ret = ublk_ctrl_start_dev(ctrl(), getpid());
+
+	printf("START_DEV: %d\n", ret);
+	return ret;
+}
+
+static int expect_err(const char *what, int ret, int want)
+{
+	if (ret == want)
+		return KSFT_PASS;
+	fprintf(stderr, "%s: %d, expected %d\n", what, ret, want);
+	return KSFT_FAIL;
+}
+
+/* START_DEV must fail with @want; if it went live, show what a read does */
+static int start_dev_expect(int want)
+{
+	struct reads r = {};
+	int ret = start_dev();
+
+	if (!ret && !reads_submit(&r, -1)) {
+		reads_reap(&r, 10);
+		reads_result(&r);
+	}
+	return expect_err("START_DEV", ret, want);
+}
+
+/* -ENODEV, or live over a canceling queue: then reads must complete */
+static int expect_enodev_or_reads(const char *what, int ret, struct reads *r)
+{
+	if (ret == -ENODEV)
+		return KSFT_PASS;
+	if (ret) {
+		fprintf(stderr, "%s: %d, expected 0 or %d\n", what, ret,
+			-ENODEV);
+		return KSFT_FAIL;
+	}
+	reads_reap(r, 10);
+	return reads_result(r);
+}
+
+static int test_stop_start(void)
+{
+	struct io_uring ring;
+	int ret;
+
+	if (add_dev(1, 0))
+		return KSFT_FAIL;
+	if (io_uring_queue_init(DEPTH, &ring, 0))
+		return KSFT_FAIL;
+	for (int tag = 0; tag < DEPTH; tag++)
+		queue_io_cmd(&ring, UBLK_U_IO_FETCH_REQ, 0, tag, 0);
+	io_uring_submit(&ring);
+
+	/* device is ready but not started: state is UBLK_S_DEV_DEAD */
+	ret = ublk_ctrl_stop_dev(ctrl());
+	printf("STOP_DEV: %d, canceled fetch cmds: %d/%d\n", ret,
+	       reap_aborts(&ring, DEPTH), DEPTH);
+
+	/* STOP_DEV canceled the attached server: no start until it exits */
+	ret = start_dev_expect(-EBUSY);
+	io_uring_queue_exit(&ring);
+	return ret;
+}
+
+static int test_partial_fetch(void)
+{
+	struct server b = { .q = 0, .first_tag = DEPTH - 1, .nr_tags = 1,
+			    .delay_us = SERVE_DELAY_US };
+	struct reads r = {};
+	int ret;
+
+	if (add_dev(1, 0) || fetch_and_die(0, 0, DEPTH - 1))
+		return KSFT_FAIL;
+
+	/* the last FETCH makes the queue ready and clears ->canceling */
+	server_start(&b);
+
+	ret = start_dev();
+	if (!ret && reads_submit(&r, -1))
+		ret = KSFT_FAIL;
+	else
+		ret = expect_enodev_or_reads("START_DEV", ret, &r);
+
+	cleanup();
+	pthread_join(b.thread, NULL);
+	return ret;
+}
+
+static int test_stop_restart(void)
+{
+	struct server s = { .q = 0, .nr_tags = DEPTH };
+	struct reads r = {};
+	int ret;
+
+	if (add_dev_noopen(1, 0))
+		return KSFT_FAIL;
+
+	/* no server is attached, so there is nothing to cancel */
+	ret = ublk_ctrl_stop_dev(ctrl());
+	printf("STOP_DEV: %d\n", ret);
+	if (open_cdev())
+		return KSFT_FAIL;
+
+	server_start(&s);
+	if (start_dev()) {
+		ret = KSFT_FAIL;
+	} else if (reads_submit(&r, -1)) {
+		ret = KSFT_FAIL;
+	} else {
+		reads_reap(&r, 10);
+		ret = reads_result(&r);
+		if (r.ok != NR_READS)
+			ret = KSFT_FAIL;
+	}
+
+	cleanup();
+	pthread_join(s.thread, NULL);
+	return ret;
+}
+
+struct start_arg {
+	pthread_barrier_t go;
+	int ret, reads_ok;
+};
+
+static void *race_start_fn(void *data)
+{
+	struct start_arg *a = data;
+	struct reads r = {};
+
+	pthread_barrier_wait(&a->go);
+	a->ret = ublk_ctrl_start_dev(ctrl(), getpid());
+	a->reads_ok = 1;
+	/*
+	 * The disk may be gone already if STOP_DEV came right after, and
+	 * reads may fail then: they only have to complete.
+	 */
+	if (!a->ret && !reads_submit(&r, -1)) {
+		reads_reap(&r, 5);
+		a->reads_ok = r.done == NR_READS;
+		if (a->reads_ok)
+			reads_put(&r);
+		else
+			reads_result(&r);
+	}
+	ctrl_put();
+	return NULL;
+}
+
+static int test_race_start(void)
+{
+	int live = 0, ebusy = 0;
+
+	open_tries = 4;
+	for (int i = 0; i < RACE_LOOPS; i++) {
+		struct server s = { .q = 0, .nr_tags = DEPTH };
+		struct start_arg a = {};
+		pthread_t t;
+
+		if (add_dev(1, 0))
+			return KSFT_FAIL;
+		server_start(&s);
+		pthread_barrier_init(&a.go, NULL, 2);
+		pthread_create(&t, NULL, race_start_fn, &a);
+		pthread_barrier_wait(&a.go);
+		ublk_ctrl_stop_dev(ctrl());
+		pthread_join(t, NULL);
+		cleanup();
+		pthread_join(s.thread, NULL);
+
+		if (a.ret == 0)
+			live++;
+		else if (a.ret == -EBUSY)
+			ebusy++;
+		if ((a.ret && a.ret != -EBUSY) || !a.reads_ok) {
+			fprintf(stderr, "loop %d: START_DEV %d, reads %s\n",
+				i, a.ret, a.reads_ok ? "ok" : "failed");
+			return KSFT_FAIL;
+		}
+	}
+	printf("%d loops: START_DEV won %d, got -EBUSY %d\n", RACE_LOOPS,
+	       live, ebusy);
+	return KSFT_PASS;
+}
+
+static int test_race_fetch(void)
+{
+	int live = 0, ebusy = 0;
+
+	for (int i = 0; i < RACE_LOOPS; i++) {
+		/* vary who goes first: STOP_DEV, or the server's open */
+		struct server s = { .q = 0, .nr_tags = DEPTH, .race = 1,
+				    .race_delay_us = (i % 10) * 50 };
+		struct reads r = {};
+		int ret;
+
+		if (add_dev_noopen(1, 0))
+			return KSFT_FAIL;
+		server_start(&s);
+		pthread_barrier_wait(&s.go);
+		ublk_ctrl_stop_dev(ctrl());
+		pthread_barrier_wait(&s.fetched);
+
+		ret = ublk_ctrl_start_dev(ctrl(), getpid());
+		if (!ret) {
+			/* nothing was taken, so every read has to succeed */
+			live++;
+			if (reads_submit(&r, -1))
+				return KSFT_FAIL;
+			reads_reap(&r, 5);
+			if (r.ok != NR_READS) {
+				fprintf(stderr, "loop %d: live, but ", i);
+				reads_result(&r);
+				return KSFT_FAIL;
+			}
+			reads_put(&r);
+		} else if (ret == -EBUSY) {
+			ebusy++;
+		} else {
+			fprintf(stderr, "loop %d: START_DEV %d\n", i, ret);
+			return KSFT_FAIL;
+		}
+		cleanup();
+		pthread_join(s.thread, NULL);
+	}
+	printf("%d loops: START_DEV worked %d, got -EBUSY %d\n", RACE_LOOPS,
+	       live, ebusy);
+	return KSFT_PASS;
+}
+
+static int test_race_async_fetch(void)
+{
+	io_cmd_sqe_flags = IOSQE_ASYNC;
+	for (int i = 0; i < ASYNC_LOOPS; i++) {
+		struct async_arg a;
+		pthread_t t;
+
+		if (add_dev(1, 0))
+			return KSFT_FAIL;
+		pthread_barrier_init(&a.go, NULL, 2);
+		pthread_barrier_init(&a.stopped, NULL, 2);
+		pthread_create(&t, NULL, race_async_fn, &a);
+		pthread_barrier_wait(&a.go);
+		/* vary where STOP_DEV lands relative to the FETCHes */
+		usleep((i % 5) * 1000);
+		ublk_ctrl_stop_dev(ctrl());
+		pthread_barrier_wait(&a.stopped);
+		pthread_join(t, NULL);
+		cleanup();
+	}
+	printf("%d loops done\n", ASYNC_LOOPS);
+	return KSFT_PASS;
+}
+
+/*
+ * The stopped server is gone: reopen /dev/ublkcN, start a new server and
+ * the device, and read from it.
+ */
+static int restart_server(void)
+{
+	struct server s = { .q = 0, .nr_tags = DEPTH };
+	struct reads r = {};
+	int ret;
+
+	close_cdev();
+	if (open_cdev())
+		return KSFT_FAIL;
+	server_start(&s);
+	usleep(200000);
+	if (s.failed) {
+		fprintf(stderr, "new server: FETCH failed %d\n", s.failed);
+		cleanup();
+		pthread_join(s.thread, NULL);
+		return KSFT_FAIL;
+	}
+	/* -EEXIST until the old disk is freed, e.g. after a udev probe */
+	for (int i = 0; i < 100; i++) {
+		ret = start_dev();
+		if (ret != -EEXIST)
+			break;
+		usleep(50000);
+	}
+	if (ret || reads_submit(&r, -1)) {
+		fprintf(stderr, "new server: START_DEV %d\n", ret);
+		ret = KSFT_FAIL;
+	} else {
+		reads_reap(&r, 10);
+		ret = reads_result(&r);
+		if (r.ok != NR_READS)
+			ret = KSFT_FAIL;
+	}
+	cleanup();
+	pthread_join(s.thread, NULL);
+	return ret;
+}
+
+/*
+ * stop_attached: STOP_DEV while a server has the device open but has
+ * fetched nothing stops that server too: its FETCH gets ABORT and
+ * START_DEV -EBUSY. Once it is gone, a new server can start the device.
+ */
+static int test_stop_attached(void)
+{
+	struct server s = { .q = 0, .nr_tags = DEPTH };
+	int ret;
+
+	if (add_dev(1, 0))
+		return KSFT_FAIL;
+	ret = ublk_ctrl_stop_dev(ctrl());
+	printf("STOP_DEV: %d\n", ret);
+
+	server_start(&s);
+	ret = server_wait_failed(&s);
+	printf("FETCH after STOP_DEV: %d\n", ret);
+	if (ret != UBLK_IO_RES_ABORT) {
+		cleanup();
+		pthread_join(s.thread, NULL);
+		return KSFT_FAIL;
+	}
+	pthread_join(s.thread, NULL);
+	if (expect_err("START_DEV", start_dev(), -EBUSY))
+		return KSFT_FAIL;
+	return restart_server();
+}
+
+/*
+ * stop_live_restart: STOP_DEV on a live device; once its server is gone,
+ * a new server has to fetch and start it again.
+ */
+static int test_stop_live_restart(void)
+{
+	struct server s = { .q = 0, .nr_tags = DEPTH };
+	int ret;
+
+	if (add_dev(1, 0))
+		return KSFT_FAIL;
+	server_start(&s);
+	if (expect_err("START_DEV", start_dev(), 0))
+		return KSFT_FAIL;
+	ret = ublk_ctrl_stop_dev(ctrl());
+	printf("STOP_DEV: %d\n", ret);
+	pthread_join(s.thread, NULL);
+	return restart_server();
+}
+
+/* first CPU which blk-mq maps to hw queue @q */
+static int queue_cpu(int q)
+{
+	char path[96];
+	FILE *f;
+	int cpu = -1;
+
+	snprintf(path, sizeof(path), "/sys/block/ublkb%d/mq/%d/cpu_list",
+		 dev_id, q);
+	for (int i = 0; i < 100 && !(f = fopen(path, "r")); i++)
+		usleep(50000);
+	if (!f)
+		return -1;
+	if (fscanf(f, "%d", &cpu) != 1)
+		cpu = -1;
+	fclose(f);
+	return cpu;
+}
+
+static int test_recovery(void)
+{
+	struct server q1 = { .q = 1, .nr_tags = DEPTH };
+	struct io_uring ring;
+	struct reads r = {};
+	int ret, cpu;
+
+	if (add_dev(NR_QUEUES, UBLK_F_USER_RECOVERY))
+		return KSFT_FAIL;
+
+	/* the first server: fetch everything, start, then die */
+	if (io_uring_queue_init(NR_QUEUES * DEPTH, &ring, 0))
+		return KSFT_FAIL;
+	for (int q = 0; q < NR_QUEUES; q++)
+		for (int tag = 0; tag < DEPTH; tag++)
+			queue_io_cmd(&ring, UBLK_U_IO_FETCH_REQ, q, tag, 0);
+	io_uring_submit(&ring);
+	if (start_dev())
+		return KSFT_FAIL;
+	cpu = queue_cpu(0);
+	if (cpu < 0) {
+		printf("no CPU maps to queue 0, can't aim the reads\n");
+		return KSFT_SKIP;
+	}
+
+	io_uring_queue_exit(&ring);
+	close_cdev();
+	for (int i = 0; i < 100 && dev_state() != UBLK_S_DEV_QUIESCED; i++)
+		usleep(50000);
+	printf("server exited, state %d (QUIESCED is %d)\n", dev_state(),
+	       UBLK_S_DEV_QUIESCED);
+
+	for (int i = 0; i < 100; i++) {
+		ret = ublk_ctrl_start_user_recovery(ctrl());
+		if (ret != -EBUSY)
+			break;
+		usleep(50000);
+	}
+	printf("START_USER_RECOVERY: %d\n", ret);
+	if (ret || open_cdev())
+		return KSFT_FAIL;
+
+	/* queue 0 gets ready, then its task dies before queue 1 is ready */
+	if (fetch_and_die(0, 0, DEPTH))
+		return KSFT_FAIL;
+
+	printf("reads on cpu %d, which maps to queue 0\n", cpu);
+	if (reads_submit(&r, cpu))
+		return KSFT_FAIL;
+
+	server_start(&q1);
+	ret = ublk_ctrl_end_user_recovery(ctrl(), getpid());
+	printf("END_USER_RECOVERY: %d\n", ret);
+	if (ret && ret != -ENODEV) {
+		fprintf(stderr, "END_USER_RECOVERY: %d, expected 0 or %d\n",
+			ret, -ENODEV);
+		ret = KSFT_FAIL;
+	} else {
+		ret = KSFT_PASS;
+	}
+
+	/*
+	 * After -ENODEV, requests held back on queue 0 wait for STOP_DEV to
+	 * fail them. Close the disk before DEL_DEV, which waits for its last
+	 * reference.
+	 */
+	reads_reap(&r, 1);
+	ublk_ctrl_stop_dev(ctrl());
+	reads_reap(&r, 5);
+	if (reads_result(&r))
+		ret = KSFT_FAIL;
+	cleanup();
+	pthread_join(q1.thread, NULL);
+	return ret;
+}
+
+static const struct {
+	const char *name;
+	int (*fn)(void);
+} modes[] = {
+	{ "stop_start",		test_stop_start },
+	{ "partial_fetch",	test_partial_fetch },
+	{ "recovery",		test_recovery },
+	{ "stop_restart",	test_stop_restart },
+	{ "race_start",		test_race_start },
+	{ "race_fetch",		test_race_fetch },
+	{ "race_async_fetch",	test_race_async_fetch },
+	{ "stop_attached",	test_stop_attached },
+	{ "stop_live_restart",	test_stop_live_restart },
+};
+
+int main(int argc, char **argv)
+{
+	int (*fn)(void) = NULL;
+	int ret;
+
+	for (int i = 0; argc == 2 && i < ARRAY_SIZE(modes); i++)
+		if (!strcmp(argv[1], modes[i].name))
+			fn = modes[i].fn;
+	if (!fn) {
+		fprintf(stderr, "usage: %s MODE, modes:", argv[0]);
+		for (int i = 0; i < ARRAY_SIZE(modes); i++)
+			fprintf(stderr, " %s", modes[i].name);
+		fprintf(stderr, "\n");
+		return KSFT_FAIL;
+	}
+
+	/* keep the output of a run that ends in an oops */
+	setvbuf(stdout, NULL, _IOLBF, 0);
+
+	if (access(CTRL_DEV, F_OK)) {
+		perror(CTRL_DEV);
+		return KSFT_SKIP;
+	}
+	printf("%s\n", argv[1]);
+	ret = fn();
+	cleanup();
+	ctrl_put();
+	return ret;
+}
-- 
2.55.0


  parent reply	other threads:[~2026-10-01 12:55 UTC|newest]

Thread overview: 18+ messages / expand[flat|nested]  mbox.gz  Atom feed  top
2026-10-01 12:54 [PATCH 0/8] ublk: don't dispatch to canceled io commands Ming Lei
2026-10-01 12:54 ` [PATCH 1/8] ublk: keep a canceled FETCH round canceling until the server is gone Ming Lei
2026-10-01 12:54 ` [PATCH 2/8] ublk: mark the io command cancelable before publishing it Ming Lei
2026-10-01 12:54 ` [PATCH 3/8] ublk: mark the batch fetch command cancelable before linking it Ming Lei
2026-10-01 12:54 ` [PATCH 4/8] ublk: reset the FETCH round in release also without a disk Ming Lei
2026-10-01 12:54 ` [PATCH 5/8] ublk: reset the FETCH round under ub->mutex Ming Lei
2026-10-01 12:54 ` [PATCH 6/8] ublk: let STOP_DEV cancel the server's commands before its release Ming Lei
2026-10-01 12:54 ` [PATCH 7/8] selftests: ublk: move the control command helpers into ctrl.c Ming Lei
2026-10-01 12:54 ` Ming Lei [this message]
2026-10-05 16:23 ` [PATCH] ublk: refuse to go live after an io command was canceled Josef Bacik
2026-10-06 14:14   ` Ming Lei
2026-10-05 16:23     ` [PATCH v2] " Josef Bacik
2026-10-05 18:50 ` [PATCH 0/8] ublk: don't dispatch to canceled io commands Josef Bacik
2026-10-06 16:10 ` [PATCH 0/4] ublk: fix UBLK_CMD_QUIESCE_DEV leaving commands behind Josef Bacik
2026-10-06 13:05   ` [PATCH 1/4] ublk: don't cancel commands in QUIESCE_DEV on a device that isn't live Josef Bacik
2026-10-06 14:49   ` [PATCH 2/4] ublk: drop QUIESCE_DEV's wait for an idle command Josef Bacik
2026-10-06 14:50   ` [PATCH 3/4] ublk: give the command back from COMMIT_AND_FETCH on a canceling queue Josef Bacik
2026-10-06 14:50   ` [PATCH 4/4] ublk: keep canceling in QUIESCE_DEV until the server's commands are taken Josef Bacik

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=20261001125422.1364260-9-tom.leiming@gmail.com \
    --to=tom.leiming@gmail.com \
    --cc=axboe@kernel.dk \
    --cc=csander@purestorage.com \
    --cc=josef@toxicpanda.com \
    --cc=linux-block@vger.kernel.org \
    /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