Implement buffered disk_mock_t mode, extract ringloop_mock.cpp
This commit is contained in:
@@ -8,7 +8,6 @@
|
||||
|
||||
#include <sys/eventfd.h>
|
||||
|
||||
#include "malloc_or_die.h"
|
||||
#include "ringloop.h"
|
||||
|
||||
ring_loop_t::ring_loop_t(int qd, bool multithreaded, bool sqe128)
|
||||
@@ -229,230 +228,3 @@ 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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user