diff options
27 files changed, 1474 insertions, 236 deletions
diff --git a/Documentation/block/ublk.rst b/Documentation/block/ublk.rst index 0413dcd9ef69..28300fee22bf 100644 --- a/Documentation/block/ublk.rst +++ b/Documentation/block/ublk.rst @@ -382,17 +382,17 @@ Zero copy --------- ublk zero copy relies on io_uring's fixed kernel buffer, which provides -two APIs: `io_buffer_register_bvec()` and `io_buffer_unregister_bvec`. +two APIs: `io_buffer_register_request()` and `io_buffer_unregister`. ublk adds IO command of `UBLK_IO_REGISTER_IO_BUF` to call -`io_buffer_register_bvec()` for ublk server to register client request +`io_buffer_register_request()` for ublk server to register client request buffer into io_uring buffer table, then ublk server can submit io_uring IOs with the registered buffer index. IO command of `UBLK_IO_UNREGISTER_IO_BUF` -calls `io_buffer_unregister_bvec()` to unregister the buffer, which is -guaranteed to be live between calling `io_buffer_register_bvec()` and -`io_buffer_unregister_bvec()`. Any io_uring operation which supports this -kind of kernel buffer will grab one reference of the buffer until the -operation is completed. +calls `io_buffer_unregister()` to unregister the buffer, which is guaranteed +to be live between calling `io_buffer_register_request()` and +`io_buffer_unregister()`. Any io_uring operation which supports this kind of +kernel buffer will grab one reference of the buffer until the operation is +completed. ublk server implementing zero copy or user copy has to be CAP_SYS_ADMIN and be trusted, because it is ublk server's responsibility to make sure IO buffer diff --git a/Documentation/filesystems/fuse/fuse-io-uring.rst b/Documentation/filesystems/fuse/fuse-io-uring.rst index d73dd0dbd238..29f98057500d 100644 --- a/Documentation/filesystems/fuse/fuse-io-uring.rst +++ b/Documentation/filesystems/fuse/fuse-io-uring.rst @@ -11,6 +11,9 @@ and works. For generic details about FUSE see fuse.rst. This document also covers the current interface, which is still in development and might change. +For the userspace protocol, see +Documentation/filesystems/fuse/uapi/fuse-uapi-io-uring.rst. + Limitations =========== As of now not all requests types are supported through io-uring, userspace @@ -95,5 +98,34 @@ Sending requests with CQEs | <fuse_unlink() | | <sys_unlink() | - - +Buffer pools +============ + +Without a buffer pool, every entry needs to pass a dedicated payload buffer +large enough for the maximum payload size. A buffer pool decouples entries +from payload buffers. The server hands the kernel one contiguous buffer pool +of memory and when the kernel sends the server a request, it indicates the +offset into the pool for that request's payload. Internally, the kernel is +able to manage/optimize the buffer pool memory however it likes. + +A server may also register the pool region with io_uring as a fixed buffer. +The backing pages are then pinned once, avoiding per-request pinning and +address translation. This also allows servers to use the same registered +buffers for subsequent backing store I/O through io-uring, keeping data +in the same pinned pages without additional pinning / mapping overhead. + +Zero-copy +========= + +Zero-copy lets the server read from / write to the client's pages (pinned +user pages for direct I/O, or page-cache folios for buffered I/O) without an +intermediary payload copy. This requires CAP_SYS_ADMIN privileges. + +When a fuse request arrives for a file that opted into zero-copy, the kernel +registers the relevant pages (pinned user pages for direct i/o or underlying +page cache folios for buffered i/o) into a sparse slot in the server's +io_uring registered buffer table. The server can then operate on these pages +directly using io-uring fixed buffer operations (eg read_fixed / write_fixed) +and the kernel unregisters these pages when the request completes. +Non-page-backed args (eg op out headers) will go through the payload buffer as +normal. diff --git a/Documentation/filesystems/fuse/index.rst b/Documentation/filesystems/fuse/index.rst index 393a845214da..3dada6c4057a 100644 --- a/Documentation/filesystems/fuse/index.rst +++ b/Documentation/filesystems/fuse/index.rst @@ -12,3 +12,4 @@ FUSE (Filesystem in Userspace) Technical Documentation fuse-io fuse-io-uring fuse-passthrough + uapi/fuse-uapi-io-uring diff --git a/Documentation/filesystems/fuse/uapi/fuse-uapi-io-uring.rst b/Documentation/filesystems/fuse/uapi/fuse-uapi-io-uring.rst new file mode 100644 index 000000000000..8367be7ea29d --- /dev/null +++ b/Documentation/filesystems/fuse/uapi/fuse-uapi-io-uring.rst @@ -0,0 +1,126 @@ +.. SPDX-License-Identifier: GPL-2.0 + +===================================== +FUSE-over-io-uring uapi documentation +===================================== + +Commands +======== + +``enum fuse_uring_cmd``: + +``FUSE_IO_URING_CMD_ADD_QUEUE`` + Create a queue identified by ``fuse_uring_cmd_req.qid``. Queue-wide + options are passed in ``fuse_uring_cmd_req.flags``: + + ``FUSE_URING_ZERO_COPY`` + Enable zero-copy on this queue. Requires ``CAP_SYS_ADMIN`` and a buffer + pool, which is added separately via ``ADD_BUFPOOL`` before registering + entries (see `Zero-copy`_). + +``FUSE_IO_URING_CMD_ADD_BUFPOOL`` + Register the payload buffer pool for an existing queue. The server provides + a single contiguous region in ``fuse_uring_cmd_req.bufpool.uaddr`` / + ``.len``. This command must be issued after ``ADD_QUEUE`` and before + registering any payload-carrying entries on that queue. + ``fuse_uring_cmd_req.flags`` must be 0. Submitting this command with + ``IORING_URING_CMD_FIXED`` marks the pool as registered, which avoids per + i/o pinning/unpinning and mapping overhead (see `Buffer pools`_). + +``FUSE_IO_URING_CMD_REGISTER`` + Register a ring entry (a long-lived SQE that carries the request header + iovec). For a zero-copy queue, ``fuse_uring_cmd_req.ent_zero_copy_buf_index`` + indicates the reserved registered buffer table slot this entry uses for + zero-copy (see `Zero-copy`_). + +``FUSE_IO_URING_CMD_COMMIT_AND_FETCH`` + Commit the reply for a completed request and fetch the next one. The + request is identified by ``fuse_uring_cmd_req.commit_id`` (the value the + kernel reported in ``fuse_uring_ent_in_out.commit_id``). + +Structures +========== + +``struct fuse_uring_cmd_req`` (80-byte SQE command area): + +============================ ================================================== +Field Meaning +============================ ================================================== +``flags`` Command-specific flags (see each command). +``commit_id`` Request id, for ``COMMIT_AND_FETCH``. +``qid`` Queue index. +``bufpool.uaddr`` Pool base address, for ``ADD_BUFPOOL``. +``bufpool.len`` Pool length in bytes, for ``ADD_BUFPOOL``. +``bufpool.reserved`` Must be 0, for ``ADD_BUFPOOL``. +``ent_zero_copy_buf_index`` Per-entry zero-copy slot, for ``REGISTER``. +============================ ================================================== + +``struct fuse_uring_ent_in_out`` (reported by the kernel per request): + +============================ ================================================== +Field Meaning +============================ ================================================== +``flags`` ``FUSE_URING_ENT_ZERO_COPY`` if zero-copied. +``commit_id`` Id to echo back in ``COMMIT_AND_FETCH``. +``payload_sz`` Total payload size in bytes (see `Zero-copy`_). +``offset`` Payload buffer offset within the pool. +============================ ================================================== + +Buffer pools +============ +Setup: + +* Issue ``ADD_QUEUE`` for the qid. +* Issue ``ADD_BUFPOOL`` with ``bufpool.uaddr`` and ``bufpool.len`` pointing + at the region. +* Register entries with ``REGISTER``. + +For every request that has a payload, the kernel reports where the payload +lives in ``struct fuse_uring_ent_in_out`` (part of +``struct fuse_uring_req_header``): + +``offset`` + Byte offset, within the pool region, for this request's payload buffer. + The server adds this to the pool base address to locate the payload. + +``payload_sz`` + Number of payload bytes for this request. + +To use registered buffers, the server registers the pool region with io_uring +and submits ``ADD_BUFPOOL`` with ``IORING_URING_CMD_FIXED`` set in +``sqe->uring_cmd_flags`` and the index of the registered bufpool in +``sqe->buf_index``. Every SQE the server submits afterwards must follow the +same fixed-buffer protocol, carrying ``IORING_URING_CMD_FIXED`` and that same +``sqe->buf_index``. The same registered buffer can be reused for the server's +backing-store I/O as well (e.g. ``IORING_OP_READ_FIXED`` / +``IORING_OP_WRITE_FIXED``). + +Zero-copy +========= +Requirements: + +* The server must be privileged (``CAP_SYS_ADMIN``). +* A zero-copy queue: ``ADD_QUEUE`` with the ``FUSE_URING_ZERO_COPY`` flag set. +* A buffer pool: ``ADD_BUFPOOL``. +* For each entry, ``REGISTER`` with ``ent_zero_copy_buf_index`` set to the + index this entry uses in the server's io_uring registered-buffer table. + This is where the kernel registers the request's pages for the server to + access (it is separate from the payload pool). On a non-zero-copy queue this + field must be 0. + +Zero-copy is selected per open file. The server sets the open-file flag in +the ``FUSE_OPEN`` / ``FUSE_CREATE`` reply: + +``FOPEN_IO_URING_ZERO_COPY`` + Reads/writes on this open file should use zero-copy. + +For a request that is zero-copied, the kernel sets ``FUSE_URING_ENT_ZERO_COPY`` +in ``fuse_uring_ent_in_out.flags`` and places the request's pages at the +entry's ``ent_zero_copy_buf_index``. The server then issues +``IORING_OP_READ_FIXED`` / ``IORING_OP_WRITE_FIXED`` against that index to +transfer the data directly to/from the client's pages. + +For such a request, ``payload_sz`` includes the zero-copied page bytes +(transferred via the registered buffer at ``ent_zero_copy_buf_index``). Any +non-page-backed args (e.g. op headers) are still copied through the pool +payload buffer at ``offset``. diff --git a/drivers/block/ublk_drv.c b/drivers/block/ublk_drv.c index 4d17ed264da1..6c5bec7da97c 100644 --- a/drivers/block/ublk_drv.c +++ b/drivers/block/ublk_drv.c @@ -1699,8 +1699,8 @@ ublk_auto_buf_register(const struct ublk_queue *ubq, struct request *req, { int ret; - ret = io_buffer_register_bvec(cmd, req, ublk_io_release, - io->buf.auto_reg.index, issue_flags); + ret = io_buffer_register_request(cmd, req, ublk_io_release, + io->buf.auto_reg.index, issue_flags); if (ret) { if (io->buf.auto_reg.flags & UBLK_AUTO_BUF_REG_FALLBACK) { ublk_auto_buf_reg_fallback(ubq, req->tag); @@ -1909,7 +1909,7 @@ static noinline void ublk_batch_dispatch_fail(struct ublk_queue *ubq, ublk_io_unlock(io); if (index != -1) - io_buffer_unregister_bvec(data->cmd, index, + io_buffer_unregister(data->cmd, index, data->issue_flags); } @@ -3192,8 +3192,8 @@ static int ublk_register_io_buf(struct io_uring_cmd *cmd, if (!req) return -EINVAL; - ret = io_buffer_register_bvec(cmd, req, ublk_io_release, index, - issue_flags); + ret = io_buffer_register_request(cmd, req, ublk_io_release, index, + issue_flags); if (ret) { ublk_put_req_ref(io, req); return ret; @@ -3224,8 +3224,8 @@ ublk_daemon_register_io_buf(struct io_uring_cmd *cmd, if (!ublk_dev_support_zero_copy(ub) || !blk_rq_has_data(req)) return -EINVAL; - ret = io_buffer_register_bvec(cmd, req, ublk_io_release, index, - issue_flags); + ret = io_buffer_register_request(cmd, req, ublk_io_release, index, + issue_flags); if (ret) return ret; @@ -3240,7 +3240,7 @@ static int ublk_unregister_io_buf(struct io_uring_cmd *cmd, if (!(ub->dev_info.flags & UBLK_F_SUPPORT_ZERO_COPY)) return -EINVAL; - return io_buffer_unregister_bvec(cmd, index, issue_flags); + return io_buffer_unregister(cmd, index, issue_flags); } static int ublk_check_fetch_buf(const struct ublk_device *ub, __u64 buf_addr) @@ -3384,7 +3384,7 @@ static int ublk_ch_uring_cmd_local(struct io_uring_cmd *cmd, goto out; /* - * io_buffer_unregister_bvec() doesn't access the ubq or io, + * io_buffer_unregister() doesn't access the ubq or io, * so no need to validate the q_id, tag, or task */ if (_IOC_NR(cmd_op) == UBLK_IO_UNREGISTER_IO_BUF) @@ -3456,7 +3456,7 @@ static int ublk_ch_uring_cmd_local(struct io_uring_cmd *cmd, req = ublk_fill_io_cmd(io, cmd); ublk_apply_io_buf(ub, io, cmd, addr, &auto_buf, &buf_idx); if (buf_idx != UBLK_INVALID_BUF_IDX) - io_buffer_unregister_bvec(cmd, buf_idx, issue_flags); + io_buffer_unregister(cmd, buf_idx, issue_flags); compl = ublk_need_complete_req(ub, io); if (req_op(req) == REQ_OP_ZONE_APPEND) @@ -3801,7 +3801,7 @@ static int ublk_batch_commit_io(struct ublk_queue *ubq, } if (buf_idx != UBLK_INVALID_BUF_IDX) - io_buffer_unregister_bvec(data->cmd, buf_idx, data->issue_flags); + io_buffer_unregister(data->cmd, buf_idx, data->issue_flags); if (req_op(req) == REQ_OP_ZONE_APPEND) req->__sector = ublk_batch_zone_lba(uc, elem); if (compl) diff --git a/fs/fuse/args.h b/fs/fuse/args.h index ecfe51a192af..5173264a1261 100644 --- a/fs/fuse/args.h +++ b/fs/fuse/args.h @@ -42,6 +42,8 @@ struct fuse_args { bool is_pinned:1; bool invalidate_vmap:1; bool abort_on_kill:1; + /* server requested io-uring zero-copy for this op */ + bool zero_copy:1; struct fuse_in_arg in_args[4]; struct fuse_arg out_args[2]; void (*end)(struct fuse_args *args, int error); diff --git a/fs/fuse/cuse.c b/fs/fuse/cuse.c index 3c15b5ba16d7..4079cf8e5974 100644 --- a/fs/fuse/cuse.c +++ b/fs/fuse/cuse.c @@ -530,7 +530,8 @@ static int cuse_channel_open(struct inode *inode, struct file *file) INIT_LIST_HEAD(&cc->list); - cc->fc.chan->initialized = 1; + /* Pairs with smp_load_acquire() readers of fch->initialized */ + smp_store_release(&cc->fc.chan->initialized, 1); rc = cuse_send_init(cc); if (rc) { fuse_dev_put(fud); @@ -653,6 +654,11 @@ static void __exit cuse_exit(void) { misc_deregister(&cuse_miscdev); class_destroy(cuse_class); + /* + * Wait for pending call_rcu() callbacks that call back into + * this module via fc->release (cuse_fc_release). + */ + rcu_barrier(); } module_init(cuse_init); diff --git a/fs/fuse/dev.c b/fs/fuse/dev.c index 5763a7cd3b37..4fec31fc0b84 100644 --- a/fs/fuse/dev.c +++ b/fs/fuse/dev.c @@ -75,17 +75,23 @@ void fuse_chan_set_initialized(struct fuse_chan *fch, struct fuse_chan_param *pa fch->minor = param->minor; fch->max_write = param->max_write; fch->max_pages = param->max_pages; + + if (param->io_uring_enabled) + fuse_uring_conn_init(fch); } - /* Make sure stores before this are seen on another CPU */ - smp_wmb(); - fch->initialized = 1; + /* Pairs with smp_load_acquire() readers of fch->initialized */ + smp_store_release(&fch->initialized, 1); wake_up_all(&fch->blocked_waitq); } static bool fuse_block_alloc(struct fuse_chan *fch, bool for_background) { - return !fch->initialized || (for_background && fch->blocked) || + /* Pairs with smp_store_release() in fuse_chan_set_initialized() */ + if (!smp_load_acquire(&fch->initialized)) + return true; + + return (for_background && fch->blocked) || (fch->io_uring && fch->connected && !fuse_uring_ready(fch)); } @@ -120,9 +126,6 @@ static struct fuse_req *fuse_get_req(struct fuse_chan *fch, bool for_background) goto out; } - /* Matches smp_wmb() in fuse_chan_set_initialized() */ - smp_rmb(); - err = -ENOTCONN; if (!fch->connected) goto out; @@ -210,10 +213,13 @@ EXPORT_SYMBOL_GPL(fuse_req_hash); /* * A new request is available, wake fiq->waitq */ -static void fuse_dev_wake_and_unlock(struct fuse_iqueue *fiq) +static void fuse_dev_wake_and_unlock(struct fuse_iqueue *fiq, bool sync) __releases(fiq->lock) { - wake_up(&fiq->waitq); + if (sync) + wake_up_sync(&fiq->waitq); + else + wake_up(&fiq->waitq); kill_fasync(&fiq->fasync, SIGIO, POLL_IN); spin_unlock(&fiq->lock); } @@ -230,7 +236,7 @@ void fuse_dev_queue_forget(struct fuse_iqueue *fiq, if (fiq->connected) { fiq->forget_list_tail->next = forget; fiq->forget_list_tail = forget; - fuse_dev_wake_and_unlock(fiq); + fuse_dev_wake_and_unlock(fiq, false); } else { kfree(forget); spin_unlock(&fiq->lock); @@ -240,7 +246,8 @@ void fuse_dev_queue_forget(struct fuse_iqueue *fiq, void fuse_dev_queue_interrupt(struct fuse_iqueue *fiq, struct fuse_req *req) { spin_lock(&fiq->lock); - if (list_empty(&req->intr_entry)) { + /* Repeat FR_SENT test after obtaining the lock to prevent race with fuse_resend() */ + if (list_empty(&req->intr_entry) && test_bit(FR_SENT, &req->flags)) { list_add_tail(&req->intr_entry, &fiq->interrupts); /* * Pairs with smp_mb() implied by test_and_set_bit() @@ -251,7 +258,7 @@ void fuse_dev_queue_interrupt(struct fuse_iqueue *fiq, struct fuse_req *req) list_del_init(&req->intr_entry); spin_unlock(&fiq->lock); } else { - fuse_dev_wake_and_unlock(fiq); + fuse_dev_wake_and_unlock(fiq, false); } } else { spin_unlock(&fiq->lock); @@ -281,11 +288,13 @@ EXPORT_SYMBOL_GPL(fuse_request_assign_unique); static void fuse_dev_queue_req(struct fuse_iqueue *fiq, struct fuse_req *req) { + bool sync = test_and_clear_bit(FR_SYNC_WAKEUP, &req->flags); + spin_lock(&fiq->lock); if (fiq->connected) { fuse_request_assign_unique_locked(fiq, req); list_add_tail(&req->list, &fiq->pending); - fuse_dev_wake_and_unlock(fiq); + fuse_dev_wake_and_unlock(fiq, sync); } else { spin_unlock(&fiq->lock); req->out.h.error = -ENOTCONN; @@ -397,7 +406,8 @@ void fuse_chan_max_background_set(struct fuse_chan *fch, unsigned int val) fch->max_background = val; fch->blocked = fch->num_background >= fch->max_background; if (!fch->blocked) - wake_up(&fch->blocked_waitq); + wake_up_nr(&fch->blocked_waitq, + fch->max_background - fch->num_background); spin_unlock(&fch->bg_lock); } @@ -411,11 +421,6 @@ void fuse_chan_set_fc(struct fuse_chan *fch, struct fuse_conn *fc) fch->conn = fc; } -void fuse_chan_io_uring_enable(struct fuse_chan *fch) -{ - fch->io_uring = 1; -} - void fuse_pqueue_init(struct fuse_pqueue *fpq) { spin_lock_init(&fpq->lock); @@ -725,7 +730,7 @@ static void request_wait_answer(struct fuse_req *req) if (req->args->abort_on_kill) { fuse_chan_abort(fch, false); - return; + goto wait_for_finish; } if (test_bit(FR_URING, &req->flags)) @@ -736,6 +741,7 @@ static void request_wait_answer(struct fuse_req *req) return; } +wait_for_finish: /* * Either request is already in userspace, or it was forced. * Wait it out. @@ -752,6 +758,11 @@ static void __fuse_request_send(struct fuse_req *req) /* acquire extra reference, since request is still needed after fuse_request_end() */ __fuse_get_request(req); + /* + * This is a synchronous request: the caller will block waiting for + * the answer. Hint the scheduler via wake_up_sync(). + */ + set_bit(FR_SYNC_WAKEUP, &req->flags); fuse_send_one(fiq, req); request_wait_answer(req); @@ -1249,11 +1260,25 @@ int fuse_copy_folio(struct fuse_copy_state *cs, struct folio **foliop, if (folio) { size = folio_size(folio); - if (zeroing && count < size) - folio_zero_range(folio, 0, size); + if (zeroing && count < size) { + /* + * When the copy is skipped the folio already holds the + * payload, so only the bytes outside [offset, offset + + * count) may be zeroed. + * + * Otherwise, the whole folio is cleared first so that a + * failed copy leaves zeros rather than stale folio + * contents. + */ + if (cs->skip_folio_copy) + folio_zero_segments(folio, 0, offset, + offset + count, size); + else + folio_zero_range(folio, 0, size); + } } - while (count) { + while (!cs->skip_folio_copy && count) { if (cs->write && cs->pipebufs && folio) { /* * Can't control lifetime of pipe buffers, so always @@ -1346,6 +1371,10 @@ int fuse_copy_args(struct fuse_copy_state *cs, unsigned numargs, for (i = 0; !err && i < numargs; i++) { struct fuse_arg *arg = &args[i]; if (i == numargs - 1 && argpages) + /* + * if cs->skip_folio_copy is set, this just does any + * needed zeroing. No copying is involved. + */ err = fuse_copy_folios(cs, arg->size, zeroing); else err = fuse_copy_one(cs, arg->value, arg->size); @@ -1760,7 +1789,7 @@ out: void fuse_chan_resend(struct fuse_chan *fch) { struct fuse_dev *fud; - struct fuse_req *req, *next; + struct fuse_req *req; struct fuse_iqueue *fiq = &fch->iq; LIST_HEAD(to_queue); unsigned int i; @@ -1775,24 +1804,20 @@ void fuse_chan_resend(struct fuse_chan *fch) struct fuse_pqueue *fpq = &fud->pq; spin_lock(&fpq->lock); - for (i = 0; i < FUSE_PQ_HASH_SIZE; i++) - list_splice_tail_init(&fpq->processing[i], &to_queue); + for (i = 0; i < FUSE_PQ_HASH_SIZE; i++) { + struct list_head *this_queue = &fpq->processing[i]; + + list_for_each_entry(req, this_queue, list) + clear_bit(FR_SENT, &req->flags); + list_splice_tail_init(this_queue, &to_queue); + } spin_unlock(&fpq->lock); } spin_unlock(&fch->lock); - list_for_each_entry_safe(req, next, &to_queue, list) { - set_bit(FR_PENDING, &req->flags); - clear_bit(FR_SENT, &req->flags); - /* mark the request as resend request */ - req->in.h.unique |= FUSE_UNIQUE_RESEND; - } - spin_lock(&fiq->lock); if (!fiq->connected) { spin_unlock(&fiq->lock); - list_for_each_entry(req, &to_queue, list) - clear_bit(FR_PENDING, &req->flags); fuse_dev_end_requests(&to_queue); return; } @@ -1801,12 +1826,16 @@ void fuse_chan_resend(struct fuse_chan *fch) * intr_entry on fiq->interrupts after the request is re-queued. */ list_for_each_entry(req, &to_queue, list) { + set_bit(FR_PENDING, &req->flags); + /* mark the request as resend request */ + req->in.h.unique |= FUSE_UNIQUE_RESEND; + if (test_bit(FR_INTERRUPTED, &req->flags)) list_del_init(&req->intr_entry); } /* iq and pq requests are both oldest to newest */ list_splice(&to_queue, &fiq->pending); - fuse_dev_wake_and_unlock(fiq); + fuse_dev_wake_and_unlock(fiq, false); } /* Look up request on processing list by unique ID */ @@ -1888,7 +1917,8 @@ static ssize_t fuse_dev_do_write(struct fuse_dev *fud, * initialized and connected state */ err = -EINVAL; - if (!fch->initialized || !fch->connected) + /* Pairs with smp_store_release() in fuse_chan_set_initialized() */ + if (!smp_load_acquire(&fch->initialized) || !fch->connected) goto copy_finish; /* Don't try to move folios (yet) */ diff --git a/fs/fuse/dev.h b/fs/fuse/dev.h index aed69fd14c41..8d25378c0918 100644 --- a/fs/fuse/dev.h +++ b/fs/fuse/dev.h @@ -22,6 +22,7 @@ struct fuse_chan_param { unsigned int minor; unsigned int max_write; unsigned int max_pages; + bool io_uring_enabled; }; struct fuse_chan *fuse_chan_new(void); @@ -34,7 +35,6 @@ void fuse_chan_max_background_set(struct fuse_chan *fch, unsigned int val); unsigned int fuse_chan_num_waiting(struct fuse_chan *fch); void fuse_chan_set_fc(struct fuse_chan *fch, struct fuse_conn *fc); void fuse_chan_set_initialized(struct fuse_chan *fch, struct fuse_chan_param *param); -void fuse_chan_io_uring_enable(struct fuse_chan *fch); ssize_t fuse_chan_send(struct fuse_chan *fch, struct fuse_args *args); int fuse_chan_send_bg(struct fuse_chan *fch, struct fuse_args *args, gfp_t gfp_flags); int fuse_chan_send_notify_reply(struct fuse_chan *fch, struct fuse_args *args, u64 unique); diff --git a/fs/fuse/dev_uring.c b/fs/fuse/dev_uring.c index 77c8cec43d9c..c6dd420c4034 100644 --- a/fs/fuse/dev_uring.c +++ b/fs/fuse/dev_uring.c @@ -9,6 +9,7 @@ #include "dev_uring_i.h" #include "fuse_trace.h" +#include <linux/bitmap.h> #include <linux/fs.h> #include <linux/io_uring/cmd.h> @@ -21,6 +22,8 @@ MODULE_PARM_DESC(enable_uring, #define FUSE_URING_IOV_HEADERS 0 #define FUSE_URING_IOV_PAYLOAD 1 +#define FUSE_URING_ADD_QUEUE_FLAGS (FUSE_URING_ZERO_COPY) + bool fuse_uring_enabled(void) { return enable_uring; @@ -30,6 +33,11 @@ struct fuse_uring_pdu { struct fuse_ring_ent *ent; }; +struct fuse_zero_copy_bvs { + unsigned int nr_bvs; + struct bio_vec bvs[]; +}; + static const struct fuse_iqueue_ops fuse_io_uring_ops; enum fuse_uring_header_type { @@ -41,6 +49,32 @@ enum fuse_uring_header_type { FUSE_URING_HEADER_RING_ENT, }; +static inline bool bufpool_enabled(struct fuse_ring_queue *queue) +{ + return queue->payload_mode == FUSE_PAYLOAD_BUFPOOL; +} + +static inline bool bufpool_registered(struct fuse_ring_queue *queue) +{ + return queue->bufpool && queue->bufpool->registered; +} + +/* + * For a registered bufpool, every sqe that drives a payload import (REGISTER, + * COMMIT_AND_FETCH) must carry the registered buffer index of the pool. + * This also must be called from the command's issue handler, where cmd->sqe is + * still valid + */ +static inline bool fuse_uring_cmd_index_ok(struct io_uring_cmd *cmd, + struct fuse_ring_queue *queue) +{ + if (!bufpool_registered(queue)) + return true; + + return (cmd->flags & IORING_URING_CMD_FIXED) && + READ_ONCE(cmd->sqe->buf_index) == queue->bufpool->registered_index; +} + static void uring_cmd_set_ring_ent(struct io_uring_cmd *cmd, struct fuse_ring_ent *ring_ent) { @@ -86,8 +120,36 @@ static void fuse_uring_flush_bg(struct fuse_ring_queue *queue) } } +static bool can_zero_copy_req(struct fuse_ring_ent *ent, struct fuse_req *req) +{ + struct fuse_args *args = req->args; + + if (!ent->queue->zero_copy || !args->zero_copy) + return false; + + if (args->opcode != FUSE_READ && args->opcode != FUSE_WRITE) + return false; + + return args->in_pages || args->out_pages; +} + +static void zero_copy_unregister(struct io_uring_cmd *cmd, + struct fuse_ring_ent *ent, + unsigned int issue_flags) +{ + if (ent->zero_copied) { + int err = io_buffer_unregister(cmd, ent->zero_copy_index, + issue_flags); + + if (err) + pr_warn_ratelimited("qid=%d zero-copy unregister failed: %d\n", + ent->queue->qid, err); + ent->zero_copied = false; + } +} + static void fuse_uring_req_end(struct fuse_ring_ent *ent, struct fuse_req *req, - int error) + int error, unsigned int issue_flags) { struct fuse_ring_queue *queue = ent->queue; struct fuse_ring *ring = queue->ring; @@ -107,6 +169,8 @@ static void fuse_uring_req_end(struct fuse_ring_ent *ent, struct fuse_req *req, spin_unlock(&queue->lock); + zero_copy_unregister(ent->cmd, ent, issue_flags); + if (error) req->out.h.error = error; @@ -204,7 +268,7 @@ void fuse_uring_destruct(struct fuse_chan *fch) return; for (qid = 0; qid < ring->nr_queues; qid++) { - struct fuse_ring_queue *queue = ring->queues[qid]; + struct fuse_ring_queue *queue = READ_ONCE(ring->queues[qid]); struct fuse_ring_ent *ent, *next; if (!queue) @@ -222,8 +286,9 @@ void fuse_uring_destruct(struct fuse_chan *fch) } kfree(queue->fpq.processing); + kfree(queue->bufpool); kfree(queue); - ring->queues[qid] = NULL; + WRITE_ONCE(ring->queues[qid], NULL); } kfree(ring->queues); @@ -238,7 +303,6 @@ static struct fuse_ring *fuse_uring_create(struct fuse_chan *fch) { struct fuse_ring *ring; size_t nr_queues = num_possible_cpus(); - struct fuse_ring *res = NULL; size_t max_payload_size; ring = kzalloc_obj(*ring, GFP_KERNEL_ACCOUNT); @@ -258,12 +322,6 @@ static struct fuse_ring *fuse_uring_create(struct fuse_chan *fch) spin_unlock(&fch->lock); goto out_err; } - if (fch->ring) { - /* race, another thread created the ring in the meantime */ - spin_unlock(&fch->lock); - res = fch->ring; - goto out_err; - } init_waitqueue_head(&ring->stop_waitq); @@ -278,11 +336,18 @@ static struct fuse_ring *fuse_uring_create(struct fuse_chan *fch) out_err: kfree(ring->queues); kfree(ring); - return res; + return NULL; +} + +void fuse_uring_conn_init(struct fuse_chan *fch) +{ + if (fuse_uring_create(fch)) + fch->io_uring = 1; } static struct fuse_ring_queue *fuse_uring_create_queue(struct fuse_ring *ring, - int qid) + int qid, bool zero_copy, + bool fail_if_exists) { struct fuse_chan *fch = ring->chan; struct fuse_ring_queue *queue; @@ -290,16 +355,17 @@ static struct fuse_ring_queue *fuse_uring_create_queue(struct fuse_ring *ring, queue = kzalloc_obj(*queue, GFP_KERNEL_ACCOUNT); if (!queue) - return NULL; |
