31
#include "qemu/error-report.h"
32
#include "qemu/main-loop.h"
33
#include "system/block-backend.h"
34
+#include "system/iothread.h"
35
36
#include <fuse.h>
37
#include <fuse_lowlevel.h>
38
39
#include "standard-headers/linux/fuse.h"
40
+#include <sys/ioctl.h>
41
42
#if defined(CONFIG_FALLOCATE_ZERO_RANGE)
43
#include <linux/falloc.h>
120
sizeof(((FuseRequestInHeaderBuf *)0)->tail) !=
121
sizeof(FuseRequestInHeader));
122
121
-typedef struct FuseExport {
122
- BlockExport common;
123
+typedef struct FuseExport FuseExport;
124
124
- struct fuse_session *fuse_session;
125
- unsigned int in_flight; /* atomic */
126
- bool mounted, fd_handler_set_up;
125
+/*
126
+ * One FUSE "queue", representing one FUSE FD from which requests are fetched
127
+ * and processed. Each queue is tied to an AioContext.
128
+ */
129
+typedef struct FuseQueue {
130
+ FuseExport *exp;
131
+
132
+ AioContext *ctx;
133
+ int fuse_fd;
134
135
/*
136
* Cached buffer to receive the data of WRITE requests. Cached because:
147
* via blk_blockalign() and thus need to be freed via qemu_vfree().
148
*/
149
void *req_write_data_cached;
150
+} FuseQueue;
151
+
152
+struct FuseExport {
153
+ BlockExport common;
154
+
155
+ struct fuse_session *fuse_session;
156
+ unsigned int in_flight; /* atomic */
157
+ bool mounted, fd_handler_set_up;
158
159
/*
160
* Set when there was an unrecoverable error and no requests should be read
163
*/
164
bool halted;
165
151
- int fuse_fd;
166
+ int num_queues;
167
+ FuseQueue *queues;
168
+ /*
169
+ * True if this export should follow the generic export's AioContext.
170
+ * Will be false if the queues' AioContexts have been explicitly set by the
171
+ * user, i.e. are expected to stay in those contexts.
172
+ * (I.e. is always false if there is more than one queue.)
173
+ */
174
+ bool follow_aio_context;
175
176
char *mountpoint;
177
bool writable;
183
mode_t st_mode;
184
uid_t st_uid;
185
gid_t st_gid;
163
-} FuseExport;
186
+};
187
188
/*
189
* Verify that the size of FuseRequestInHeaderBuf.head plus the data
202
static void init_exports_table(void);
203
204
static int mount_fuse_export(FuseExport *exp, Error **errp);
205
+static int clone_fuse_fd(int fd, Error **errp);
206
207
static bool is_regular_file(const char *path, Error **errp);
208
209
static void read_from_fuse_fd(void *opaque);
210
static void coroutine_fn
187
-fuse_co_process_request(FuseExport *exp, const FuseRequestInHeader *in_hdr,
211
+fuse_co_process_request(FuseQueue *q, const FuseRequestInHeader *in_hdr,
212
const void *data_buffer);
213
static int fuse_write_err(int fd, const struct fuse_in_header *in_hdr, int err);
214
240
return;
241
}
242
219
- aio_set_fd_handler(exp->common.ctx, exp->fuse_fd,
220
- read_from_fuse_fd, NULL, NULL, NULL, exp);
243
+ for (int i = 0; i < exp->num_queues; i++) {
244
+ aio_set_fd_handler(exp->queues[i].ctx, exp->queues[i].fuse_fd,
245
+ read_from_fuse_fd, NULL, NULL, NULL,
246
+ &exp->queues[i]);
247
+ }
248
exp->fd_handler_set_up = true;
249
}
250
253
*/
254
static void fuse_detach_handlers(FuseExport *exp)
255
{
229
- aio_set_fd_handler(exp->common.ctx, exp->fuse_fd,
230
- NULL, NULL, NULL, NULL, NULL);
256
+ for (int i = 0; i < exp->num_queues; i++) {
257
+ aio_set_fd_handler(exp->queues[i].ctx, exp->queues[i].fuse_fd,
258
+ NULL, NULL, NULL, NULL, NULL);
259
+ }
260
exp->fd_handler_set_up = false;
261
}
262
271
272
/* Refresh AioContext in case it changed */
273
exp->common.ctx = blk_get_aio_context(exp->common.blk);
274
+ if (exp->follow_aio_context) {
275
+ assert(exp->num_queues == 1);
276
+ exp->queues[0].ctx = exp->common.ctx;
277
+ }
278
+
279
fuse_attach_handlers(exp);
280
}
281
307
assert(blk_exp_args->type == BLOCK_EXPORT_TYPE_FUSE);
308
309
if (multithread) {
276
- error_setg(errp, "FUSE export does not support multi-threading");
277
- return -EINVAL;
310
+ /* Guaranteed by common export code */
311
+ assert(mt_count >= 1);
312
+
313
+ exp->follow_aio_context = false;
314
+ exp->num_queues = mt_count;
315
+ exp->queues = g_new(FuseQueue, mt_count);
316
+
317
+ for (size_t i = 0; i < mt_count; i++) {
318
+ exp->queues[i] = (FuseQueue) {
319
+ .exp = exp,
320
+ .ctx = multithread[i],
321
+ .fuse_fd = -1,
322
+ };
323
+ }
324
+ } else {
325
+ /* Guaranteed by common export code */
326
+ assert(mt_count == 0);
327
+
328
+ exp->follow_aio_context = true;
329
+ exp->num_queues = 1;
330
+ exp->queues = g_new(FuseQueue, 1);
331
+ exp->queues[0] = (FuseQueue) {
332
+ .exp = exp,
333
+ .ctx = exp->common.ctx,
334
+ .fuse_fd = -1,
335
+ };
336
}
337
338
/* For growable and writable exports, take the RESIZE permission */
344
ret = blk_set_perm(exp->common.blk, blk_perm | BLK_PERM_RESIZE,
345
blk_shared_perm, errp);
346
if (ret < 0) {
289
- return ret;
347
+ goto fail;
348
}
349
}
350
420
421
g_hash_table_insert(exports, g_strdup(exp->mountpoint), NULL);
422
365
- exp->fuse_fd = fuse_session_fd(exp->fuse_session);
366
- ret = qemu_fcntl_addfl(exp->fuse_fd, O_NONBLOCK);
423
+ assert(exp->num_queues >= 1);
424
+ exp->queues[0].fuse_fd = fuse_session_fd(exp->fuse_session);
425
+ ret = qemu_fcntl_addfl(exp->queues[0].fuse_fd, O_NONBLOCK);
426
if (ret < 0) {
427
error_setg_errno(errp, -ret, "Failed to make FUSE FD non-blocking");
428
goto fail;
429
}
430
431
+ for (int i = 1; i < exp->num_queues; i++) {
432
+ int fd = clone_fuse_fd(exp->queues[0].fuse_fd, errp);
433
+ if (fd < 0) {
434
+ ret = fd;
435
+ goto fail;
436
+ }
437
+ exp->queues[i].fuse_fd = fd;
438
+ }
439
+
440
fuse_attach_handlers(exp);
441
return 0;
442
529
/**
530
* Allocate a buffer to receive WRITE data, or take the cached one.
531
*/
464
-static void *get_write_data_buffer(FuseExport *exp)
532
+static void *get_write_data_buffer(FuseQueue *q)
533
{
466
- if (exp->req_write_data_cached) {
467
- void *cached = exp->req_write_data_cached;
468
- exp->req_write_data_cached = NULL;
534
+ if (q->req_write_data_cached) {
535
+ void *cached = q->req_write_data_cached;
536
+ q->req_write_data_cached = NULL;
537
return cached;
538
} else {
471
- return blk_blockalign(exp->common.blk, FUSE_MAX_WRITE_BYTES);
539
+ return blk_blockalign(q->exp->common.blk, FUSE_MAX_WRITE_BYTES);
540
}
541
}
542
543
/**
544
* Release a WRITE data buffer, possibly reusing it for a subsequent request.
545
*/
478
-static void release_write_data_buffer(FuseExport *exp, void **buffer)
546
+static void release_write_data_buffer(FuseQueue *q, void **buffer)
547
{
548
if (!*buffer) {
549
return;
550
}
551
484
- if (!exp->req_write_data_cached) {
485
- exp->req_write_data_cached = *buffer;
552
+ if (!q->req_write_data_cached) {
553
+ q->req_write_data_cached = *buffer;
554
} else {
555
qemu_vfree(*buffer);
556
}
596
}
597
}
598
599
+/**
600
+ * Clone the given /dev/fuse file descriptor, yielding a second FD from which
601
+ * requests can be pulled for the associated filesystem. Returns an FD on
602
+ * success, and -errno on error.
603
+ */
604
+static int clone_fuse_fd(int fd, Error **errp)
605
+{
606
+ uint32_t src_fd = fd;
607
+ int new_fd;
608
+ int ret;
609
+
610
+ /*
611
+ * The name "/dev/fuse" is fixed, see libfuse's lib/fuse_loop_mt.c
612
+ * (fuse_clone_chan()).
613
+ */
614
+ new_fd = open("/dev/fuse", O_RDWR | O_CLOEXEC | O_NONBLOCK);
615
+ if (new_fd < 0) {
616
+ ret = -errno;
617
+ error_setg_errno(errp, errno, "Failed to open /dev/fuse");
618
+ return ret;
619
+ }
620
+
621
+ ret = ioctl(new_fd, FUSE_DEV_IOC_CLONE, &src_fd);
622
+ if (ret < 0) {
623
+ ret = -errno;
624
+ error_setg_errno(errp, errno, "Failed to clone FUSE FD");
625
+ close(new_fd);
626
+ return ret;
627
+ }
628
+
629
+ return new_fd;
630
+}
631
+
632
/**
633
* Try to read a single request from the FUSE FD.
533
- * Takes a FuseExport pointer in `opaque`.
634
+ * Takes a FuseQueue pointer in `opaque`.
635
*
636
* Assumes the export's in-flight counter has already been incremented.
637
*
639
*/
640
static void coroutine_fn co_read_from_fuse_fd(void *opaque)
641
{
541
- FuseExport *exp = opaque;
542
- int fuse_fd = exp->fuse_fd;
642
+ FuseQueue *q = opaque;
643
+ int fuse_fd = q->fuse_fd;
644
+ FuseExport *exp = q->exp;
645
ssize_t ret;
646
FuseRequestInHeaderBuf in_hdr_buf;
647
const FuseRequestInHeader *in_hdr;
653
goto no_request;
654
}
655
554
- data_buffer = get_write_data_buffer(exp);
656
+ data_buffer = get_write_data_buffer(q);
657
658
/* Construct the I/O vector to hold the FUSE request */
659
iov[0] = (struct iovec) { &in_hdr_buf.head, sizeof(in_hdr_buf.head) };
714
memcpy(in_hdr_buf.tail, data_buffer, len);
715
}
716
615
- release_write_data_buffer(exp, &data_buffer);
717
+ release_write_data_buffer(q, &data_buffer);
718
}
719
618
- fuse_co_process_request(exp, in_hdr, data_buffer);
720
+ fuse_co_process_request(q, in_hdr, data_buffer);
721
722
no_request:
621
- release_write_data_buffer(exp, &data_buffer);
723
+ release_write_data_buffer(q, &data_buffer);
724
fuse_dec_in_flight(exp);
725
}
726
727
/**
728
* Try to read and process a single request from the FUSE FD.
729
* (To be used as a handler for when the FUSE FD becomes readable.)
628
- * Takes a FuseExport pointer in `opaque`.
730
+ * Takes a FuseQueue pointer in `opaque`.
731
*/
732
static void read_from_fuse_fd(void *opaque)
733
{
632
- FuseExport *exp = opaque;
734
+ FuseQueue *q = opaque;
735
Coroutine *co;
736
635
- co = qemu_coroutine_create(co_read_from_fuse_fd, exp);
737
+ co = qemu_coroutine_create(co_read_from_fuse_fd, q);
738
/* Decremented by co_read_from_fuse_fd() */
637
- fuse_inc_in_flight(exp);
739
+ fuse_inc_in_flight(q->exp);
740
qemu_coroutine_enter(co);
741
}
742
761
{
762
FuseExport *exp = container_of(blk_exp, FuseExport, common);
763
764
+ for (int i = 0; i < exp->num_queues; i++) {
765
+ FuseQueue *q = &exp->queues[i];
766
+
767
+ /* Queue 0's FD belongs to the FUSE session */
768
+ if (i > 0 && q->fuse_fd >= 0) {
769
+ close(q->fuse_fd);
770
+ }
771
+ qemu_vfree(q->req_write_data_cached);
772
+ }
773
+ g_free(exp->queues);
774
+
775
if (exp->fuse_session) {
776
if (exp->mounted) {
777
fuse_session_unmount(exp->fuse_session);
780
fuse_session_destroy(exp->fuse_session);
781
}
782
670
- qemu_vfree(exp->req_write_data_cached);
783
g_free(exp->mountpoint);
784
}
785
1456
* Process a FUSE request, incl. writing the response.
1457
*/
1458
static void coroutine_fn
1347
-fuse_co_process_request(FuseExport *exp, const FuseRequestInHeader *in_hdr,
1459
+fuse_co_process_request(FuseQueue *q, const FuseRequestInHeader *in_hdr,
1460
const void *data_buffer)
1461
{
1462
FuseRequestOutHeader out_hdr;
1463
+ FuseExport *exp = q->exp;
1464
/* For read requests: Data to be returned */
1465
void *out_data_buffer = NULL;
1466
ssize_t ret;
1584
}
1585
1586
if (out_data_buffer) {
1474
- fuse_write_buf_response(exp->fuse_fd, &out_hdr.common, out_data_buffer);
1587
+ fuse_write_buf_response(q->fuse_fd, &out_hdr.common, out_data_buffer);
1588
qemu_vfree(out_data_buffer);
1589
} else {
1477
- fuse_write_response(exp->fuse_fd, &out_hdr);
1590
+ fuse_write_response(q->fuse_fd, &out_hdr);
1591
}
1592
}
1593