14#include "fuse_kernel.h"
15#include "fuse_uring_i.h"
19#include <sys/sysinfo.h>
28#include <linux/sched.h>
30#include <sys/eventfd.h>
33#define FUSE_URING_MAX_SQE128_CMD_DATA 80
36 struct fuse_ring_queue *ring_queue;
39 struct fuse_uring_req_header *req_header;
41 size_t req_payload_sz;
44 uint64_t req_commit_id;
46 enum fuse_uring_cmd last_cmd;
52struct fuse_ring_queue {
54 struct fuse_ring_pool *ring_pool;
62 pthread_mutex_t ring_lock;
66 struct fuse_ring_ent ent[];
73 struct fuse_session *se;
82 size_t max_req_payload_sz;
85 size_t queue_mem_size;
87 unsigned int started_threads;
88 unsigned int failed_threads;
93 pthread_cond_t thread_start_cond;
94 pthread_mutex_t thread_start_mutex;
97 struct fuse_ring_queue *queues;
101fuse_ring_queue_size(
const size_t q_depth)
103 const size_t req_size =
sizeof(
struct fuse_ring_ent) * q_depth;
105 return sizeof(
struct fuse_ring_queue) + req_size;
108static struct fuse_ring_queue *
112 ((
char *)fuse_ring->queues) + (qid * fuse_ring->queue_mem_size);
120static void *fuse_uring_get_sqe_cmd(
struct io_uring_sqe *sqe)
122 return (
void *)&sqe->cmd[0];
126 const unsigned int qid,
127 const uint64_t commit_id)
130 req->commit_id = commit_id;
135fuse_uring_sqe_prepare(
struct io_uring_sqe *sqe,
struct fuse_ring_ent *req,
139 sqe->opcode = IORING_OP_URING_CMD;
145 sqe->flags = IOSQE_FIXED_FILE;
152 io_uring_sqe_set_data(sqe, req);
154 sqe->cmd_op = cmd_op;
159 struct fuse_ring_queue *queue,
160 struct fuse_ring_ent *ring_ent)
163 struct fuse_session *se = ring_pool->se;
165 struct fuse_out_header *out = (
struct fuse_out_header *)&rrh->in_out;
166 struct fuse_uring_ent_in_out *ent_in_out =
167 (
struct fuse_uring_ent_in_out *)&rrh->ring_ent_in_out;
168 struct io_uring_sqe *sqe;
170 if (pthread_self() != queue->tid) {
171 pthread_mutex_lock(&queue->ring_lock);
175 sqe = io_uring_get_sqe(&queue->ring);
185 fuse_log(FUSE_LOG_ERR,
"Failed to get a ring SQEs\n");
190 ring_ent->last_cmd = FUSE_IO_URING_CMD_COMMIT_AND_FETCH;
191 fuse_uring_sqe_prepare(sqe, ring_ent, ring_ent->last_cmd);
192 fuse_uring_sqe_set_req_data(fuse_uring_get_sqe_cmd(sqe), queue->qid,
193 ring_ent->req_commit_id);
196 fuse_log(FUSE_LOG_DEBUG,
" unique: %" PRIu64
", result=%d\n",
197 out->unique, ent_in_out->payload_sz);
200 if (!queue->cqe_processing)
201 io_uring_submit(&queue->ring);
204 pthread_mutex_unlock(&queue->ring_lock);
212 struct fuse_ring_ent *ring_ent;
215 if (!req->flags.is_uring)
218 ring_ent = container_of(req,
struct fuse_ring_ent, req);
220 *payload = ring_ent->op_payload;
221 *payload_sz = ring_ent->req_payload_sz;
233int send_reply_uring(
fuse_req_t req,
int error,
const void *arg,
size_t argsize)
236 struct fuse_ring_ent *ring_ent =
237 container_of(req,
struct fuse_ring_ent, req);
239 struct fuse_out_header *out = (
struct fuse_out_header *)&rrh->in_out;
240 struct fuse_uring_ent_in_out *ent_in_out =
241 (
struct fuse_uring_ent_in_out *)&rrh->ring_ent_in_out;
243 struct fuse_ring_queue *queue = ring_ent->ring_queue;
245 size_t max_payload_sz = ring_pool->max_req_payload_sz;
247 if (argsize > max_payload_sz) {
248 fuse_log(FUSE_LOG_ERR,
"argsize %zu exceeds buffer size %zu",
249 argsize, max_payload_sz);
251 }
else if (argsize) {
252 if (arg != ring_ent->op_payload)
253 memcpy(ring_ent->op_payload, arg, argsize);
255 ent_in_out->payload_sz = argsize;
258 out->unique = req->unique;
260 res = fuse_uring_commit_sqe(ring_pool, queue, ring_ent);
270 struct fuse_ring_ent *ring_ent =
271 container_of(req,
struct fuse_ring_ent, req);
273 struct fuse_ring_queue *queue = ring_ent->ring_queue;
276 struct fuse_out_header *out = (
struct fuse_out_header *)&rrh->in_out;
277 struct fuse_uring_ent_in_out *ent_in_out =
278 (
struct fuse_uring_ent_in_out *)&rrh->ring_ent_in_out;
279 size_t max_payload_sz = ring_ent->req_payload_sz;
280 struct fuse_bufvec dest_vec = FUSE_BUFVEC_INIT(max_payload_sz);
283 dest_vec.
buf[0].
mem = ring_ent->op_payload;
284 dest_vec.
buf[0].
size = max_payload_sz;
288 out->error = res < 0 ? res : 0;
289 out->unique = req->unique;
291 ent_in_out->payload_sz = res > 0 ? res : 0;
293 res = fuse_uring_commit_sqe(ring_pool, queue, ring_ent);
305 struct fuse_ring_ent *ring_ent =
306 container_of(req,
struct fuse_ring_ent, req);
308 struct fuse_ring_queue *queue = ring_ent->ring_queue;
311 struct fuse_out_header *out = (
struct fuse_out_header *)&rrh->in_out;
312 struct fuse_uring_ent_in_out *ent_in_out =
313 (
struct fuse_uring_ent_in_out *)&rrh->ring_ent_in_out;
314 size_t max_buf = ring_pool->max_req_payload_sz;
319 for (
int idx = 1; idx < count; idx++) {
320 struct iovec *cur = &iov[idx];
322 if (len + cur->iov_len > max_buf) {
324 "iov[%d] exceeds buffer size %zu",
330 memcpy(ring_ent->op_payload + len, cur->iov_base, cur->iov_len);
334 ent_in_out->payload_sz = len;
337 out->unique = req->unique;
340 return fuse_uring_commit_sqe(ring_pool, queue, ring_ent);
343static int fuse_queue_setup_io_uring(
struct io_uring *ring,
size_t qid,
344 size_t depth,
int fd,
int evfd)
347 struct io_uring_params params = {0};
348 int files[2] = { fd, evfd };
352 params.flags = IORING_SETUP_SQE128;
357 params.flags |= IORING_SETUP_SUBMIT_ALL;
360 params.flags |= IORING_SETUP_CQSIZE;
361 params.cq_entries = depth * 2;
368 params.flags |= IORING_SETUP_SINGLE_ISSUER;
371 params.flags |= IORING_SETUP_TASKRUN_FLAG;
374 params.flags |= IORING_SETUP_COOP_TASKRUN;
377 rc = io_uring_queue_init_params(depth, ring, ¶ms);
379 fuse_log(FUSE_LOG_ERR,
"Failed to setup qid %zu: %d (%s)\n",
380 qid, rc, strerror(-rc));
384 rc = io_uring_register_files(ring, files, 1);
388 "Failed to register files for ring idx %zu: %s",
389 qid, strerror(errno));
396static void fuse_session_destruct_uring(
struct fuse_ring_pool *fuse_ring)
398 for (
size_t qid = 0; qid < fuse_ring->nr_queues; qid++) {
399 struct fuse_ring_queue *queue =
400 fuse_uring_get_queue(fuse_ring, qid);
402 if (queue->tid != 0) {
403 uint64_t value = 1ULL;
406 rc = write(queue->eventfd, &value,
sizeof(value));
407 if (rc !=
sizeof(value))
409 "Wrote to eventfd=%d err=%s: rc=%d\n",
410 queue->eventfd, strerror(errno), rc);
411 pthread_cancel(queue->tid);
412 pthread_join(queue->tid, NULL);
416 if (queue->eventfd >= 0) {
417 close(queue->eventfd);
421 if (queue->ring.ring_fd != -1)
422 io_uring_queue_exit(&queue->ring);
424 for (
size_t idx = 0; idx < fuse_ring->queue_depth; idx++) {
425 struct fuse_ring_ent *ent = &queue->ent[idx];
427 numa_free(ent->op_payload, ent->req_payload_sz);
428 numa_free(ent->req_header, queue->req_header_sz);
431 pthread_mutex_destroy(&queue->ring_lock);
434 free(fuse_ring->queues);
435 pthread_cond_destroy(&fuse_ring->thread_start_cond);
436 pthread_mutex_destroy(&fuse_ring->thread_start_mutex);
440static int fuse_uring_register_ent(
struct fuse_ring_queue *queue,
441 struct fuse_ring_ent *ent)
443 struct io_uring_sqe *sqe;
445 sqe = io_uring_get_sqe(&queue->ring);
451 fuse_log(FUSE_LOG_ERR,
"Failed to get all ring SQEs");
455 ent->last_cmd = FUSE_IO_URING_CMD_REGISTER;
456 fuse_uring_sqe_prepare(sqe, ent, ent->last_cmd);
459 ent->iov[0].iov_base = ent->req_header;
460 ent->iov[0].iov_len = queue->req_header_sz;
462 ent->iov[1].iov_base = ent->op_payload;
463 ent->iov[1].iov_len = ent->req_payload_sz;
465 sqe->addr = (uint64_t)(ent->iov);
469 fuse_uring_sqe_set_req_data(fuse_uring_get_sqe_cmd(sqe), queue->qid, 0);
475static int fuse_uring_register_queue(
struct fuse_ring_queue *queue)
478 unsigned int sq_ready;
479 struct io_uring_sqe *sqe;
482 for (
size_t idx = 0; idx < ring_pool->queue_depth; idx++) {
483 struct fuse_ring_ent *ent = &queue->ent[idx];
485 res = fuse_uring_register_ent(queue, ent);
490 sq_ready = io_uring_sq_ready(&queue->ring);
491 if (sq_ready != ring_pool->queue_depth) {
493 "SQE ready mismatch, expected %zu got %u\n",
494 ring_pool->queue_depth, sq_ready);
499 sqe = io_uring_get_sqe(&queue->ring);
501 fuse_log(FUSE_LOG_ERR,
"Failed to get eventfd SQE");
505 io_uring_prep_poll_add(sqe, queue->eventfd, POLLIN);
506 io_uring_sqe_set_data(sqe, (
void *)(uintptr_t)queue->eventfd);
513static struct fuse_ring_pool *fuse_create_ring(
struct fuse_session *se)
516 const size_t nr_queues = get_nprocs_conf();
517 size_t payload_sz = se->bufsize - FUSE_BUFFER_HEADER_SIZE;
521 fuse_log(FUSE_LOG_DEBUG,
"starting io-uring q-depth=%d\n",
524 fuse_ring = calloc(1,
sizeof(*fuse_ring));
525 if (fuse_ring == NULL) {
526 fuse_log(FUSE_LOG_ERR,
"Allocating the ring failed\n");
530 queue_sz = fuse_ring_queue_size(se->uring.q_depth);
531 fuse_ring->queues = calloc(1, queue_sz * nr_queues);
532 if (fuse_ring->queues == NULL) {
533 fuse_log(FUSE_LOG_ERR,
"Allocating the queues failed\n");
538 fuse_ring->nr_queues = nr_queues;
539 fuse_ring->queue_depth = se->uring.q_depth;
540 fuse_ring->max_req_payload_sz = payload_sz;
541 fuse_ring->queue_mem_size = queue_sz;
548 for (
size_t qid = 0; qid < nr_queues; qid++) {
549 struct fuse_ring_queue *queue =
550 fuse_uring_get_queue(fuse_ring, qid);
552 queue->ring.ring_fd = -1;
553 queue->numa_node = numa_node_of_cpu(qid);
555 queue->ring_pool = fuse_ring;
557 pthread_mutex_init(&queue->ring_lock, NULL);
560 pthread_cond_init(&fuse_ring->thread_start_cond, NULL);
561 pthread_mutex_init(&fuse_ring->thread_start_mutex, NULL);
562 sem_init(&fuse_ring->init_sem, 0, 0);
568 fuse_session_destruct_uring(fuse_ring);
573static void fuse_uring_resubmit(
struct fuse_ring_queue *queue,
574 struct fuse_ring_ent *ent)
576 struct io_uring_sqe *sqe;
578 sqe = io_uring_get_sqe(&queue->ring);
586 queue->ring_pool->se->error = -EIO;
587 fuse_log(FUSE_LOG_ERR,
"Failed to get a ring SQEs\n");
592 fuse_uring_sqe_prepare(sqe, ent, ent->last_cmd);
594 switch (ent->last_cmd) {
595 case FUSE_IO_URING_CMD_REGISTER:
596 sqe->addr = (uint64_t)(ent->iov);
598 fuse_uring_sqe_set_req_data(fuse_uring_get_sqe_cmd(sqe),
601 case FUSE_IO_URING_CMD_COMMIT_AND_FETCH:
602 fuse_uring_sqe_set_req_data(fuse_uring_get_sqe_cmd(sqe),
603 queue->qid, ent->req_commit_id);
606 fuse_log(FUSE_LOG_ERR,
"Unknown command type: %d\n",
608 queue->ring_pool->se->error = -EINVAL;
615static void fuse_uring_handle_cqe(
struct fuse_ring_queue *queue,
616 struct io_uring_cqe *cqe)
618 struct fuse_ring_ent *ent = io_uring_cqe_get_data(cqe);
622 "cqe=%p io_uring_cqe_get_data returned NULL\n", cqe);
626 struct fuse_req *req = &ent->req;
630 struct fuse_in_header *in = (
struct fuse_in_header *)&rrh->in_out;
631 struct fuse_uring_ent_in_out *ent_in_out = &rrh->ring_ent_in_out;
633 ent->req_commit_id = ent_in_out->commit_id;
634 if (unlikely(ent->req_commit_id == 0)) {
639 fuse_log(FUSE_LOG_ERR,
"Received invalid commit_id=0\n");
643 memset(&req->flags, 0,
sizeof(req->flags));
644 memset(&req->u, 0,
sizeof(req->u));
645 req->flags.is_uring = 1;
648 req->interrupted = 0;
651 fuse_session_process_uring_cqe(fuse_ring->se, req, in, &rrh->op_in,
652 ent->op_payload, ent_in_out->payload_sz);
655static int fuse_uring_queue_handle_cqes(
struct fuse_ring_queue *queue)
658 struct fuse_session *se = ring_pool->se;
659 size_t num_completed = 0;
660 struct io_uring_cqe *cqe;
662 struct fuse_ring_ent *ent;
665 io_uring_for_each_cqe(&queue->ring, head, cqe) {
671 if (unlikely(err != 0)) {
672 if (err > 0 && ((uintptr_t)io_uring_cqe_get_data(cqe) ==
673 (
unsigned int)queue->eventfd)) {
683 ent = io_uring_cqe_get_data(cqe);
684 fuse_uring_resubmit(queue, ent);
691 if (err != -ENOTCONN) {
692 se->error = cqe->res;
700 fuse_uring_handle_cqe(queue, cqe);
705 io_uring_cq_advance(&queue->ring, num_completed);
707 return ret == 0 ? 0 : num_completed;
714static void fuse_uring_set_thread_core(
int qid)
721 rc = sched_setaffinity(0,
sizeof(cpu_set_t), &mask);
723 fuse_log(FUSE_LOG_ERR,
"Failed to bind qid=%d to its core: %s\n",
724 qid, strerror(errno));
727 const int policy = SCHED_IDLE;
728 const struct sched_param param = {
729 .sched_priority = sched_get_priority_min(policy),
735 rc = sched_setscheduler(0, policy, ¶m);
737 fuse_log(FUSE_LOG_ERR,
"Failed to set scheduler: %s\n",
745static int fuse_uring_init_queue(
struct fuse_ring_queue *queue)
748 struct fuse_session *se = ring->se;
750 size_t page_sz = sysconf(_SC_PAGESIZE);
752 queue->eventfd = eventfd(0, EFD_CLOEXEC);
753 if (queue->eventfd < 0) {
756 "Failed to create eventfd for qid %d: %s\n",
757 queue->qid, strerror(errno));
761 res = fuse_queue_setup_io_uring(&queue->ring, queue->qid,
762 ring->queue_depth, se->fd,
765 fuse_log(FUSE_LOG_ERR,
"qid=%d io_uring init failed\n",
773 for (
size_t idx = 0; idx < ring->queue_depth; idx++) {
774 struct fuse_ring_ent *ring_ent = &queue->ent[idx];
775 struct fuse_req *req = &ring_ent->req;
777 ring_ent->ring_queue = queue;
783 ring_ent->req_header =
784 numa_alloc_local(queue->req_header_sz);
785 if (!ring_ent->req_header)
787 ring_ent->req_payload_sz = ring->max_req_payload_sz;
789 ring_ent->op_payload =
790 numa_alloc_local(ring_ent->req_payload_sz);
791 if (!ring_ent->op_payload)
795 pthread_mutex_init(&req->lock, NULL);
796 req->flags.is_uring = 1;
801 res = fuse_uring_register_queue(queue);
805 "Grave fuse-uring error on preparing SQEs, aborting\n");
811 return queue->ring.ring_fd;
814static void *fuse_uring_thread(
void *arg)
816 struct fuse_ring_queue *queue = arg;
818 struct fuse_session *se = ring_pool->se;
820 char thread_name[16] = { 0 };
822 snprintf(thread_name, 16,
"fuse-ring-%d", queue->qid);
823 thread_name[15] =
'\0';
824 fuse_set_thread_name(thread_name);
826 fuse_uring_set_thread_core(queue->qid);
828 err = fuse_uring_init_queue(queue);
829 pthread_mutex_lock(&ring_pool->thread_start_mutex);
831 ring_pool->failed_threads++;
832 ring_pool->started_threads++;
833 pthread_cond_broadcast(&ring_pool->thread_start_cond);
834 pthread_mutex_unlock(&ring_pool->thread_start_mutex);
837 fuse_log(FUSE_LOG_ERR,
"qid=%d queue setup failed\n",
842 sem_wait(&ring_pool->init_sem);
845 while (!atomic_load_explicit(&se->mt_exited, memory_order_relaxed)) {
846 io_uring_submit_and_wait(&queue->ring, 1);
848 pthread_mutex_lock(&queue->ring_lock);
849 queue->cqe_processing =
true;
850 err = fuse_uring_queue_handle_cqes(queue);
851 queue->cqe_processing =
false;
852 pthread_mutex_unlock(&queue->ring_lock);
865static int fuse_uring_start_ring_threads(
struct fuse_ring_pool *ring)
869 for (
size_t qid = 0; qid < ring->nr_queues; qid++) {
870 struct fuse_ring_queue *queue = fuse_uring_get_queue(ring, qid);
872 rc = pthread_create(&queue->tid, NULL, fuse_uring_thread, queue);
880static int fuse_uring_sanity_check(
struct fuse_session *se)
882 if (se->uring.q_depth == 0) {
883 fuse_log(FUSE_LOG_ERR,
"io-uring queue depth must be > 0\n");
888 FUSE_URING_MAX_SQE128_CMD_DATA,
889 "SQE128_CMD_DATA has 80B cmd data");
894int fuse_uring_start(
struct fuse_session *se)
899 fuse_uring_sanity_check(se);
901 fuse_ring = fuse_create_ring(se);
902 if (fuse_ring == NULL) {
903 err = -EADDRNOTAVAIL;
907 se->uring.pool = fuse_ring;
910 sem_init(&fuse_ring->init_sem, 0, 0);
911 pthread_cond_init(&fuse_ring->thread_start_cond, NULL);
912 pthread_mutex_init(&fuse_ring->thread_start_mutex, NULL);
914 err = fuse_uring_start_ring_threads(fuse_ring);
921 pthread_mutex_lock(&fuse_ring->thread_start_mutex);
922 while (fuse_ring->started_threads < fuse_ring->nr_queues)
923 pthread_cond_wait(&fuse_ring->thread_start_cond,
924 &fuse_ring->thread_start_mutex);
926 if (fuse_ring->failed_threads != 0)
927 err = -EADDRNOTAVAIL;
928 pthread_mutex_unlock(&fuse_ring->thread_start_mutex);
934 fuse_session_destruct_uring(fuse_ring);
935 se->uring.pool = NULL;
940int fuse_uring_stop(
struct fuse_session *se)
947 fuse_session_destruct_uring(ring);
952void fuse_uring_wake_ring_threads(
struct fuse_session *se)
957 for (
size_t qid = 0; qid < ring->nr_queues; qid++)
958 sem_post(&ring->init_sem);
ssize_t fuse_buf_copy(struct fuse_bufvec *dst, struct fuse_bufvec *src, enum fuse_buf_copy_flags flags)
void fuse_log(enum fuse_log_level level, const char *fmt,...) __attribute__((format(printf
void fuse_session_exit(struct fuse_session *se)
struct fuse_req * fuse_req_t
int fuse_req_get_payload(fuse_req_t req, char **payload, size_t *payload_sz, void **mr)