[umka_os] Introduce dedicated I/O thread
Not as nice as io_uring but portable.
This commit is contained in:
@@ -7,35 +7,127 @@
|
||||
Copyright (C) 2023 Ivan Baravy <dunkaist@gmail.com>
|
||||
*/
|
||||
|
||||
#include <stdatomic.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
#include <inttypes.h>
|
||||
#include "umka.h"
|
||||
#include "umkaio.h"
|
||||
#include "io_async.h"
|
||||
|
||||
#define IOT_QUEUE_DEPTH 1
|
||||
|
||||
struct iot_cmd iot_cmd_buf[IOT_QUEUE_DEPTH];
|
||||
|
||||
static void *
|
||||
thread_io(void *arg) {
|
||||
(void)arg;
|
||||
for (size_t i = 0; i < IOT_QUEUE_DEPTH; i++) {
|
||||
iot_cmd_buf[i].status = UMKA_CMD_STATUS_EMPTY;
|
||||
iot_cmd_buf[i].type = 0;
|
||||
pthread_cond_init(&iot_cmd_buf[i].iot_cond, NULL);
|
||||
pthread_mutex_init(&iot_cmd_buf[i].iot_mutex, NULL);
|
||||
pthread_mutex_lock(&iot_cmd_buf[i].iot_mutex);
|
||||
pthread_mutex_init(&iot_cmd_buf[i].mutex, NULL);
|
||||
}
|
||||
|
||||
struct iot_cmd *cmd = iot_cmd_buf;
|
||||
ssize_t ret;
|
||||
while (1) {
|
||||
pthread_cond_wait(&cmd->iot_cond, &cmd->iot_mutex);
|
||||
// status must be ready
|
||||
switch (cmd->type) {
|
||||
case IOT_CMD_TYPE_READ:
|
||||
ret = read(cmd->read.arg.fd, cmd->read.arg.buf, cmd->read.arg.count);
|
||||
atomic_store_explicit(&cmd->read.ret.val, ret, memory_order_release);
|
||||
break;
|
||||
case IOT_CMD_TYPE_WRITE:
|
||||
cmd->read.ret.val = write(cmd->read.arg.fd, cmd->read.arg.buf,
|
||||
cmd->read.arg.count);
|
||||
break;
|
||||
default:
|
||||
break;
|
||||
}
|
||||
|
||||
atomic_store_explicit(&cmd->status, UMKA_CMD_STATUS_DONE, memory_order_release);
|
||||
}
|
||||
|
||||
return NULL;
|
||||
}
|
||||
|
||||
static uint32_t
|
||||
io_async_submit_wait_test() {
|
||||
// appdata_t *app;
|
||||
// __asm__ __volatile__ ("":"=b"(app)::);
|
||||
// struct io_uring_queue *q = app->wait_param;
|
||||
int done = pthread_mutex_trylock(&iot_cmd_buf[0].mutex);
|
||||
return done;
|
||||
}
|
||||
|
||||
static uint32_t
|
||||
io_async_complete_wait_test() {
|
||||
// appdata_t *app;
|
||||
// __asm__ __volatile__ ("":"=b"(app)::);
|
||||
// struct io_uring_queue *q = app->wait_param;
|
||||
int status = atomic_load_explicit(&iot_cmd_buf[0].status, memory_order_acquire);
|
||||
return status == UMKA_CMD_STATUS_DONE;
|
||||
}
|
||||
|
||||
ssize_t
|
||||
io_async_read(int fd, void *buf, size_t count, void *arg) {
|
||||
(void)arg;
|
||||
|
||||
kos_wait_events(io_async_submit_wait_test, NULL);
|
||||
// status must be empty
|
||||
struct iot_cmd *cmd = iot_cmd_buf;
|
||||
cmd->read.arg.fd = fd;
|
||||
cmd->read.arg.buf = buf;
|
||||
cmd->read.arg.count = count;
|
||||
atomic_store_explicit(&cmd->status, UMKA_CMD_STATUS_READY, memory_order_release);
|
||||
|
||||
pthread_cond_signal(&cmd->iot_cond);
|
||||
kos_wait_events(io_async_complete_wait_test, NULL);
|
||||
|
||||
ssize_t res = atomic_load_explicit(&cmd->read.ret.val, memory_order_acquire);
|
||||
|
||||
atomic_store_explicit(&cmd->status, UMKA_CMD_STATUS_EMPTY, memory_order_release);
|
||||
pthread_mutex_unlock(&cmd->mutex);
|
||||
|
||||
return res;
|
||||
}
|
||||
|
||||
ssize_t
|
||||
io_async_write(int fd, const void *buf, size_t count, void *arg) {
|
||||
(void)fd;
|
||||
(void)buf;
|
||||
(void)count;
|
||||
(void)arg;
|
||||
return -1;
|
||||
}
|
||||
|
||||
struct umka_io *
|
||||
io_init(int *running) {
|
||||
struct umka_io *io = malloc(sizeof(struct umka_io));
|
||||
io->running = running;
|
||||
io->async = io_async_init();
|
||||
if (running) {
|
||||
pthread_create(&io->iot, NULL, thread_io, NULL);
|
||||
}
|
||||
return io;
|
||||
}
|
||||
|
||||
void
|
||||
io_close(struct umka_io *io) {
|
||||
io_async_close(io->async);
|
||||
free(io);
|
||||
}
|
||||
|
||||
ssize_t
|
||||
io_read(int fd, void *buf, size_t count, struct umka_io *io) {
|
||||
ssize_t res;
|
||||
if (!*io->running) {
|
||||
if (!io->running || !*io->running) {
|
||||
res = read(fd, buf, count);
|
||||
} else {
|
||||
res = io_async_read(fd, buf, count, io->async);
|
||||
res = io_async_read(fd, buf, count, NULL);
|
||||
}
|
||||
return res;
|
||||
}
|
||||
@@ -43,10 +135,10 @@ io_read(int fd, void *buf, size_t count, struct umka_io *io) {
|
||||
ssize_t
|
||||
io_write(int fd, const void *buf, size_t count, struct umka_io *io) {
|
||||
ssize_t res;
|
||||
if (!*io->running) {
|
||||
if (!io->running || !*io->running) {
|
||||
res = write(fd, buf, count);
|
||||
} else {
|
||||
res = io_async_write(fd, buf, count, io->async);
|
||||
res = io_async_write(fd, buf, count, NULL);
|
||||
}
|
||||
return res;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user