aos/lib/block_driver/block_driver.c

197 lines
6.3 KiB
C

#include <aos/ump_chan.h>
#include <aos/aos_rpc.h>
#include <aos/deferred.h>
#include <drivers/block_driver.h>
// #define BLOCK_DRIVER_CLIENT_PERFORMANCE
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
#include <aos/performance.h>
static struct performance_context pcontext;
#endif
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
#endif
struct block_driver_state {
struct ump_send_chan *send_chan;
struct ump_recv_chan *recv_chan;
bool request_ongoing;
struct block_driver_request *current_request;
void *read_buf;
errval_t result;
};
static struct block_driver_state state = {
.recv_chan = NULL,
.send_chan = NULL,
};
// use a waitset to make the calls blocking
struct waitset ws;
errval_t block_driver_init(void) {
errval_t err;
waitset_init(&ws);
// connect to the block driver
struct capref cap;
struct aos_rpc *rpc = aos_rpc_get_init_channel();
// retry several times, waiting inbetween in case the server is still starting up
for (int attempts = 0; attempts < 50; attempts++) {
err = aos_rpc_ump_connect(rpc, UMP_SERVER_BLOCK_DRIVER, &cap);
if (err_no(err) != LIB_ERR_UMP_NOT_REGISTERED) break;
barrelfish_usleep(100000);
}
if (err_is_fail(err)) return err;
// Create the two uni-directional channels
struct ump_send_chan *send_chan;
struct ump_recv_chan *recv_chan;
err = ump_chan_init(UMP_ROLE_CLIENT, &send_chan, &recv_chan, sizeof(struct block_driver_result), cap, &ws);
if (err_is_fail(err)) return err;
state.recv_chan = recv_chan;
state.send_chan = send_chan;
state.request_ongoing = false;
state.current_request = NULL;
state.read_buf = NULL;
state.result = SYS_ERR_OK;
return SYS_ERR_OK;
}
static void handle_payload(void *arg, size_t payload_size, void *payload) {
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
perf_add_now(&pcontext, "payload");
#endif
state.request_ongoing = false;
}
static void handle_response(void *arg, size_t header_size, void *header, size_t payload_size) {
struct block_driver_result *result = header;
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
perf_add_now(&pcontext, "header");
#endif
state.result = result->err;
if (state.current_request->action == BLOCK_DRIVER_ACTION_READ && err_is_ok(result->err)) {
if (payload_size != state.current_request->bytes) {
USER_PANIC("We should always receive as much as requested or an error");
}
// receive the payload in the buffer passed to the blocking function
ump_recv_payload(state.recv_chan, state.read_buf, handle_payload, NULL);
return;
} else if (state.current_request->action == BLOCK_DRIVER_ACTION_COUNT_HANDLES && err_is_ok(result->err)) {
if (payload_size != sizeof(size_t)) {
USER_PANIC("We should always receive a size_t result for this request or an error");
}
// receive the payload in the buffer passed to the blocking function
ump_recv_payload(state.recv_chan, state.read_buf, handle_payload, NULL);
return;
}
ump_recv_payload(state.recv_chan, NULL, handle_payload, NULL);
}
static void handle_send_completed(void *arg, struct ump_send_queue_entry *entry) {
free((void *)entry->header);
free(entry);
}
static errval_t send_blocking_request(enum block_driver_action action, uint64_t directory_entry_id, uint32_t sector_number, size_t offset, size_t size, size_t payload_size, const void *payload) {
assert(state.recv_chan != NULL);
assert(state.send_chan != NULL);
assert(!state.request_ongoing);
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
perf_init(&pcontext, "block_driver_client");
perf_add_now(&pcontext, "start");
#endif
errval_t err;
state.request_ongoing = true;
// send the read request
struct block_driver_request *request = malloc(sizeof(struct block_driver_request));
if (request == NULL) return LIB_ERR_MALLOC_FAIL;
struct ump_send_queue_entry *entry = malloc(sizeof(struct ump_send_queue_entry));
if (entry == NULL) return LIB_ERR_MALLOC_FAIL;
request->action = action;
request->block_number = sector_number;
request->bytes = size;
request->offset = offset;
request->directory_entry_id = directory_entry_id;
state.current_request = request;
assert(state.send_chan->send_queue_head == NULL);
ump_send(
state.send_chan,
entry,
sizeof(struct block_driver_request),
request,
payload_size,
payload,
handle_send_completed,
NULL
);
// register response handler
ump_recv_header(state.recv_chan, handle_response, NULL);
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
perf_add_now(&pcontext, "dispatching");
#endif
// wait for the response
while (state.request_ongoing) {
err = event_dispatch(&ws);
if (err_is_fail(err)) {
DEBUG_ERR(err, "in event_dispatch");
abort();
}
}
#ifdef BLOCK_DRIVER_CLIENT_PERFORMANCE
perf_add_now(&pcontext, "done");
perf_print(&pcontext);
#endif
return state.result;
}
errval_t block_driver_read_object(uint32_t sector_number, size_t offset, size_t size, void *dst) {
state.read_buf = dst;
return send_blocking_request(BLOCK_DRIVER_ACTION_READ, 0, sector_number, offset, size, 0, NULL);
}
errval_t block_driver_write_object(uint32_t sector_number, size_t offset, size_t size, const void *src) {
return send_blocking_request(BLOCK_DRIVER_ACTION_WRITE, 0, sector_number, offset, size, size, src);
}
errval_t block_driver_lock(void) {
return send_blocking_request(BLOCK_DRIVER_ACTION_LOCK, 0, 0, 0, 0, 0, NULL);
}
errval_t block_driver_unlock(void) {
return send_blocking_request(BLOCK_DRIVER_ACTION_UNLOCK, 0, 0, 0, 0, 0, NULL);
}
errval_t block_driver_register_handle(uint64_t directory_entry_id) {
return send_blocking_request(BLOCK_DRIVER_ACTION_REGISTER_HANDLE, directory_entry_id, 0, 0, 0, 0, NULL);
}
errval_t block_driver_unregister_handle(uint64_t directory_entry_id) {
return send_blocking_request(BLOCK_DRIVER_ACTION_UNREGISTER_HANDLE, directory_entry_id, 0, 0, 0, 0, NULL);
}
errval_t block_driver_count_handles(uint64_t directory_entry_id, size_t *count) {
state.read_buf = count;
return send_blocking_request(BLOCK_DRIVER_ACTION_COUNT_HANDLES, directory_entry_id, 0, 0, 0, 0, NULL);
}