Skip to content

Commit c1fc280

Browse files
bsberndhbirth
authored andcommitted
fuse: {io-uring}: Use a bitmap which queues are available
Signed-off-by: Bernd Schubert <bschubert@ddn.com> (cherry picked from commit c7fc2de)
1 parent 9449bfe commit c1fc280

2 files changed

Lines changed: 109 additions & 0 deletions

File tree

fs/fuse/dev_uring.c

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,9 @@ static inline void io_uring_cmd_private_sz_check(size_t cmd_sz)
4343
)
4444
#endif
4545

46+
/* Number of queued fuse requests until a queue is considered full */
47+
#define FUSE_URING_QUEUE_THRESHOLD 5
48+
4649
bool fuse_uring_enabled(void)
4750
{
4851
return enable_uring;
@@ -244,6 +247,23 @@ static void fuse_uring_destruct_q_masks(struct fuse_ring *ring)
244247
}
245248
}
246249

250+
static void fuse_ring_destruct_q_masks(struct fuse_ring *ring)
251+
{
252+
free_cpumask_var(ring->avail_q_mask);
253+
if (ring->per_numa_avail_q_mask) {
254+
for (int node = 0; node < ring->nr_numa_nodes; node++)
255+
free_cpumask_var(ring->per_numa_avail_q_mask[node]);
256+
kfree(ring->per_numa_avail_q_mask);
257+
}
258+
259+
free_cpumask_var(ring->registered_q_mask);
260+
if (ring->numa_registered_q_mask) {
261+
for (int node = 0; node < ring->nr_numa_nodes; node++)
262+
free_cpumask_var(ring->numa_registered_q_mask[node]);
263+
kfree(ring->numa_registered_q_mask);
264+
}
265+
}
266+
247267
void fuse_uring_destruct(struct fuse_conn *fc)
248268
{
249269
struct fuse_ring *ring = fc->ring;
@@ -279,6 +299,7 @@ void fuse_uring_destruct(struct fuse_conn *fc)
279299
}
280300

281301
fuse_uring_destruct_q_masks(ring);
302+
fuse_ring_destruct_q_masks(ring);
282303
kfree(ring->queues);
283304
kfree(ring);
284305
fc->ring = NULL;
@@ -319,6 +340,38 @@ static int fuse_uring_create_q_masks(struct fuse_ring *ring, size_t nr_queues)
319340
return 0;
320341
}
321342

343+
static int fuse_ring_create_q_masks(struct fuse_ring *ring, int nr_queues)
344+
{
345+
if (!zalloc_cpumask_var(&ring->avail_q_mask, GFP_KERNEL_ACCOUNT))
346+
return -ENOMEM;
347+
348+
if (!zalloc_cpumask_var(&ring->registered_q_mask, GFP_KERNEL_ACCOUNT))
349+
return -ENOMEM;
350+
351+
ring->per_numa_avail_q_mask = kmalloc_array(ring->nr_numa_nodes,
352+
sizeof(struct cpumask *),
353+
GFP_KERNEL_ACCOUNT);
354+
if (!ring->per_numa_avail_q_mask)
355+
return -ENOMEM;
356+
for (int node = 0; node < ring->nr_numa_nodes; node++)
357+
if (!zalloc_cpumask_var(&ring->per_numa_avail_q_mask[node],
358+
GFP_KERNEL_ACCOUNT))
359+
return -ENOMEM;
360+
361+
ring->numa_registered_q_mask = kmalloc_array(ring->nr_numa_nodes,
362+
sizeof(struct cpumask *),
363+
GFP_KERNEL_ACCOUNT);
364+
if (!ring->numa_registered_q_mask)
365+
return -ENOMEM;
366+
for (int node = 0; node < ring->nr_numa_nodes; node++) {
367+
if (!zalloc_cpumask_var(&ring->numa_registered_q_mask[node],
368+
GFP_KERNEL_ACCOUNT))
369+
return -ENOMEM;
370+
}
371+
372+
return 0;
373+
}
374+
322375
/*
323376
* Basic ring setup for this connection based on the provided configuration
324377
*/
@@ -348,6 +401,10 @@ static struct fuse_ring *fuse_uring_create(struct fuse_conn *fc)
348401
if (err)
349402
goto out_err;
350403

