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
3 changes: 2 additions & 1 deletion src/common/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ load('//:copts.bzl', 'PAIN_COPTS', 'PAIN_LINKOPTS', 'PAIN_TEST_COPTS')

cc_library(
name = "pain_common",
srcs = glob(["*.h", "*.cc"]),
srcs = glob(["*.h", "*.cc", "rsm/*.h", "rsm/*.cc"]),
copts = PAIN_COPTS,
linkopts = PAIN_LINKOPTS,
deps = [
Expand All @@ -12,6 +12,7 @@ cc_library(
"@boost.smart_ptr",
"@rocksdb",
"@braft",
"@magic_enum",
],
visibility = ["//visibility:public"],
)
Expand Down
39 changes: 39 additions & 0 deletions src/common/rsm/bridge.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
#pragma once

#include <braft/raft.h>
#include <pain/base/future.h>
#include <pain/base/types.h>
#include <functional>
#include "common/rsm/container_op.h"
#include "common/rsm/op.h"
#include "common/rsm/rsm.h"

namespace pain::common {

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

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

} // namespace pain::common
6 changes: 4 additions & 2 deletions src/deva/container.h → src/common/rsm/container.h
Original file line number Diff line number Diff line change
Expand Up @@ -5,14 +5,16 @@
#include <string>
#include <string_view>
#include <boost/intrusive_ptr.hpp>
#include "common/rsm/op_factory.h"

namespace pain::deva {
namespace pain::common {

class Container {
public:
virtual ~Container() = default;
virtual Status save_snapshot(std::string_view path, std::vector<std::string>* files) = 0;
virtual Status load_snapshot(std::string_view path) = 0;
virtual OpFactory* op_factory() = 0;

private:
std::atomic<int> _use_count = 0;
Expand All @@ -30,4 +32,4 @@ class Container {

using ContainerPtr = boost::intrusive_ptr<Container>;

} // namespace pain::deva
} // namespace pain::common
13 changes: 13 additions & 0 deletions src/common/rsm/container_op.cc
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
#include "common/rsm/container_op.h"
#include "common/rsm/op_factory.h"

namespace pain::common {

OpPtr decode(int32_t version, uint32_t op_type, IOBuf* buf, RsmPtr rsm) {
auto op_factory = rsm->container()->op_factory();
auto op = op_factory->create(op_type, version, rsm);
op->decode(buf);
return op;
}

} // namespace pain::common
22 changes: 11 additions & 11 deletions src/deva/container_op.h → src/common/rsm/container_op.h
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
#pragma once

#include <braft/raft.h>
#include <pain/base/macro.h>
#include <pain/base/plog.h>
#include <pain/base/types.h>
#include <functional>
#include "deva/macro.h"
#include "deva/op.h"
#include "deva/rsm.h"
#include "common/rsm/op.h"
#include "common/rsm/rsm.h"

namespace pain::deva {
namespace pain::common {

class OpClosure : public braft::Closure {
public:
Expand All @@ -34,14 +34,14 @@ class OpClosure : public braft::Closure {
std::shared_ptr<opentelemetry::trace::Span> _span;
};

OpPtr decode(int32_t version, OpType op_type, IOBuf* buf, RsmPtr rsm);
OpPtr decode(int32_t version, uint32_t 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(int32_t version,
OpType type,
uint32_t type,
RsmPtr rsm,
Request request,
Response* response = nullptr,
Expand All @@ -57,18 +57,18 @@ class ContainerOp : public Op {
}
}

OpType type() const override {
uint32_t type() const override {
return _type;
}

void apply() override {
SPAN(span);
SPAN("common", 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(_version, self, &buf);
pain::common::encode(_version, self, &buf);
task.data = &buf;
task.done = new OpClosure(self, span);
task.expected_term = -1;
Expand Down Expand Up @@ -109,12 +109,12 @@ class ContainerOp : public Op {

protected:
int32_t _version;
OpType _type;
uint32_t _type;
RsmPtr _rsm;
Request _request;
Response* _response;
Response _internal_response;
OnFinish _finish;
};

} // namespace pain::deva
} // namespace pain::common
8 changes: 4 additions & 4 deletions src/deva/op.cc → src/common/rsm/op.cc
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
#include "deva/op.h"
#include "common/rsm/op.h"
#include <pain/base/plog.h>
#include <functional>
#include "butil/iobuf.h"
#include "butil/time.h"

namespace pain::deva {
namespace pain::common {

void encode(int32_t version, OpPtr op, IOBuf* buf) {
OpMeta op_meta = {};
Expand All @@ -19,7 +19,7 @@ void encode(int32_t version, OpPtr op, IOBuf* buf) {
buf->append(meta);
}

OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(int32_t, OpType, IOBuf*)> decode) {
OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(int32_t, uint32_t, IOBuf*)> decode) {
OpMeta op_meta = {};
static_assert(sizeof(op_meta) == 64, "OpMeta size must be 64byte"); // NOLINT(readability-magic-numbers)
uint32_t op_size = 0;
Expand All @@ -35,4 +35,4 @@ OpPtr decode(IOBuf* buf, std::move_only_function<OpPtr(int32_t, OpType, IOBuf*)>
return op;
}

} // namespace pain::deva
} // namespace pain::common
76 changes: 76 additions & 0 deletions src/common/rsm/op.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
#pragma once
#include <pain/base/types.h>
#include <atomic>
#include <cstdint>
#include <fmt/format.h>
#include <boost/intrusive_ptr.hpp>
#include <magic_enum/magic_enum.hpp>

#define DEFINE_RSM_OP(code, name, need_apply) k##name = ((code << 1) | (need_apply ? 1 : 0))

namespace pain::common {

// enum class OpType : uint32_t {
// kInvalid = 0,
// // DevaOp: 1 ~ 100
// DEFINE_DEVA_OP(1, CreateFile, true),
// DEFINE_DEVA_OP(2, CreateDir, true),
// DEFINE_DEVA_OP(3, RemoveFile, true),
// DEFINE_DEVA_OP(4, SealFile, true),
// DEFINE_DEVA_OP(5, CreateChunk, true),
// DEFINE_DEVA_OP(6, CheckInChunk, true),
// DEFINE_DEVA_OP(7, SealChunk, true),
// DEFINE_DEVA_OP(8, SealAndNewChunk, true),
// DEFINE_DEVA_OP(9, ReadDir, false),
// DEFINE_DEVA_OP(10, GetFileInfo, false),
// DEFINE_DEVA_OP(20, ManusyaHeartbeat, false),
// DEFINE_DEVA_OP(21, ListManusya, false),
// DEFINE_DEVA_OP(100, MaxDevaOp, true),
// };

inline constexpr bool need_apply(uint32_t type) {
return (type & 1) == 1;
}

struct OpMeta {
int32_t version; // op version
uint32_t type;
uint64_t timestamp;
uint32_t size;
char reserved[40]; // NOLINT(readability-magic-numbers)
};

static_assert(sizeof(OpMeta) == 64, "OpMeta size must be 64byte"); // NOLINT(readability-magic-numbers)

class Op;
using OpPtr = boost::intrusive_ptr<Op>;
class Op {
public:
Op() = default;
virtual ~Op() = default;
virtual uint32_t type() const = 0;
virtual void apply() = 0;
virtual void on_apply(int64_t index) = 0;
virtual void on_finish(Status status) = 0;
virtual void encode(IOBuf* buf) = 0;
virtual void decode(IOBuf* buf) = 0;

private:
std::atomic<int> _use_count = 0;
friend void intrusive_ptr_add_ref(Op* op) {
++op->_use_count;
}
friend void intrusive_ptr_release(Op* op) {
if (op->_use_count.fetch_sub(1) == 1) {
delete op;
}
}
};

class Rsm;
using RsmPtr = boost::intrusive_ptr<Rsm>;

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

} // namespace pain::common
14 changes: 14 additions & 0 deletions src/common/rsm/op_factory.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
#pragma once

#include "common/rsm/op.h"

namespace pain::common {

class OpFactory {
public:
OpFactory() = default;
virtual ~OpFactory() = default;
virtual OpPtr create(uint32_t op_type, int32_t version, RsmPtr rsm) = 0;
};

} // namespace pain::common
Loading
Loading