diff --git a/src/blockstore/blockstore.cpp b/src/blockstore/blockstore.cpp index 74e1c306..15ca8d35 100644 --- a/src/blockstore/blockstore.cpp +++ b/src/blockstore/blockstore.cpp @@ -6,7 +6,7 @@ #include "blockstore_impl.h" #include "v1/impl.h" -blockstore_i* blockstore_i::create(blockstore_config_t & config, ring_loop_t *ringloop, timerfd_manager_t *tfd) +blockstore_i* blockstore_i::create(blockstore_config_t & config, ring_loop_i *ringloop, timerfd_manager_t *tfd) { auto meta_format = stoull_full(config["meta_format"]); if (meta_format == BLOCKSTORE_META_FORMAT_HEAP) diff --git a/src/blockstore/blockstore.h b/src/blockstore/blockstore.h index 222d2fee..536c1588 100644 --- a/src/blockstore/blockstore.h +++ b/src/blockstore/blockstore.h @@ -176,7 +176,7 @@ typedef std::map blockstore_config_t; class __attribute__((visibility("default"))) blockstore_i { public: - static blockstore_i* create(blockstore_config_t & config, ring_loop_t *ringloop, timerfd_manager_t *tfd); + static blockstore_i* create(blockstore_config_t & config, ring_loop_i *ringloop, timerfd_manager_t *tfd); virtual ~blockstore_i() = default; diff --git a/src/blockstore/blockstore_disk.cpp b/src/blockstore/blockstore_disk.cpp index 8e7d8fc2..726dc485 100644 --- a/src/blockstore/blockstore_disk.cpp +++ b/src/blockstore/blockstore_disk.cpp @@ -102,6 +102,15 @@ void blockstore_disk_t::parse_config(std::map & config disable_data_fsync = config["disable_data_fsync"] == "true" || config["disable_data_fsync"] == "1" || config["disable_data_fsync"] == "yes"; disable_meta_fsync = config["disable_meta_fsync"] == "true" || config["disable_meta_fsync"] == "1" || config["disable_meta_fsync"] == "yes"; disable_journal_fsync = config["disable_journal_fsync"] == "true" || config["disable_journal_fsync"] == "1" || config["disable_journal_fsync"] == "yes"; + if (mock_mode) + { + data_device_size = parse_size(config["data_device_size"]); + data_device_sect = parse_size(config["data_device_sect"]); + meta_device_size = parse_size(config["meta_device_size"]); + meta_device_sect = parse_size(config["meta_device_sect"]); + journal_device_size = parse_size(config["journal_device_size"]); + journal_device_sect = parse_size(config["journal_device_sect"]); + } // Validate if (!data_block_size) { @@ -328,12 +337,15 @@ static int bs_openmode(const std::string & mode) void blockstore_disk_t::open_data() { - data_fd = open(data_device.c_str(), bs_openmode(data_io) | O_RDWR); + data_fd = mock_mode ? MOCK_DATA_FD : open(data_device.c_str(), bs_openmode(data_io) | O_RDWR); if (data_fd == -1) { throw std::runtime_error("Failed to open data device "+data_device+": "+std::string(strerror(errno))); } - check_size(data_fd, &data_device_size, &data_device_sect, "data device"); + if (!mock_mode) + { + check_size(data_fd, &data_device_size, &data_device_sect, "data device"); + } if (disk_alignment % data_device_sect) { throw std::runtime_error( @@ -345,7 +357,7 @@ void blockstore_disk_t::open_data() { throw std::runtime_error("data_offset exceeds device size = "+std::to_string(data_device_size)); } - if (!disable_flock && flock(data_fd, LOCK_EX|LOCK_NB) != 0) + if (!mock_mode && !disable_flock && flock(data_fd, LOCK_EX|LOCK_NB) != 0) { throw std::runtime_error(std::string("Failed to lock data device: ") + strerror(errno)); } @@ -355,17 +367,20 @@ void blockstore_disk_t::open_meta() { if (meta_device != data_device || meta_io != data_io) { - meta_fd = open(meta_device.c_str(), bs_openmode(meta_io) | O_RDWR); + meta_fd = mock_mode ? MOCK_META_FD : open(meta_device.c_str(), bs_openmode(meta_io) | O_RDWR); if (meta_fd == -1) { throw std::runtime_error("Failed to open metadata device "+meta_device+": "+std::string(strerror(errno))); } - check_size(meta_fd, &meta_device_size, &meta_device_sect, "metadata device"); + if (!mock_mode) + { + check_size(meta_fd, &meta_device_size, &meta_device_sect, "metadata device"); + } if (meta_offset >= meta_device_size) { throw std::runtime_error("meta_offset exceeds device size = "+std::to_string(meta_device_size)); } - if (!disable_flock && meta_device != data_device && flock(meta_fd, LOCK_EX|LOCK_NB) != 0) + if (!mock_mode && !disable_flock && meta_device != data_device && flock(meta_fd, LOCK_EX|LOCK_NB) != 0) { throw std::runtime_error(std::string("Failed to lock metadata device: ") + strerror(errno)); } @@ -393,13 +408,20 @@ void blockstore_disk_t::open_journal() { if (journal_device != meta_device || journal_io != meta_io) { - journal_fd = open(journal_device.c_str(), bs_openmode(journal_io) | O_RDWR); + journal_fd = mock_mode ? MOCK_JOURNAL_FD : open(journal_device.c_str(), bs_openmode(journal_io) | O_RDWR); if (journal_fd == -1) { throw std::runtime_error("Failed to open journal device "+journal_device+": "+std::string(strerror(errno))); } - check_size(journal_fd, &journal_device_size, &journal_device_sect, "journal device"); - if (!disable_flock && journal_device != meta_device && flock(journal_fd, LOCK_EX|LOCK_NB) != 0) + if (!mock_mode) + { + check_size(journal_fd, &journal_device_size, &journal_device_sect, "journal device"); + } + if (journal_offset >= journal_device_size) + { + throw std::runtime_error("journal_offset exceeds device size = "+std::to_string(journal_device_size)); + } + if (!mock_mode && !disable_flock && journal_device != meta_device && flock(journal_fd, LOCK_EX|LOCK_NB) != 0) { throw std::runtime_error(std::string("Failed to lock journal device: ") + strerror(errno)); } @@ -425,12 +447,15 @@ void blockstore_disk_t::open_journal() void blockstore_disk_t::close_all() { - if (data_fd >= 0) - close(data_fd); - if (meta_fd >= 0 && meta_fd != data_fd) - close(meta_fd); - if (journal_fd >= 0 && journal_fd != meta_fd) - close(journal_fd); + if (!mock_mode) + { + if (data_fd >= 0) + close(data_fd); + if (meta_fd >= 0 && meta_fd != data_fd) + close(meta_fd); + if (journal_fd >= 0 && journal_fd != meta_fd) + close(journal_fd); + } data_fd = meta_fd = journal_fd = -1; } @@ -438,6 +463,10 @@ void blockstore_disk_t::close_all() // so it's not a big deal that we can only run it synchronously. int blockstore_disk_t::trim_data(std::function is_free) { + if (mock_mode) + { + return -EINVAL; + } int r = 0; uint64_t j = 0, i = 0; uint64_t discarded = 0; diff --git a/src/blockstore/blockstore_disk.h b/src/blockstore/blockstore_disk.h index 7bd89bb1..55de03cd 100644 --- a/src/blockstore/blockstore_disk.h +++ b/src/blockstore/blockstore_disk.h @@ -17,6 +17,10 @@ // Lower byte of checksum type is its length #define BLOCKSTORE_CSUM_CRC32C 0x104 +#define MOCK_DATA_FD 1000 +#define MOCK_META_FD 1001 +#define MOCK_JOURNAL_FD 1002 + class allocator_t; struct blockstore_disk_t @@ -66,6 +70,8 @@ struct blockstore_disk_t uint32_t clean_entry_bitmap_size = 0; uint32_t clean_entry_size = 0, clean_dyn_size = 0; // for meta_v1/2 + bool mock_mode = false; + void parse_config(std::map & config); void open_data(); void open_meta(); diff --git a/src/blockstore/blockstore_impl.cpp b/src/blockstore/blockstore_impl.cpp index 22583ec1..229753ca 100644 --- a/src/blockstore/blockstore_impl.cpp +++ b/src/blockstore/blockstore_impl.cpp @@ -5,11 +5,12 @@ #include "blockstore_internal.h" #include "crc32c.h" -blockstore_impl_t::blockstore_impl_t(blockstore_config_t & config, ring_loop_t *ringloop, timerfd_manager_t *tfd) +blockstore_impl_t::blockstore_impl_t(blockstore_config_t & config, ring_loop_i *ringloop, timerfd_manager_t *tfd, bool mock_mode) { assert(sizeof(blockstore_op_private_t) <= BS_OP_PRIVATE_DATA_SIZE); this->tfd = tfd; this->ringloop = ringloop; + dsk.mock_mode = mock_mode; ring_consumer.loop = [this]() { loop(); }; ringloop->register_consumer(&ring_consumer); initialized = 0; diff --git a/src/blockstore/blockstore_impl.h b/src/blockstore/blockstore_impl.h index 2462962e..4cf71a83 100644 --- a/src/blockstore/blockstore_impl.h +++ b/src/blockstore/blockstore_impl.h @@ -124,7 +124,7 @@ class blockstore_impl_t: public blockstore_i bool fsyncing_data = false; bool live = false, queue_stall = false; - ring_loop_t *ringloop; + ring_loop_i *ringloop; timerfd_manager_t *tfd; bool stop_sync_submitted; @@ -187,7 +187,7 @@ class blockstore_impl_t: public blockstore_i public: - blockstore_impl_t(blockstore_config_t & config, ring_loop_t *ringloop, timerfd_manager_t *tfd); + blockstore_impl_t(blockstore_config_t & config, ring_loop_i *ringloop, timerfd_manager_t *tfd, bool mock_mode = false); ~blockstore_impl_t(); void parse_config(blockstore_config_t & config); diff --git a/src/util/ringloop.cpp b/src/util/ringloop.cpp index 10af62f4..862ec715 100644 --- a/src/util/ringloop.cpp +++ b/src/util/ringloop.cpp @@ -8,6 +8,7 @@ #include +#include "malloc_or_die.h" #include "ringloop.h" ring_loop_t::ring_loop_t(int qd, bool multithreaded, bool sqe128) @@ -228,3 +229,230 @@ int ring_loop_t::register_eventfd() loop(); return ring_eventfd; } + +ring_loop_mock_t::ring_loop_mock_t(int qd, std::function submit_cb) +{ + this->submit_cb = std::move(submit_cb); + sqes.resize(qd); + ring_datas.resize(qd); + free_ring_datas.reserve(qd); + submit_ring_datas.reserve(qd); + completed_ring_datas.reserve(qd); + for (size_t i = 0; i < ring_datas.size(); i++) + { + free_ring_datas.push_back(ring_datas.data() + i); + } + in_loop = false; +} + +void ring_loop_mock_t::register_consumer(ring_consumer_t *consumer) +{ + unregister_consumer(consumer); + consumers.push_back(consumer); +} + +void ring_loop_mock_t::unregister_consumer(ring_consumer_t *consumer) +{ + for (int i = 0; i < consumers.size(); i++) + { + if (consumers[i] == consumer) + { + consumers.erase(consumers.begin()+i, consumers.begin()+i+1); + break; + } + } +} + +void ring_loop_mock_t::wakeup() +{ + loop_again = true; +} + +void ring_loop_mock_t::set_immediate(const std::function & cb) +{ + immediate_queue.push_back(cb); + wakeup(); +} + +unsigned ring_loop_mock_t::space_left() +{ + return free_ring_datas.size(); +} + +bool ring_loop_mock_t::has_work() +{ + return loop_again; +} + +bool ring_loop_mock_t::has_sendmsg_zc() +{ + return false; +} + +int ring_loop_mock_t::register_eventfd() +{ + return -1; +} + +io_uring_sqe* ring_loop_mock_t::get_sqe() +{ + if (free_ring_datas.size() == 0) + { + return NULL; + } + ring_data_t *d = free_ring_datas.back(); + free_ring_datas.pop_back(); + submit_ring_datas.push_back(d); + io_uring_sqe *sqe = &sqes[d - ring_datas.data()]; + *sqe = { 0 }; + io_uring_sqe_set_data(sqe, d); + return sqe; +} + +int ring_loop_mock_t::submit() +{ + for (size_t i = 0; i < submit_ring_datas.size(); i++) + { + submit_cb(&sqes[submit_ring_datas[i] - ring_datas.data()]); + } + submit_ring_datas.clear(); + return 0; +} + +int ring_loop_mock_t::wait() +{ + return 0; +} + +unsigned ring_loop_mock_t::save() +{ + return submit_ring_datas.size(); +} + +void ring_loop_mock_t::restore(unsigned sqe_tail) +{ + while (submit_ring_datas.size() > sqe_tail) + { + free_ring_datas.push_back(submit_ring_datas.back()); + submit_ring_datas.pop_back(); + } +} + +void ring_loop_mock_t::loop() +{ + if (in_loop) + { + return; + } + in_loop = true; + submit(); + while (completed_ring_datas.size()) + { + ring_data_t *d = completed_ring_datas.back(); + completed_ring_datas.pop_back(); + if (d->callback) + { + struct ring_data_t dl; + dl.iov = d->iov; + dl.res = d->res; + dl.more = dl.prev = false; + dl.callback.swap(d->callback); + free_ring_datas.push_back(d); + dl.callback(&dl); + } + else + { + fprintf(stderr, "Warning: empty callback in SQE\n"); + free_ring_datas.push_back(d); + } + } + do + { + loop_again = false; + for (int i = 0; i < consumers.size(); i++) + { + consumers[i]->loop(); + if (immediate_queue.size()) + { + immediate_queue2.swap(immediate_queue); + for (auto & cb: immediate_queue2) + cb(); + immediate_queue2.clear(); + } + } + } while (loop_again); + in_loop = false; +} + +void ring_loop_mock_t::mark_completed(ring_data_t *data) +{ + completed_ring_datas.push_back(data); + wakeup(); +} + +disk_mock_t::disk_mock_t(ring_loop_mock_t *loop, size_t size) +{ + this->loop = loop; + this->size = size; + this->data = (uint8_t*)malloc_or_die(size); +} + +disk_mock_t::~disk_mock_t() +{ + free(data); +} + +bool disk_mock_t::submit(io_uring_sqe *sqe) +{ + ring_data_t *userdata = (ring_data_t*)sqe->user_data; + if (sqe->opcode == IORING_OP_READV) + { + size_t off = sqe->off; + iovec *v = (iovec*)sqe->addr; + size_t n = sqe->len; + for (size_t i = 0; i < n; i++) + { + size_t cur = (off + v[i].iov_len > size ? size-off : v[i].iov_len); + if (trace) + printf("read %zu+%zu to %jx\n", off, cur, (uint64_t)v[i].iov_base); + memcpy(v[i].iov_base, data + off, cur); + off += v[i].iov_len; + } + userdata->res = off - sqe->off; + } + else if (sqe->opcode == IORING_OP_WRITEV) + { + // Simple "immediate" mode. We should also implement "buffered" mode + size_t off = sqe->off; + iovec *v = (iovec*)sqe->addr; + size_t n = sqe->len; + for (size_t i = 0; i < n; i++) + { + if (off >= size) + { + off = sqe->off - EINVAL; // :D + break; + } + size_t cur = (off + v[i].iov_len > size ? size-off : v[i].iov_len); + if (trace) + printf("write %zu+%zu from %jx\n", off, cur, (uint64_t)v[i].iov_base); + memcpy(data + off, v[i].iov_base, cur); + off += v[i].iov_len; + } + userdata->res = off - sqe->off; + } + else if (sqe->opcode == IORING_OP_FSYNC) + { + userdata->res = 0; + } + else + { + return false; + } + // Execution variability should also be introduced: + // 1) reads submitted in parallel to writes (not after completing the write) should return old or new data randomly + // 2) parallel operation completions should be delivered in random order + // 3) when fsync is enabled, write cache should be sometimes lost during a simulated power outage + loop->mark_completed(userdata); + return true; +} diff --git a/src/util/ringloop.h b/src/util/ringloop.h index c397461f..5996ded5 100644 --- a/src/util/ringloop.h +++ b/src/util/ringloop.h @@ -32,7 +32,27 @@ struct ring_consumer_t std::function loop; }; -class __attribute__((visibility("default"))) ring_loop_t +class __attribute__((visibility("default"))) ring_loop_i +{ +public: + virtual ~ring_loop_i() = default; + virtual void register_consumer(ring_consumer_t *consumer) = 0; + virtual void unregister_consumer(ring_consumer_t *consumer) = 0; + virtual int register_eventfd() = 0; + virtual io_uring_sqe* get_sqe() = 0; + virtual void set_immediate(const std::function & cb) = 0; + virtual int submit() = 0; + virtual int wait() = 0; + virtual unsigned space_left() = 0; + virtual bool has_work() = 0; + virtual bool has_sendmsg_zc() = 0; + virtual void loop() = 0; + virtual void wakeup() = 0; + virtual unsigned save() = 0; + virtual void restore(unsigned sqe_tail) = 0; +}; + +class __attribute__((visibility("default"))) ring_loop_t: public ring_loop_i { std::vector> immediate_queue, immediate_queue2; std::vector consumers; @@ -54,7 +74,7 @@ public: int register_eventfd(); io_uring_sqe* get_sqe(); - inline void set_immediate(const std::function cb) + inline void set_immediate(const std::function & cb) { immediate_queue.push_back(cb); wakeup(); @@ -84,3 +104,51 @@ public: unsigned save(); void restore(unsigned sqe_tail); }; + +class ring_loop_mock_t: public ring_loop_i +{ + std::vector> immediate_queue, immediate_queue2; + std::vector consumers; + std::vector sqes; + std::vector ring_datas; + std::vector free_ring_datas; + std::vector submit_ring_datas; + std::vector completed_ring_datas; + std::function submit_cb; + bool in_loop; + bool loop_again; + bool support_zc = false; + +public: + ring_loop_mock_t(int qd, std::function submit_cb); + + void register_consumer(ring_consumer_t *consumer); + void unregister_consumer(ring_consumer_t *consumer); + void wakeup(); + void set_immediate(const std::function & cb); + unsigned space_left(); + bool has_work(); + bool has_sendmsg_zc(); + + int register_eventfd(); + io_uring_sqe* get_sqe(); + int submit(); + int wait(); + void loop(); + unsigned save(); + void restore(unsigned sqe_tail); + + void mark_completed(ring_data_t *data); +}; + +class disk_mock_t +{ + uint8_t *data = NULL; + size_t size = 0; + ring_loop_mock_t *loop = NULL; +public: + bool trace = false; + disk_mock_t(ring_loop_mock_t *loop, size_t size); + ~disk_mock_t(); + bool submit(io_uring_sqe *sqe); +}; diff --git a/src/util/timerfd_manager.cpp b/src/util/timerfd_manager.cpp index cafb10c3..e64ebc88 100644 --- a/src/util/timerfd_manager.cpp +++ b/src/util/timerfd_manager.cpp @@ -4,6 +4,7 @@ #include #include #include +#include #include #include #include @@ -15,21 +16,27 @@ timerfd_manager_t::timerfd_manager_t(std::functionset_fd_handler = set_fd_handler; wait_state = 0; - timerfd = timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK); - if (timerfd < 0) + if (set_fd_handler) { - throw std::runtime_error(std::string("timerfd_create: ") + strerror(errno)); + timerfd = timerfd_create(CLOCK_MONOTONIC, TFD_NONBLOCK); + if (timerfd < 0) + { + throw std::runtime_error(std::string("timerfd_create: ") + strerror(errno)); + } + set_fd_handler(timerfd, false, [this](int fd, int events) + { + handle_readable(); + }); } - set_fd_handler(timerfd, false, [this](int fd, int events) - { - handle_readable(); - }); } timerfd_manager_t::~timerfd_manager_t() { - set_fd_handler(timerfd, false, NULL); - close(timerfd); + if (timerfd >= 0) + { + set_fd_handler(timerfd, false, NULL); + close(timerfd); + } } void timerfd_manager_t::inc_timer(timerfd_timer_t & t) @@ -52,7 +59,14 @@ int timerfd_manager_t::set_timer_us(uint64_t micros, bool repeat, std::function< { int timer_id = id++; timespec start; - clock_gettime(CLOCK_MONOTONIC, &start); + if (timerfd >= 0) + { + clock_gettime(CLOCK_MONOTONIC, &start); + } + else + { + start = cur; + } timers.push_back({ .id = timer_id, .micros = micros, @@ -101,7 +115,7 @@ again: { nearest = -1; itimerspec exp = {}; - if (timerfd_settime(timerfd, 0, &exp, NULL)) + if (timerfd >= 0 && timerfd_settime(timerfd, 0, &exp, NULL)) { throw std::runtime_error(std::string("timerfd_settime: ") + strerror(errno)); } @@ -120,7 +134,14 @@ again: } } timespec now; - clock_gettime(CLOCK_MONOTONIC, &now); + if (timerfd >= 0) + { + clock_gettime(CLOCK_MONOTONIC, &now); + } + else + { + now = cur; + } itimerspec exp = { .it_interval = { 0 }, .it_value = timers[nearest].next, @@ -142,7 +163,7 @@ again: } exp.it_value = { .tv_sec = 0, .tv_nsec = 1 }; } - if (timerfd_settime(timerfd, 0, &exp, NULL)) + if (timerfd >= 0 && timerfd_settime(timerfd, 0, &exp, NULL)) { throw std::runtime_error(std::string("timerfd_settime: ") + strerror(errno)); } @@ -178,3 +199,13 @@ void timerfd_manager_t::trigger_nearest() nearest = -1; cb(nearest_id); } + +void timerfd_manager_t::tick(timespec passed) +{ + assert(timerfd == -1); + cur.tv_sec += passed.tv_sec; + cur.tv_nsec += passed.tv_nsec; + cur.tv_sec += (cur.tv_nsec / 1000000000); + cur.tv_nsec = (cur.tv_nsec % 1000000000); + set_nearest(true); +} diff --git a/src/util/timerfd_manager.h b/src/util/timerfd_manager.h index 58c58465..e127bd08 100644 --- a/src/util/timerfd_manager.h +++ b/src/util/timerfd_manager.h @@ -19,11 +19,12 @@ struct timerfd_timer_t class __attribute__((visibility("default"))) timerfd_manager_t { int wait_state = 0; - int timerfd; + int timerfd = -1; int nearest = -1; int id = 1; int onstack = 0; std::vector timers; + timespec cur = {}; void inc_timer(timerfd_timer_t & t); void set_nearest(bool trigger_inline); @@ -37,4 +38,5 @@ public: int set_timer(uint64_t millis, bool repeat, std::function callback); int set_timer_us(uint64_t micros, bool repeat, std::function callback); void clear_timer(int timer_id); + void tick(timespec passed); };