Class: UBLK::Native::Server

Inherits:
Object
  • Object
show all
Defined in:
ext/ublk/server.c

Instance Method Summary collapse

Constructor Details

#initialize(id, queues, depth) ⇒ Object



79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
# File 'ext/ublk/server.c', line 79

static VALUE server_initialize(VALUE self, VALUE id, VALUE queues, VALUE depth)
{
  ublk_server *server;
  char path[64];
  unsigned index;

  TypedData_Get_Struct(self, ublk_server, &server_type, server);
  server->queues = NUM2UINT(queues);
  server->depth = NUM2UINT(depth);
  if (server->queues == 0 || server->queues > UBLK_MAX_NR_QUEUES) rb_raise(rb_eArgError, "invalid queue count");
  if (server->depth == 0 || server->depth > UBLK_MAX_QUEUE_DEPTH) rb_raise(rb_eArgError, "invalid queue depth");

  snprintf(path, sizeof(path), "/dev/ublkc%u", NUM2UINT(id));
  server->fd = open(path, O_RDWR | O_CLOEXEC);
  if (server->fd < 0) rb_sys_fail(path);
  server->pid = getpid();
  server->queue = ALLOC_N(ublk_queue, server->queues);
  memset(server->queue, 0, server->queues * sizeof(ublk_queue));
  for (index = 0; index < server->queues; index++) server->queue[index].event_fd = -1;
  return self;
}

Instance Method Details

#closeObject



360
361
362
363
364
365
366
367
368
369
370
371
372
373
# File 'ext/ublk/server.c', line 360

static VALUE server_close(VALUE self)
{
  ublk_server *server;
  unsigned index;
  uint64_t value = 1;
  TypedData_Get_Struct(self, ublk_server, &server_type, server);
  if (atomic_exchange(&server->closed, 1)) return Qnil;
  for (index = 0; index < server->queues; index++) {
    if (server->queue[index].event_fd >= 0 &&
        write(server->queue[index].event_fd, &value, sizeof(value)) < 0 && errno != EAGAIN)
      rb_sys_fail("eventfd write");
  }
  return Qnil;
}

#releaseObject



375
376
377
378
379
380
381
382
383
# File 'ext/ublk/server.c', line 375

static VALUE server_release(VALUE self)
{
  ublk_server *server;
  TypedData_Get_Struct(self, ublk_server, &server_type, server);
  ublk_check_pid(server->pid);
  if (!atomic_load(&server->closed)) rb_raise(rb_eRuntimeError, "close ublk server before releasing it");
  server_release_resources(server);
  return Qnil;
}

#run(queue_id, target, ready_queue) ⇒ Object



324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
# File 'ext/ublk/server.c', line 324

static VALUE server_run(VALUE self, VALUE queue_id, VALUE target, VALUE ready_queue)
{
  ublk_server *server = get_server(self);
  unsigned qid = NUM2UINT(queue_id);
  ublk_queue *queue;

  if (qid >= server->queues) rb_raise(rb_eArgError, "invalid queue id");
  queue = &server->queue[qid];
  if (queue->ready) rb_raise(rb_eRuntimeError, "queue is already running");
  queue_setup(server, qid);
  rb_funcall(ready_queue, id_push, 1, queue_id);

  while (!atomic_load(&server->closed)) {
    struct wait_args wait = {&queue->ring, NULL, 0};
    uint64_t data;
    unsigned tag;
    int result;

    rb_thread_call_without_gvl(wait_without_gvl, &wait, RUBY_UBF_IO, NULL);
    if (wait.result < 0) rb_syserr_fail(-wait.result, "io_uring_wait_cqe");
    data = io_uring_cqe_get_data64(wait.cqe);
    result = wait.cqe->res;
    io_uring_cqe_seen(&queue->ring, wait.cqe);
    if (data == UBLK_STOP_DATA || result == -ENODEV) break;
    if (result < 0) rb_syserr_fail(-result, "ublk fetch request");

    tag = (unsigned)data - 1;
    if (tag >= server->depth) rb_raise(rb_eRuntimeError, "kernel returned an invalid ublk tag");
    result = process_request(server, qid, tag, target);
    prep_io_command(server, queue, qid, tag, UBLK_IO_COMMIT_AND_FETCH_REQ, result);
    result = submit_ring(&queue->ring);
    if (result < 0) rb_syserr_fail(-result, "io_uring_submit(commit)");
  }
  return Qnil;
}