Add mocks for blockstore integration tests: timerfd, ring_loop_mock_t and disk_mock_t
This commit is contained in:
@@ -8,6 +8,7 @@
|
||||
|
||||
#include <sys/eventfd.h>
|
||||
|
||||
#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<void(io_uring_sqe *)> 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<void()> & 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;
|
||||
}
|
||||
|
||||
+70
-2
@@ -32,7 +32,27 @@ struct ring_consumer_t
|
||||
std::function<void(void)> 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<void()> & 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<std::function<void()>> immediate_queue, immediate_queue2;
|
||||
std::vector<ring_consumer_t*> consumers;
|
||||
@@ -54,7 +74,7 @@ public:
|
||||
int register_eventfd();
|
||||
|
||||
io_uring_sqe* get_sqe();
|
||||
inline void set_immediate(const std::function<void()> cb)
|
||||
inline void set_immediate(const std::function<void()> & 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<std::function<void()>> immediate_queue, immediate_queue2;
|
||||
std::vector<ring_consumer_t*> consumers;
|
||||
std::vector<io_uring_sqe> sqes;
|
||||
std::vector<ring_data_t> ring_datas;
|
||||
std::vector<ring_data_t *> free_ring_datas;
|
||||
std::vector<ring_data_t *> submit_ring_datas;
|
||||
std::vector<ring_data_t *> completed_ring_datas;
|
||||
std::function<void(io_uring_sqe *)> submit_cb;
|
||||
bool in_loop;
|
||||
bool loop_again;
|
||||
bool support_zc = false;
|
||||
|
||||
public:
|
||||
ring_loop_mock_t(int qd, std::function<void(io_uring_sqe *)> submit_cb);
|
||||
|
||||
void register_consumer(ring_consumer_t *consumer);
|
||||
void unregister_consumer(ring_consumer_t *consumer);
|
||||
void wakeup();
|
||||
void set_immediate(const std::function<void()> & 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);
|
||||
};
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
#include <sys/timerfd.h>
|
||||
#include <sys/poll.h>
|
||||
#include <sys/epoll.h>
|
||||
#include <assert.h>
|
||||
#include <unistd.h>
|
||||
#include <errno.h>
|
||||
#include <string.h>
|
||||
@@ -15,21 +16,27 @@ timerfd_manager_t::timerfd_manager_t(std::function<void(int, bool, std::function
|
||||
{
|
||||
this->set_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);
|
||||
}
|
||||
|
||||
@@ -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<timerfd_timer_t> 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<void(int)> callback);
|
||||
int set_timer_us(uint64_t micros, bool repeat, std::function<void(int)> callback);
|
||||
void clear_timer(int timer_id);
|
||||
void tick(timespec passed);
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user