From: Jens Axboe <axboe@kernel.dk>
To: io-uring@vger.kernel.org
Cc: linux-arm-kernel@lists.infradead.org,
linux-kernel@vger.kernel.org, tglx@kernel.org, mingo@redhat.com,
peterz@infradead.org, Jens Axboe <axboe@kernel.dk>
Subject: [PATCH 10/15] io-wq: support handing a task identity to an idle worker
Date: Fri, 11 Sep 2026 09:41:00 -0600 [thread overview]
Message-ID: <20260911154148.644489-11-axboe@kernel.dk> (raw)
In-Reply-To: <20260911154148.644489-1-axboe@kernel.dk>
Add the io-wq side of handing a blocked submitter's identity to an idle
worker. io_wq_handoff_claim() picks an idle worker that can take over
the identity of the current task, and arranges for it to run a caller
supplied function when it wakes instead of continuing as a worker.
Only a worker inside the idle sleep of the worker loop is claimable,
being on the free list isn't enough as it may be sleeping elsewhere.
This is tracked in acct->nr_iosleep. Redundant workers exit through
the normal idle timeout.
Signed-off-by: Jens Axboe <axboe@kernel.dk>
---
io_uring/io-wq.c | 263 ++++++++++++++++++++++++++++++++++++++++++++---
io_uring/io-wq.h | 14 +++
kernel/fork.c | 3 +-
3 files changed, 267 insertions(+), 13 deletions(-)
diff --git a/io_uring/io-wq.c b/io_uring/io-wq.c
index 2ca223e47d41..3d4eb4992d5b 100644
--- a/io_uring/io-wq.c
+++ b/io_uring/io-wq.c
@@ -20,6 +20,7 @@
#include <linux/sched/sysctl.h>
#include <uapi/linux/io_uring.h>
#include <linux/kcov.h>
+#include <linux/thread_handoff.h>
#include "io-wq.h"
#include "slist.h"
@@ -32,6 +33,15 @@ enum {
IO_WORKER_F_UP = 0, /* up and active */
IO_WORKER_F_RUNNING = 1, /* account as running */
IO_WORKER_F_FREE = 2, /* worker on free list */
+ IO_WORKER_F_IDLE_SLEEP = 3, /* in the idle sleep of the worker loop */
+};
+
+/* worker->handoff handshake */
+enum {
+ IO_WORKER_HANDOFF_NONE = 0,
+ IO_WORKER_HANDOFF_PROMOTE, /* claimed, set under ->workers_lock */
+ IO_WORKER_HANDOFF_DONE, /* claimer took over the worker */
+ IO_WORKER_HANDOFF_FINISHED, /* identity moved, demoted released */
};
enum {
@@ -64,6 +74,10 @@ struct io_worker {
struct callback_head create_work;
int init_retries;
+ atomic_t handoff;
+ io_wq_handoff_fn *handoff_fn;
+ struct task_struct *handoff_task;
+
union {
struct rcu_head rcu;
struct delayed_work work;
@@ -88,6 +102,9 @@ struct io_wq_acct {
unsigned max_workers;
atomic_t nr_running;
+ /* workers in the idle sleep, claimable for a handoff */
+ atomic_t nr_iosleep;
+
/**
* The list of free workers. Protected by #workers_lock
* (write) and RCU (read).
@@ -151,6 +168,8 @@ static bool io_acct_cancel_pending_work(struct io_wq *wq,
struct io_wq_acct *acct,
struct io_cb_cancel_data *match);
static void create_worker_cb(struct callback_head *cb);
+static void create_worker_cont(struct callback_head *cb);
+static bool io_task_work_match(struct callback_head *cb, void *data);
static void io_wq_cancel_tw_create(struct io_wq *wq);
static inline unsigned int __io_get_work_hash(unsigned int work_flags)
@@ -233,7 +252,7 @@ static bool io_task_worker_match(struct callback_head *cb, void *data)
return worker == data;
}
-static void io_worker_exit(struct io_worker *worker)
+static void __noreturn io_worker_exit(struct io_worker *worker)
{
struct io_wq *wq = worker->wq;
struct io_wq_acct *acct = io_wq_get_acct(worker);
@@ -411,7 +430,7 @@ static bool io_queue_worker_create(struct io_worker *worker,
atomic_inc(&wq->worker_refs);
init_task_work(&worker->create_work, func);
- if (!task_work_add(wq->task, &worker->create_work, TWA_SIGNAL)) {
+ if (!io_wq_task_work_add(wq->task, &worker->create_work, TWA_SIGNAL)) {
/*
* EXIT may have been set after checking it above, check after
* adding the task_work and remove any creation item if it is
@@ -687,19 +706,48 @@ static void io_worker_handle_work(struct io_wq_acct *acct,
} while (1);
}
-static int io_wq_worker(void *data)
+/*
+ * Leave the idle sleep, check if we got claimed while in it. Serialized with
+ * the claimer by ->workers_lock, so a claim can't be missed or raced.
+ */
+static bool io_wq_worker_idle_done(struct io_wq_acct *acct,
+ struct io_worker *worker)
{
- struct io_worker *worker = data;
- struct io_wq_acct *acct = io_wq_get_acct(worker);
- struct io_wq *wq = worker->wq;
- bool exit_mask = false, last_timeout = false;
- char buf[TASK_COMM_LEN] = {};
+ bool promoted;
- set_mask_bits(&worker->flags, 0,
- BIT(IO_WORKER_F_UP) | BIT(IO_WORKER_F_RUNNING));
+ atomic_dec(&acct->nr_iosleep);
+ raw_spin_lock(&acct->workers_lock);
+ clear_bit(IO_WORKER_F_IDLE_SLEEP, &worker->flags);
+ /* the claimer may have committed already, DONE rather than PROMOTE */
+ promoted = atomic_read(&worker->handoff) != IO_WORKER_HANDOFF_NONE;
+ raw_spin_unlock(&acct->workers_lock);
- snprintf(buf, sizeof(buf), "iou-wrk-%d", wq->task->pid);
- set_task_comm(current, buf);
+ WARN_ON_ONCE(promoted && worker->handoff_task != current);
+ return promoted;
+}
+
+/* claimed, wait for the claimer to finish taking over the worker struct */
+static io_wq_handoff_fn *io_wq_worker_promoted(struct io_worker *worker)
+{
+ io_wq_handoff_fn *fn = worker->handoff_fn;
+
+ /* stop the scheduler from treating us as a worker while we wait */
+ current->flags &= ~(PF_IO_WORKER | PF_USER_WORKER);
+
+ wait_var_event(&worker->handoff,
+ atomic_read_acquire(&worker->handoff) == IO_WORKER_HANDOFF_DONE);
+
+ WARN_ON_ONCE(current->worker_private);
+ atomic_set(&worker->handoff, IO_WORKER_HANDOFF_NONE);
+ return fn;
+}
+
+/* the worker loop, only returns if claimed for a handoff */
+static io_wq_handoff_fn *io_wq_worker_run(struct io_worker *worker)
+{
+ struct io_wq_acct *acct = io_wq_get_acct(worker);
+ bool exit_mask = false, last_timeout = false;
+ struct io_wq *wq = worker->wq;
while (!test_bit(IO_WQ_BIT_EXIT, &wq->state)) {
long ret;
@@ -724,6 +772,11 @@ static int io_wq_worker(void *data)
if ((last_timeout && (exit_mask || acct->nr_workers > 1)) ||
test_bit(IO_WQ_BIT_EXIT_ON_IDLE, &wq->state)) {
acct->nr_workers--;
+ /* a handoff can't claim an exiting worker */
+ if (test_bit(IO_WORKER_F_FREE, &worker->flags)) {
+ clear_bit(IO_WORKER_F_FREE, &worker->flags);
+ hlist_nulls_del_rcu(&worker->nulls_node);
+ }
raw_spin_unlock(&acct->workers_lock);
__set_current_state(TASK_RUNNING);
break;
@@ -733,7 +786,12 @@ static int io_wq_worker(void *data)
raw_spin_unlock(&acct->workers_lock);
if (io_run_task_work())
continue;
+ /* claimable only in this sleep, the free list isn't enough */
+ set_bit(IO_WORKER_F_IDLE_SLEEP, &worker->flags);
+ atomic_inc(&acct->nr_iosleep);
ret = schedule_timeout(WORKER_IDLE_TIMEOUT);
+ if (unlikely(io_wq_worker_idle_done(acct, worker)))
+ return io_wq_worker_promoted(worker);
if (signal_pending(current)) {
struct ksignal ksig;
@@ -752,9 +810,190 @@ static int io_wq_worker(void *data)
io_worker_handle_work(acct, worker);
io_worker_exit(worker);
+}
+
+static int io_wq_worker(void *data)
+{
+ struct io_worker *worker = data;
+ struct io_wq *wq = worker->wq;
+ io_wq_handoff_fn *fn;
+ char buf[TASK_COMM_LEN] = {};
+
+ set_mask_bits(&worker->flags, 0,
+ BIT(IO_WORKER_F_UP) | BIT(IO_WORKER_F_RUNNING));
+
+ snprintf(buf, sizeof(buf), "iou-wrk-%d", wq->task->pid);
+ set_task_comm(current, buf);
+
+ /*
+ * Only returns if we got handed an identity. -EIOCBQUEUED means we got
+ * demoted again while running it, back to the worker loop.
+ */
+ for (;;) {
+ long ret;
+
+ fn = io_wq_worker_run(worker);
+ ret = fn();
+ /* what we return is what userspace gets on some archs */
+ if (ret != -EIOCBQUEUED)
+ return ret;
+ worker = current->worker_private;
+ }
+}
+
+/* find and claim an idle sleeping worker, see io_wq_worker_idle_done() */
+static struct io_worker *io_wq_acct_handoff_claim(struct io_wq *wq,
+ struct io_wq_acct *acct,
+ io_wq_handoff_fn *fn)
+{
+ struct io_worker *worker, *found = NULL;
+ struct hlist_nulls_node *n;
+
+ raw_spin_lock(&acct->workers_lock);
+ if (test_bit(IO_WQ_BIT_EXIT, &wq->state))
+ goto out_unlock;
+ hlist_nulls_for_each_entry(worker, n, &acct->free_list, nulls_node) {
+ /* only claimable inside the idle sleep of the worker loop */
+ if (!test_bit(IO_WORKER_F_IDLE_SLEEP, &worker->flags))
+ continue;
+ if (!thread_handoff_compatible(current, worker->task))
+ continue;
+ clear_bit(IO_WORKER_F_FREE, &worker->flags);
+ hlist_nulls_del_init_rcu(&worker->nulls_node);
+ worker->handoff_fn = fn;
+ worker->handoff_task = worker->task;
+ atomic_set_release(&worker->handoff, IO_WORKER_HANDOFF_PROMOTE);
+ found = worker;
+ break;
+ }
+out_unlock:
+ raw_spin_unlock(&acct->workers_lock);
+ if (found)
+ wake_up_process(found->handoff_task);
+ return found;
+}
+
+/*
+ * task_work_add() that doesn't notify a task inside a blocking inline issue,
+ * it'd interrupt the sleep. io_handoff_end() picks pending work up instead.
+ */
+int io_wq_task_work_add(struct task_struct *task, struct callback_head *cb,
+ enum task_work_notify_mode notify)
+{
+ int ret;
+
+ if (notify != TWA_SIGNAL && notify != TWA_SIGNAL_NO_IPI)
+ return task_work_add(task, cb, notify);
+
+ ret = task_work_add(task, cb, TWA_NONE);
+ if (ret)
+ return ret;
+ if (READ_ONCE(task->flags) & PF_IO_HANDOFF)
+ return 0;
+ if (notify == TWA_SIGNAL)
+ set_notify_signal(task);
+ else
+ __set_notify_signal(task);
return 0;
}
+/* claim an idle worker to hand our identity to, pairs with _commit() */
+struct task_struct *io_wq_handoff_claim(struct io_wq *wq, bool bound,
+ io_wq_handoff_fn *fn)
+{
+ struct io_worker *worker;
+
+ worker = io_wq_acct_handoff_claim(wq, io_get_acct(wq, bound), fn);
+ if (!worker)
+ worker = io_wq_acct_handoff_claim(wq, io_get_acct(wq, !bound), fn);
+ if (worker)
+ return worker->task;
+ return NULL;
+}
+
+/* we take over the worker @dst was, @dst goes on to run the handoff fn */
+void io_wq_handoff_commit(struct task_struct *dst)
+{
+ struct io_worker *worker = dst->worker_private;
+ struct task_struct *src = current;
+ struct io_wq *wq = worker->wq;
+ struct callback_head *cb;
+
+ WARN_ON_ONCE(src->worker_private);
+ WARN_ON_ONCE(worker->task != dst);
+ WARN_ON_ONCE(wq->task != src);
+
+ /* the sched hooks around @dst's wakeup cope with NULL worker_private */
+ WRITE_ONCE(dst->worker_private, NULL);
+ src->worker_private = worker;
+ WRITE_ONCE(worker->task, src);
+ src->flags |= PF_IO_WORKER | PF_USER_WORKER;
+
+ /* the user task owns the wq */
+ get_task_struct(dst);
+ WRITE_ONCE(wq->task, dst);
+
+ /* move pending worker creations along, we may block for a while */
+ while ((cb = task_work_cancel_match(src, io_task_work_match, wq))) {
+ struct io_worker *w = container_of(cb, struct io_worker,
+ create_work);
+
+ if (!task_work_add(dst, cb, TWA_SIGNAL))
+ continue;
+ io_worker_cancel_cb(w);
+ if (cb->func == create_worker_cont)
+ kfree(w);
+ }
+ put_task_struct(src);
+
+ atomic_set_release(&worker->handoff, IO_WORKER_HANDOFF_DONE);
+ wake_up_var(&worker->handoff);
+}
+
+/* worker loop entry for a demoted task, only returns on another handoff */
+io_wq_handoff_fn *io_wq_handoff_worker(void)
+{
+ struct io_worker *worker = current->worker_private;
+ char buf[TASK_COMM_LEN] = {};
+
+ WARN_ON_ONCE(!io_wq_current_is_worker());
+
+ /* the promoted task reads our state until it's done migrating it */
+ wait_var_event(&worker->handoff,
+ atomic_read_acquire(&worker->handoff) == IO_WORKER_HANDOFF_FINISHED);
+ atomic_set(&worker->handoff, IO_WORKER_HANDOFF_NONE);
+
+ snprintf(buf, sizeof(buf), "iou-wrk-%d", worker->wq->task->pid);
+ set_task_comm(current, buf);
+ set_cpus_allowed_ptr(current, worker->wq->cpu_mask);
+
+ return io_wq_worker_run(worker);
+}
+
+/* release the demoted task @tsk to run as the worker it now is */
+void io_wq_handoff_finished(struct task_struct *tsk)
+{
+ struct io_worker *worker = tsk->worker_private;
+
+ atomic_set_release(&worker->handoff, IO_WORKER_HANDOFF_FINISHED);
+ wake_up_var(&worker->handoff);
+}
+
+/* idle sleepers to keep around as handoff targets */
+#define IO_WQ_HANDOFF_SPARES 2
+
+/* true if a worker is claimable, @topup forks one if below the target */
+bool io_wq_handoff_spare(struct io_wq *wq, bool bound, bool topup)
+{
+ struct io_wq_acct *acct = io_get_acct(wq, bound);
+ unsigned int idle;
+
+ idle = atomic_read(&acct->nr_iosleep);
+ if (topup && idle < IO_WQ_HANDOFF_SPARES)
+ io_wq_create_worker(wq, acct);
+ return idle > 0;
+}
+
/*
* Called when a worker is scheduled in. Mark us as currently running.
*/
diff --git a/io_uring/io-wq.h b/io_uring/io-wq.h
index 42f00a47a9c9..98357b665e54 100644
--- a/io_uring/io-wq.h
+++ b/io_uring/io-wq.h
@@ -4,6 +4,7 @@
#include <linux/refcount.h>
#include <linux/io_uring_types.h>
+#include <linux/task_work.h>
struct io_wq;
@@ -47,6 +48,19 @@ void io_wq_set_exit_on_idle(struct io_wq *wq, bool enable);
void io_wq_enqueue(struct io_wq *wq, struct io_wq_work *work);
void io_wq_hash_work(struct io_wq_work *work, void *val);
+typedef long (io_wq_handoff_fn)(void);
+
+/* claim an idle worker, it runs @fn instead of the worker loop when woken */
+struct task_struct *io_wq_handoff_claim(struct io_wq *wq, bool bound,
+ io_wq_handoff_fn *fn);
+
+void io_wq_handoff_commit(struct task_struct *dst);
+io_wq_handoff_fn *io_wq_handoff_worker(void);
+bool io_wq_handoff_spare(struct io_wq *wq, bool bound, bool topup);
+void io_wq_handoff_finished(struct task_struct *tsk);
+int io_wq_task_work_add(struct task_struct *task, struct callback_head *cb,
+ enum task_work_notify_mode notify);
+
int io_wq_cpu_affinity(struct io_uring_task *tctx, cpumask_var_t mask);
int io_wq_max_workers(struct io_wq *wq, int *new_count);
bool io_wq_worker_stopped(void);
diff --git a/kernel/fork.c b/kernel/fork.c
index 510c8a9aa870..f31af9c5bae7 100644
--- a/kernel/fork.c
+++ b/kernel/fork.c
@@ -2688,11 +2688,12 @@ struct task_struct * __init fork_idle(int cpu)
* creating io_uring workers. It returns a created task, or an error pointer.
* The returned task is inactive, and the caller must fire it up through
* wake_up_new_task(p). All signals are blocked in the created task.
+ * CLONE_SYSVSEM as a worker may take over a user thread's identity.
*/
struct task_struct *create_io_thread(int (*fn)(void *), void *arg, int node)
{
unsigned long flags = CLONE_FS|CLONE_FILES|CLONE_SIGHAND|CLONE_THREAD|
- CLONE_IO|CLONE_VM|CLONE_UNTRACED;
+ CLONE_IO|CLONE_VM|CLONE_UNTRACED|CLONE_SYSVSEM;
struct kernel_clone_args args = {
.flags = flags,
.fn = fn,
--
2.55.0
next prev parent reply other threads:[~2026-09-11 15:42 UTC|newest]
Thread overview: 18+ messages / expand[flat|nested] mbox.gz Atom feed top
2026-09-11 15:40 [RFC PATCH 00/15] io_uring: thread identity handoff for blocking inline issue Jens Axboe
2026-09-11 15:40 ` [PATCH 01/15] kernel: add thread identity handoff Jens Axboe
2026-09-11 15:40 ` [PATCH 02/15] sched: call into io_uring when a PF_IO_HANDOFF task blocks Jens Axboe
2026-09-11 15:40 ` [PATCH 03/15] arm64: implement thread identity handoff Jens Axboe
2026-09-11 15:40 ` [PATCH 04/15] x86: " Jens Axboe
2026-09-11 15:40 ` [PATCH 05/15] io_uring/kbuf: use io_ring_submit_unlock() helper Jens Axboe
2026-09-11 15:40 ` [PATCH 06/15] io_uring: keep the tctx nodes on a list Jens Axboe
2026-09-11 15:40 ` [PATCH 07/15] io_uring: add uring_lock section depth tracking and blockable opdef flag Jens Axboe
2026-09-11 15:40 ` [PATCH 08/15] io_uring: split io_uring_enter() and io_submit_sqes() into helpers Jens Axboe
2026-09-11 15:40 ` [PATCH 09/15] io_uring: keep the submission plug on the io_submit_sqes() stack Jens Axboe
2026-09-11 15:41 ` Jens Axboe [this message]
2026-09-11 15:41 ` [PATCH 11/15] io_uring: enable handing submitter identity to an io-wq worker Jens Axboe
2026-09-11 15:41 ` [PATCH 12/15] io_uring: defer the identity migration to the end of the submission Jens Axboe
2026-09-11 15:41 ` [PATCH 13/15] io_uring: issue blockable requests inline in blocking mode Jens Axboe
2026-09-11 15:41 ` [PATCH 14/15] io_uring: add tracepoints for the handoff operation Jens Axboe
2026-09-11 15:41 ` [PATCH 15/15] io_uring: issue IOSQE_ASYNC requests inline when a handoff is possible Jens Axboe
2026-09-11 17:33 ` [RFC PATCH 00/15] io_uring: thread identity handoff for blocking inline issue Gabriel Krisman Bertazi
2026-09-11 17:51 ` Jens Axboe
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=20260911154148.644489-11-axboe@kernel.dk \
--to=axboe@kernel.dk \
--cc=io-uring@vger.kernel.org \
--cc=linux-arm-kernel@lists.infradead.org \
--cc=linux-kernel@vger.kernel.org \
--cc=mingo@redhat.com \
--cc=peterz@infradead.org \
--cc=tglx@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 an external index of several public inboxes,
see mirroring instructions on how to clone and mirror
all data and code used by this external index.