From 696a4eeb74481c596213c91e10e8bfe802542a38 Mon Sep 17 00:00:00 2001 From: allen <1007729991@qq.com> Date: Mon, 1 Sep 2025 15:47:30 +0800 Subject: [PATCH] add version for op Signed-off-by: allen <1007729991@qq.com> --- src/deva/bridge.h | 17 +++++++++++------ src/deva/container_op.cc | 13 +++++++------ src/deva/container_op.h | 17 ++++++++++++----- src/deva/deva.cc | 21 +++++++++++---------- src/deva/deva.h | 8 +++++--- src/deva/deva_service_impl.cc | 6 +++--- src/deva/op.cc | 8 ++++---- src/deva/op.h | 4 ++-- src/deva/rsm.cc | 4 ++-- 9 files changed, 57 insertions(+), 41 deletions(-) diff --git a/src/deva/bridge.h b/src/deva/bridge.h index 856286f..20f2fb4 100644 --- a/src/deva/bridge.h +++ b/src/deva/bridge.h @@ -11,10 +11,14 @@ namespace pain::deva { template -void bridge(RsmPtr rsm, const Request& request, Response* response, std::move_only_function cb) { +void bridge(int32_t version, + RsmPtr rsm, + const Request& request, + Response* response, + std::move_only_function cb) { // TODO: get rsm by partition id auto op = new ContainerOp( - 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(); @@ -22,12 +26,13 @@ void bridge(RsmPtr rsm, const Request& request, Response* response, std::move_on // Future style template -Future bridge(RsmPtr rsm, const Request& request, Response* response) { +Future bridge(int32_t version, RsmPtr rsm, const Request& request, Response* response) { Promise promise; auto future = promise.get_future(); - bridge(rsm, request, response, [promise = std::move(promise)](Status status) mutable { - promise.set_value(std::move(status)); - }); + bridge( + version, rsm, request, response, [promise = std::move(promise)](Status status) mutable { + promise.set_value(std::move(status)); + }); return future; } diff --git a/src/deva/container_op.cc b/src/deva/container_op.cc index bb0c09c..f3f9136 100644 --- a/src/deva/container_op.cc +++ b/src/deva/container_op.cc @@ -8,23 +8,24 @@ namespace pain::deva { template -OpPtr create(RsmPtr rsm) { +OpPtr create(int32_t version, RsmPtr rsm) { Request request; OpPtr op = nullptr; if constexpr (OpType < OpType::kMaxDevaOp) { - op = new ContainerOp(OpType, rsm, request, nullptr); + op = new ContainerOp(version, OpType, rsm, request, nullptr); } return op; } #define BRANCH(name) \ case OpType::k##name: \ - return create(rsm); \ + return create(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) @@ -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; } diff --git a/src/deva/container_op.h b/src/deva/container_op.h index 56d7db7..86ab517 100644 --- a/src/deva/container_op.h +++ b/src/deva/container_op.h @@ -34,13 +34,19 @@ class OpClosure : public braft::Closure { std::shared_ptr _span; }; -OpPtr decode(OpType op_type, IOBuf* buf, RsmPtr rsm); +OpPtr decode(int32_t version, OpType op_type, IOBuf* buf, RsmPtr rsm); template class ContainerOp : public Op { public: using OnFinish = std::move_only_function; - 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)), @@ -57,11 +63,12 @@ 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; @@ -69,14 +76,13 @@ class ContainerOp : public Op { } 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(container.get()); - auto status = c->process(&_request, _response, index); + auto status = c->process(_version, &_request, _response, index); on_finish(std::move(status)); } @@ -102,6 +108,7 @@ class ContainerOp : public Op { } protected: + int32_t _version; OpType _type; RsmPtr _rsm; Request _request; diff --git a/src/deva/deva.cc b/src/deva/deva.cc index 4cf0c2d..91eeba5 100644 --- a/src/deva/deva.cc +++ b/src/deva/deva.cc @@ -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) @@ -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()); @@ -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()); @@ -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; @@ -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(); } diff --git a/src/deva/deva.h b/src/deva/deva.h index 5771e0d..a3f9116 100644 --- a/src/deva/deva.h +++ b/src/deva/deva.h @@ -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 { diff --git a/src/deva/deva_service_impl.cc b/src/deva/deva_service_impl.cc index e0e551a..49aac81 100644 --- a/src/deva/deva_service_impl.cc +++ b/src/deva/deva_service_impl.cc @@ -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(_rsm, create_request, &create_response).get(); + auto status = bridge(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()); @@ -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(_rsm, create_request, &create_response).get(); + auto status = bridge(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()); @@ -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(_rsm, read_dir_request, &read_dir_response).get(); + auto status = bridge(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()); diff --git a/src/deva/op.cc b/src/deva/op.cc index 71d6b57..373a8e4 100644 --- a/src/deva/op.cc +++ b/src/deva/op.cc @@ -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; @@ -18,7 +18,7 @@ void encode(OpPtr op, IOBuf* buf) { buf->append(meta); } -OpPtr decode(IOBuf* buf, std::move_only_function decode) { +OpPtr decode(IOBuf* buf, std::move_only_function decode) { OpMeta op_meta = {}; static_assert(sizeof(op_meta) == 64, "OpMeta size must be 64byte"); // NOLINT(readability-magic-numbers) uint32_t meta_size = 0; @@ -26,7 +26,7 @@ OpPtr decode(IOBuf* buf, std::move_only_function decode) 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; } diff --git a/src/deva/op.h b/src/deva/op.h index 58a0e43..05f1c6f 100644 --- a/src/deva/op.h +++ b/src/deva/op.h @@ -67,8 +67,8 @@ class Op { class Rsm; using RsmPtr = boost::intrusive_ptr; -void encode(OpPtr op, IOBuf* buf); -OpPtr decode(IOBuf* buf, std::move_only_function decode); +void encode(int32_t version, OpPtr op, IOBuf* buf); +OpPtr decode(IOBuf* buf, std::move_only_function decode); } // namespace pain::deva diff --git a/src/deva/rsm.cc b/src/deva/rsm.cc index 3da65a0..ba7428c 100644 --- a/src/deva/rsm.cc +++ b/src/deva/rsm.cc @@ -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());