404+
err = fuse_ring_create_q_masks(ring, nr_queues);
405+
if (err)
406+
goto out_err;
407+
351408
spin_lock(&fc->lock);
352409
if (fc->ring) {
353410
/* race, another thread created the ring in the meantime */
@@ -368,6 +425,7 @@ static struct fuse_ring *fuse_uring_create(struct fuse_conn *fc)
368425

369426
out_err:
370427
fuse_uring_destruct_q_masks(ring);
428+
fuse_ring_destruct_q_masks(ring);
371429
kfree(ring->queues);
372430
kfree(ring);
373431
return res;
@@ -415,6 +473,7 @@ static struct fuse_ring_queue *fuse_uring_create_queue(struct fuse_ring *ring,
415473

416474
queue->qid = qid;
417475
queue->ring = ring;
476+
queue->numa_node = cpu_to_node(qid);
418477
spin_lock_init(&queue->lock);
419478

420479
INIT_LIST_HEAD(&queue->ent_avail_queue);
@@ -625,6 +684,16 @@ void fuse_uring_stop_queues(struct fuse_ring *ring)
625684

626685
fuse_uring_abort_end_queue_requests(queue);
627686
fuse_uring_teardown_entries(queue);
687+
688+
cpumask_clear_cpu(qid, ring->registered_q_mask);
689+
cpumask_clear_cpu(qid, ring->avail_q_mask);
690+
for (node = 0; node < ring->nr_numa_nodes; node++) {
691+
/* Clear the queue from all masks */
692+
cpumask_clear_cpu(qid,
693+
ring->numa_registered_q_mask[node]);
694+
cpumask_clear_cpu(qid,
695+
ring->per_numa_avail_q_mask[node]);
696+
}
628697
}
629698

630699
/* Reset all queue masks, we won't process any more IO */
@@ -973,9 +1042,18 @@ static int fuse_uring_send_next_to_ring(struct fuse_ring_ent *ent,
9731042
static void fuse_uring_ent_avail(struct fuse_ring_ent *ent,
9741043
struct fuse_ring_queue *queue)
9751044
{
1045+
struct fuse_ring *ring = queue->ring;
1046+
int node = queue->numa_node;
1047+
9761048
WARN_ON_ONCE(!ent->cmd);
9771049
list_move(&ent->list, &queue->ent_avail_queue);
9781050
ent->state = FRRS_AVAILABLE;
1051+
1052+
if (list_is_singular(&queue->ent_avail_queue) &&
1053+
queue->nr_reqs <= FUSE_URING_QUEUE_THRESHOLD) {
1054+
cpumask_set_cpu(queue->qid, ring->avail_q_mask);
1055+
cpumask_set_cpu(queue->qid, ring->per_numa_avail_q_mask[node]);
1056+
}
9791057
}
9801058

9811059
/* Used to find the request on SQE commit */
@@ -998,6 +1076,8 @@ static void fuse_uring_add_req_to_ring_ent(struct fuse_ring_ent *ent,
9981076
struct fuse_req *req)
9991077
{
10001078
struct fuse_ring_queue *queue = ent->queue;
1079+
struct fuse_ring *ring = queue->ring;
1080+
int node = queue->numa_node;
10011081

10021082
lockdep_assert_held(&queue->lock);
10031083

@@ -1012,6 +1092,16 @@ static void fuse_uring_add_req_to_ring_ent(struct fuse_ring_ent *ent,
10121092
ent->state = FRRS_FUSE_REQ;
10131093
list_move_tail(&ent->list, &queue->ent_w_req_queue);
10141094
fuse_uring_add_to_pq(ent, req);
1095+
1096+
/*
1097+
* If there are no more available entries, mark the queue as unavailable
1098+
* in both global and per-NUMA node masks
1099+
*/
1100+
if (list_empty(&queue->ent_avail_queue)) {
1101+
cpumask_clear_cpu(queue->qid, ring->avail_q_mask);
1102+
cpumask_clear_cpu(queue->qid,
1103+
ring->per_numa_avail_q_mask[node]);
1104+
}
10151105
}
10161106

10171107
/* Fetch the next fuse request if available */
@@ -1401,6 +1491,10 @@ static int fuse_uring_register(struct io_uring_cmd *cmd,
14011491
/* Marks the ring entry as ready */
14021492
fuse_uring_next_fuse_req(ent, queue, issue_flags);
14031493

1494+
cpumask_set_cpu(queue->qid, ring->registered_q_mask);
1495+
cpumask_set_cpu(queue->qid,
1496+
ring->numa_registered_q_mask[queue->numa_node]);
1497+
14041498
return 0;
14051499
}
14061500

fs/fuse/dev_uring_i.h

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,9 @@ struct fuse_ring_queue {
7070
/* queue id, corresponds to the cpu core */
7171
unsigned int qid;
7272

73+
/* NUMA node this queue belongs to */
74+
int numa_node;
75+
7376
/*
7477
* queue lock, taken when any value in the queue changes _and_ also
7578
* a ring entry state changes.
@@ -149,6 +152,18 @@ struct fuse_ring {
149152
/* all queue tracking */
150153
struct fuse_queue_map q_map;
151154

155+
/* Tracks which queues are available (empty) globally */
156+
cpumask_var_t avail_q_mask;
157+
158+
/* Tracks which queues are available per NUMA node */
159+
cpumask_var_t *per_numa_avail_q_mask;
160+
161+
/* Tracks which queues are registered */
162+
cpumask_var_t registered_q_mask;
163+
164+
/* Tracks which queues are registered per NUMA node */
165+
cpumask_var_t *numa_registered_q_mask;
166+
152167
wait_queue_head_t stop_waitq;
153168

154169
/* async tear down */

0 commit comments

Comments
 (0)