Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 11 additions & 6 deletions src/deva/bridge.h
Original file line number Diff line number Diff line change
Expand Up @@ -11,23 +11,28 @@
namespace pain::deva {

template <typename ContainerType, OpType OpType, typename Request, typename Response>
void bridge(RsmPtr rsm, const Request& request, Response* response, std::move_only_function<void(Status)> cb) {
void bridge(int32_t version,
RsmPtr rsm,
const Request& request,
Response* response,
std::move_only_function<void(Status)> cb) {
// TODO: get rsm by partition id
auto op = new ContainerOp<ContainerType, Request, Response>(
OpType, rsm, request, response, [cb = std::move(cb)](Status status) mutable {
version, OpType, rsm, request, response, [cb = std::move(cb)](Status status) mutable {
cb(std::move(status));
});
op->apply();
}

// Future style
template <typename ContainerType, OpType OpType, typename Request, typename Response>
Future<Status> bridge(RsmPtr rsm, const Request& request, Response* response) {
Future<Status> bridge(int32_t version, RsmPtr rsm, const Request& request, Response* response) {
Promise<Status> promise;
auto future = promise.get_future();
bridge<ContainerType, OpType>(rsm, request, response, [promise = std::move(promise)](Status status) mutable {
promise.set_value(std::move(status));
});
bridge<ContainerType, OpType>(
version, rsm, request, response, [promise = std::move(promise)](Status status) mutable {
promise.set_value(std::move(status));
});
return future;
}

Expand Down
13 changes: 7 additions & 6 deletions src/deva/container_op.cc
Original file line number Diff line number Diff line change
Expand Up @@ -8,23 +8,24 @@
namespace pain::deva {

template <OpType OpType, typename Request, typename Response>
OpPtr create(RsmPtr rsm) {
OpPtr create(int32_t version, RsmPtr rsm) {
Request request;
OpPtr op = nullptr;

if constexpr (OpType < OpType::kMaxDevaOp) {
op = new ContainerOp<Deva, Request, Response>(OpType, rsm, request, nullptr);
op = new ContainerOp<Deva, Request, Response>(version, OpType, rsm, request, nullptr);
}
return op;
}

#define BRANCH(name) \
case OpType::k##name: \
return create<OpType::k##name, proto::deva::store::name##Request, proto::deva::store::name##Response>(rsm); \
return create<OpType::k##name, proto::deva::store::name##Request, proto::deva::store::name##Response>(version, \
rsm); \
break;

OpPtr decode(OpType op_type, IOBuf* buf, RsmPtr rsm) {
auto op = [](OpType op_type, IOBuf* buf, RsmPtr rsm) -> OpPtr {
OpPtr decode(int32_t version, OpType op_type, IOBuf* buf, RsmPtr rsm) {
auto op = [](int32_t version, OpType op_type, RsmPtr rsm) -> OpPtr {
switch (op_type) {
BRANCH(CreateFile)
BRANCH(CreateDir)
Expand All @@ -39,7 +40,7 @@ OpPtr decode(OpType op_type, IOBuf* buf, RsmPtr rsm) {
BOOST_ASSERT_MSG(false, fmt::format("unknown op type: {}", op_type).c_str());
}
return nullptr;
}(op_type, buf, rsm);
}(version, op_type, rsm);
op->decode(buf);
return op;
}
Expand Down
17 changes: 12 additions & 5 deletions src/deva/container_op.h
Original file line number Diff line number Diff line change
Expand Up @@ -34,13 +34,19 @@ class OpClosure : public braft::Closure {
std::shared_ptr<opentelemetry::trace::Span> _span;
};

OpPtr decode(OpType op_type, IOBuf* buf, RsmPtr rsm);
OpPtr decode(int32_t version, OpType op_type, IOBuf* buf, RsmPtr rsm);

template <typename ContainerType, typename Request, typename Response>
class ContainerOp : public Op {
public:
using OnFinish = std::move_only_function<void(Status)>;
ContainerOp(OpType type, RsmPtr rsm, Request request, Response* response = nullptr, OnFinish finish = nullptr) :
ContainerOp(int32_t version,
OpType type,
RsmPtr rsm,
Request request,
Response* response = nullptr,
OnFinish finish = nullptr) :
_version(version),
_type(type),
_rsm(rsm),
_request(std::move(request)),
Expand All @@ -57,26 +63,26 @@ class ContainerOp : public Op {

void apply() override {
SPAN(span);
PLOG_DEBUG(("desc", "apply op")("type", _type)("version", _version));
if (need_apply(_type)) {
braft::Task task;
IOBuf buf;
OpPtr self(this);
pain::deva::encode(self, &buf);
pain::deva::encode(_version, self, &buf);
task.data = &buf;
task.done = new OpClosure(self, span);
task.expected_term = -1;
_rsm->apply(task);
} else {
on_apply(0);
}
PLOG_DEBUG(("desc", "apply op")("type", _type));
}

void on_apply(int64_t index) override {
PLOG_DEBUG(("desc", "on apply op")("type", _type)("index", index));
auto container = _rsm->container();
auto c = static_cast<ContainerType*>(container.get());
auto status = c->process(&_request, _response, index);
auto status = c->process(_version, &_request, _response, index);
on_finish(std::move(status));
}

Expand All @@ -102,6 +108,7 @@ class ContainerOp : public Op {
}

protected:
int32_t _version;
OpType _type;
RsmPtr _rsm;
Request _request;
Expand Down
21 changes: 11 additions & 10 deletions src/deva/deva.cc
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,8 @@
#include "deva/macro.h"

#define DEVA_METHOD(name) \
Status Deva::name([[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
Status Deva::name([[maybe_unused]] int32_t version, \
[[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
[[maybe_unused]] pain::proto::deva::store::name##Response* response, \
[[maybe_unused]] int64_t index)

Expand Down Expand Up @@ -39,7 +40,7 @@ Status Deva::create(const std::string& path, const UUID& id, FileType type) {

DEVA_METHOD(CreateFile) {
SPAN(span);
PLOG_DEBUG(("desc", "create_file")("index", index)("request", request->DebugString()));
PLOG_DEBUG(("desc", "create_file")("version", version)("index", index)("request", request->DebugString()));
auto& path = request->path();
auto& file_id = request->file_id();
UUID file_uuid(file_id.high(), file_id.low());
Expand All @@ -64,7 +65,7 @@ DEVA_METHOD(CreateFile) {

DEVA_METHOD(CreateDir) {
SPAN(span);
PLOG_DEBUG(("desc", "create_dir")("index", index)("request", request->DebugString()));
PLOG_DEBUG(("desc", "create_dir")("version", version)("index", index)("request", request->DebugString()));
auto& path = request->path();
auto& dir_id = request->dir_id();
UUID dir_uuid(dir_id.high(), dir_id.low());
Expand All @@ -89,7 +90,7 @@ DEVA_METHOD(CreateDir) {

DEVA_METHOD(ReadDir) {
SPAN(span);
PLOG_DEBUG(("desc", "read_dir")("index", index)("request", request->DebugString()));
PLOG_DEBUG(("desc", "read_dir")("version", version)("index", index)("request", request->DebugString()));
auto& path = request->path();
UUID parent_dir_uuid;
FileType file_type = FileType::kNone;
Expand All @@ -114,37 +115,37 @@ DEVA_METHOD(ReadDir) {

DEVA_METHOD(RemoveFile) {
SPAN(span);
PLOG_INFO(("desc", "remove_file")("index", index));
PLOG_INFO(("desc", "remove_file")("version", version)("index", index));
return Status::OK();
}

DEVA_METHOD(SealFile) {
SPAN(span);
PLOG_INFO(("desc", "seal_file")("index", index));
PLOG_INFO(("desc", "seal_file")("version", version)("index", index));
return Status::OK();
}

DEVA_METHOD(CreateChunk) {
SPAN(span);
PLOG_INFO(("desc", "remove")("index", index));
PLOG_INFO(("desc", "remove")("version", version)("index", index));
return Status::OK();
}

DEVA_METHOD(CheckInChunk) {
SPAN(span);
PLOG_INFO(("desc", "seal")("index", index));
PLOG_INFO(("desc", "seal")("version", version)("index", index));
return Status::OK();
}

DEVA_METHOD(SealChunk) {
SPAN(span);
PLOG_INFO(("desc", "create_chunk")("index", index));
PLOG_INFO(("desc", "create_chunk")("version", version)("index", index));
return Status::OK();
}

DEVA_METHOD(SealAndNewChunk) {
SPAN(span);
PLOG_INFO(("desc", "remove_chunk")("index", index));
PLOG_INFO(("desc", "remove_chunk")("version", version)("index", index));
return Status::OK();
}

Expand Down
8 changes: 5 additions & 3 deletions src/deva/deva.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,15 @@
#include "deva/container.h"
#include "deva/namespace.h"
#define DEVA_ENTRY(name) \
Status name([[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
Status name([[maybe_unused]] int32_t version, \
[[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
[[maybe_unused]] pain::proto::deva::store::name##Response* response, \
[[maybe_unused]] int64_t index); \
Status process([[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
Status process([[maybe_unused]] int32_t version, \
[[maybe_unused]] const pain::proto::deva::store::name##Request* request, \
[[maybe_unused]] pain::proto::deva::store::name##Response* response, \
[[maybe_unused]] int64_t index) { \
return name(request, response, index); \
return name(version, request, response, index); \
}

namespace pain::deva {
Expand Down
6 changes: 3 additions & 3 deletions src/deva/deva_service_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ DEVA_SERVICE_METHOD(OpenFile) {
create_request.set_atime(butil::gettimeofday_us());
create_request.set_mtime(butil::gettimeofday_us());
create_request.set_ctime(butil::gettimeofday_us());
auto status = bridge<Deva, OpType::kCreateFile>(_rsm, create_request, &create_response).get();
auto status = bridge<Deva, OpType::kCreateFile>(1, _rsm, create_request, &create_response).get();
if (!status.ok()) {
PLOG_ERROR(("desc", "failed to create file")("error", status.error_str()));
response->mutable_header()->set_status(status.error_code());
Expand Down Expand Up @@ -81,7 +81,7 @@ DEVA_SERVICE_METHOD(Mkdir) {
create_request.set_atime(butil::gettimeofday_us());
create_request.set_mtime(butil::gettimeofday_us());
create_request.set_ctime(butil::gettimeofday_us());
auto status = bridge<Deva, OpType::kCreateDir>(_rsm, create_request, &create_response).get();
auto status = bridge<Deva, OpType::kCreateDir>(1, _rsm, create_request, &create_response).get();
if (!status.ok()) {
PLOG_ERROR(("desc", "failed to create file")("error", status.error_str()));
response->mutable_header()->set_status(status.error_code());
Expand All @@ -100,7 +100,7 @@ DEVA_SERVICE_METHOD(ReadDir) {
pain::proto::deva::store::ReadDirRequest read_dir_request;
pain::proto::deva::store::ReadDirResponse read_dir_response;
read_dir_request.set_path(path);
auto status = bridge<Deva, OpType::kReadDir>(_rsm, read_dir_request, &read_dir_response).get();
auto status = bridge<Deva, OpType::kReadDir>(1, _rsm, read_dir_request, &read_dir_response).get();
if (!status.ok()) {
PLOG_ERROR(("desc", "failed to read dir")("error", status.error_str()));
response->mutable_header()->set_status(status.error_code());
Expand Down
8 changes: 4 additions & 4 deletions src/deva/op.cc
Original file line number Diff line number Diff line change
Expand Up @@ -5,10 +5,10 @@

namespace pain::deva {

void encode(OpPtr op, IOBuf* buf) {
void encode(int32_t version, OpPtr op, IOBuf* buf) {
OpMeta op_meta = {};
static_assert(sizeof(op_meta) == 64, "OpMeta size must be 64byte"); // NOLINT(readability-magic-numbers)
op_meta.version = 1;
op_meta.version = version;
op_meta.type = op->type();
op_meta.timestamp = butil::gettimeofday_us();
butil::IOBuf meta;
Expand All @@ -18,15 +18,15 @@ void encode(OpPtr op, IOBuf* buf) {
buf->append(meta);
}

OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(OpType, IOBuf*)> decode) {
OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(int32_t, OpType, IOBuf*)> decode) {
OpMeta op_meta = {};
static_assert(sizeof(op_meta) == 64, "OpMeta size must be 64byte"); // NOLINT(readability-magic-numbers)
uint32_t meta_size = 0;
buf->cutn(&op_meta, sizeof(op_meta));
meta_size = op_meta.size;
butil::IOBuf meta;
buf->cutn(&meta, meta_size);
auto op = decode(op_meta.type, &meta);
auto op = decode(op_meta.version, op_meta.type, buf);
return op;
}

Expand Down
4 changes: 2 additions & 2 deletions src/deva/op.h
Original file line number Diff line number Diff line change
Expand Up @@ -67,8 +67,8 @@ class Op {
class Rsm;
using RsmPtr = boost::intrusive_ptr<Rsm>;

void encode(OpPtr op, IOBuf* buf);
OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(OpType, IOBuf*)> decode);
void encode(int32_t version, OpPtr op, IOBuf* buf);
OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(int32_t, OpType, IOBuf*)> decode);

} // namespace pain::deva

Expand Down
4 changes: 2 additions & 2 deletions src/deva/rsm.cc
Original file line number Diff line number Diff line change
Expand Up @@ -90,8 +90,8 @@ void Rsm::on_apply(braft::Iterator& iter) {
} else {
butil::IOBuf saved_log = iter.data();
// clang-format off
auto op = decode(&saved_log, [rsm = RsmPtr(this)](OpType op_type, IOBuf* buf) {
return decode(op_type, buf, rsm);
auto op = decode(&saved_log, [rsm = RsmPtr(this)](int32_t version, OpType op_type, IOBuf* buf) {
return decode(version, op_type, buf, rsm);
});
// clang-format on
op->on_apply(iter.index());
Expand Down
Loading