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: 15 additions & 2 deletions src/common/rocksdb_store.cc
Original file line number Diff line number Diff line change
Expand Up @@ -119,14 +119,23 @@ Status RocksdbStore::check_point(const char* to, std::vector<std::string>* files
}
std::unique_ptr<rocksdb::Checkpoint> cpt_guard(cpt);

// remove to path
auto fs = braft::default_file_system();
if (fs->directory_exists(to)) {
if (!fs->delete_file(to, true)) {
PLOG_ERROR(("desc", "delete cpt dir failed") //
("path", to));
}
}

status = cpt->CreateCheckpoint(to);
if (!status.ok()) {
PLOG_ERROR(("desc", "create checkpoint failed") //
("path", to) //
("error", status.ToString()));
return convert_to_pain_status(status);
}

auto fs = braft::default_file_system();
std::unique_ptr<braft::DirReader> dir_reader(fs->directory_reader(to));

if (dir_reader == nullptr) {
Expand All @@ -143,7 +152,7 @@ Status RocksdbStore::check_point(const char* to, std::vector<std::string>* files

std::vector<std::string> snapshot_files;
while (dir_reader->next()) {
auto file_name = fmt::format("cpt/{}", dir_reader->name());
auto file_name = fmt::format("{}", dir_reader->name());
PLOG_INFO(("desc", "snapshot add file") //
("file", file_name));
snapshot_files.push_back(file_name);
Expand Down Expand Up @@ -225,6 +234,10 @@ Status RocksdbStore::recover(const char* from) {
("path", bak_path));
}

// delete old db
delete _txn_db;
_txn_db = nullptr;
_db = nullptr;
open_or_die();

return Status::OK();
Expand Down
32 changes: 16 additions & 16 deletions src/common/rocksdb_util.h
Original file line number Diff line number Diff line change
Expand Up @@ -11,37 +11,37 @@ inline Status convert_to_pain_status(const rocksdb::Status& status) {
case rocksdb::Status::Code::kOk:
return Status::OK();
case rocksdb::Status::Code::kNotFound:
return Status(ENOENT, fmt::format("NotFound:{}", status.ToString()));
return Status(ENOENT, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kCorruption:
return Status(EBADMSG, fmt::format("Corruption:{}", status.ToString()));
return Status(EBADMSG, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kNotSupported:
return Status(ENOTSUP, fmt::format("NotSupported:{}", status.ToString()));
return Status(ENOTSUP, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kInvalidArgument:
return Status(EINVAL, fmt::format("InvalidArgument:{}", status.ToString()));
return Status(EINVAL, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kIOError:
return Status(EIO, fmt::format("IOError:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kMergeInProgress:
return Status(EBUSY, fmt::format("MergeInProgress:{}", status.ToString()));
return Status(EBUSY, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kIncomplete:
return Status(EIO, fmt::format("Incomplete:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kShutdownInProgress:
return Status(EIO, fmt::format("ShutdownInProgress:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kTimedOut:
return Status(ETIMEDOUT, fmt::format("TimedOut:{}", status.ToString()));
return Status(ETIMEDOUT, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kAborted:
return Status(EIO, fmt::format("Aborted:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kBusy:
return Status(EBUSY, fmt::format("Busy:{}", status.ToString()));
return Status(EBUSY, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kExpired:
return Status(EKEYEXPIRED, fmt::format("Expired:{}", status.ToString()));
return Status(EKEYEXPIRED, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kTryAgain:
return Status(EAGAIN, fmt::format("TryAgain:{}", status.ToString()));
return Status(EAGAIN, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kCompactionTooLarge:
return Status(EIO, fmt::format("CompactionTooLarge:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
case rocksdb::Status::Code::kColumnFamilyDropped:
return Status(EIO, fmt::format("ColumnFamilyDropped:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
default:
return Status(EIO, fmt::format("Unknown:{}", status.ToString()));
return Status(EIO, fmt::format("{}", status.ToString()));
}
}

Expand Down
7 changes: 2 additions & 5 deletions src/deva/deva.cc
Original file line number Diff line number Diff line change
Expand Up @@ -150,14 +150,11 @@ DEVA_METHOD(SealAndNewChunk) {
}

Status Deva::save_snapshot(std::string_view path, std::vector<std::string>* files) {
std::ignore = path;
std::ignore = files;
return Status::OK();
return _store->check_point(path.data(), files);
}

Status Deva::load_snapshot(std::string_view path) {
std::ignore = path;
return Status::OK();
return _store->recover(path.data());
}

} // namespace pain::deva
4 changes: 3 additions & 1 deletion src/deva/mock/deva_machine.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ class DevaMachine {
BOOST_ASSERT_MSG(false, "Fail to parse address");
}

node_options.election_timeout_ms = 5000; // NOLINT
node_options.election_timeout_ms = 1000; // NOLINT
node_options.node_owns_fsm = false;
node_options.snapshot_interval_s = 30; // NOLINT
std::string prefix = fmt::format("local://{}", data_path);
Expand All @@ -31,6 +31,8 @@ class DevaMachine {
node_options.snapshot_uri = fmt::format("{}/{}/snapshot", prefix, group);

std::string rocksdb_path = fmt::format("{}/{}/db", data_path, group);
// remove rocksdb path
std::filesystem::remove_all(rocksdb_path);
common::RocksdbStorePtr store;
auto status = common::RocksdbStore::open(rocksdb_path.c_str(), &store);
BOOST_ASSERT_MSG(status.ok(), "Fail to open rocksdb store");
Expand Down
29 changes: 27 additions & 2 deletions src/deva/mock/mock_deva.cc
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
#include "deva/mock/mock_deva.h"
#include <braft/cli.h>
#include <braft/route_table.h>
#include <pain/base/path.h>
#include <pain/base/plog.h>
Expand All @@ -11,7 +12,8 @@ namespace pain::deva::mock {
MockDeva::MockDeva() :
_group("test_group"),
_data_path("/tmp/deva_mock_data_XXXXXX"),
_node_addrs({"127.0.0.1:8200", "127.0.0.1:8201", "127.0.0.1:8202"}) {
_node_addrs({"127.0.0.1:8200", "127.0.0.1:8201", "127.0.0.1:8202"}),
_do_not_remove_data_path(false) {
_node_conf = fmt::format("{}", fmt::join(_node_addrs, ","));
make_temp_dir_or_die(&_data_path);
PLOG_INFO(("data_path", _data_path));
Expand All @@ -29,10 +31,13 @@ MockDeva::MockDeva() :

MockDeva::~MockDeva() {
stop();
std::filesystem::remove_all(_data_path);
if (!_do_not_remove_data_path) {
std::filesystem::remove_all(_data_path);
}
}

Status MockDeva::start() {
PLOG_INFO(("desc", "start")("node_count", _deva_machine.size()));
for (size_t i = 0; i < _deva_machine.size(); i++) {
auto status = start(i);
if (!status.ok()) {
Expand All @@ -43,25 +48,29 @@ Status MockDeva::start() {
}

Status MockDeva::start(int index) {
PLOG_INFO(("desc", "start")("index", index));
_deva_machine[index] = std::make_unique<DevaMachine>(
_data_paths[index].c_str(), _group.c_str(), _node_addrs[index].c_str(), _node_conf.c_str());
return _deva_machine[index]->start();
}

void MockDeva::stop(int index) {
if (_deva_machine[index]) {
PLOG_INFO(("desc", "stop")("index", index));
_deva_machine[index]->stop();
_deva_machine[index].reset();
}
}

void MockDeva::stop() {
PLOG_INFO(("desc", "stop"));
for (size_t i = 0; i < _deva_machine.size(); i++) {
stop(i);
}
}

Status MockDeva::wait_for_leader(std::string* addr, int timeout_ms) {
PLOG_INFO(("desc", "wait_for_leader")("timeout_ms", timeout_ms));
constexpr int sleep_interval_ms = 500;
for (int i = 0; i < timeout_ms / sleep_interval_ms; i++) {
auto status = braft::rtb::refresh_leader(_group, sleep_interval_ms);
Expand All @@ -81,4 +90,20 @@ Status MockDeva::wait_for_leader(std::string* addr, int timeout_ms) {
return Status(ENODEV, "No leader");
}

Status MockDeva::snapshot(int index) {
PLOG_INFO(("desc", "snapshot")("index", index));
return braft::cli::snapshot(_group, _node_addrs[index], braft::cli::CliOptions());
}

Status MockDeva::snapshot() {
PLOG_INFO(("desc", "snapshot"));
for (size_t i = 0; i < _deva_machine.size(); i++) {
auto status = snapshot(i);
if (!status.ok()) {
return status;
}
}
return Status::OK();
}

} // namespace pain::deva::mock
8 changes: 8 additions & 0 deletions src/deva/mock/mock_deva.h
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,21 @@ class MockDeva {
return _data_paths[index];
}

Status snapshot(int index);
Status snapshot();

void do_not_remove_data_path() {
_do_not_remove_data_path = true;
}

private:
std::string _group;
std::string _node_conf;
std::string _data_path;
std::vector<std::string> _data_paths;
std::vector<std::string> _node_addrs;
std::vector<std::unique_ptr<DevaMachine>> _deva_machine;
bool _do_not_remove_data_path = false;
};

} // namespace pain::deva::mock
2 changes: 1 addition & 1 deletion src/deva/namespace.cc
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,7 @@ void Namespace::list(const UUID& parent, std::list<DirEntry>* entries) const {
Status Namespace::parse_path(const char* path, std::list<std::string_view>* components) const {
const char* p = path;
if (*p != '/') {
return Status(EINVAL, "Invalid path");
return Status(EINVAL, fmt::format("Invalid path:{}", path));
}

while (*p != '\0') {
Expand Down
13 changes: 9 additions & 4 deletions src/deva/op.cc
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
#include "deva/op.h"
#include <pain/base/plog.h>
#include <functional>
#include "butil/iobuf.h"
#include "butil/time.h"
Expand All @@ -21,11 +22,15 @@ void encode(int32_t version, OpPtr op, IOBuf* buf) {
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;
uint32_t op_size = 0;
buf->cutn(&op_meta, sizeof(op_meta));
meta_size = op_meta.size;
butil::IOBuf meta;
buf->cutn(&meta, meta_size);
op_size = op_meta.size;
if (buf->size() < op_size) {
PLOG_ERROR(("desc", "op size is too small") //
("op_size", op_size) //
("buf_size", buf->size()));
return nullptr;
}
auto op = decode(op_meta.version, op_meta.type, buf);
return op;
}
Expand Down
31 changes: 29 additions & 2 deletions src/deva/rsm.cc
Original file line number Diff line number Diff line change
Expand Up @@ -39,10 +39,16 @@ Rsm::Rsm(const butil::EndPoint& address,
_node_options.fsm = this;
}
Rsm::~Rsm() {
PLOG_INFO(("desc", "destructor rsm") //
("group", _group) //
("address", butil::endpoint2str(_address).c_str()));
delete _node;
}

int Rsm::start() {
PLOG_INFO(("desc", "start rsm") //
("group", _group) //
("address", butil::endpoint2str(_address).c_str()));
braft::Node* node = new braft::Node(_group, braft::PeerId(_address));
if (node->init(_node_options) != 0) {
LOG(ERROR) << "Fail to init raft node";
Expand Down Expand Up @@ -108,25 +114,46 @@ struct SnapshotArg {
};

void* Rsm::save_snapshot(void* arg) {
PLOG_INFO(("desc", "save_snapshot"));
SnapshotArg* sa = (SnapshotArg*)arg;
std::unique_ptr<SnapshotArg> arg_guard(sa);
brpc::ClosureGuard done_guard(sa->done);
std::string snapshot_path = sa->writer->get_path() + "/data";
std::string snapshot_path = sa->writer->get_path();
std::vector<std::string> files;
auto status = sa->container->save_snapshot(snapshot_path, &files);
if (!status.ok()) {
sa->done->status() = status;
LOG(ERROR) << "Fail to save snapshot to " << snapshot_path;
return nullptr;
}
for (const auto& file : files) {
PLOG_INFO(("desc", "add_file to snapshot")("file", file));
sa->writer->add_file(file);
}
PLOG_INFO(("desc", "save_snapshot done"));
return nullptr;
}

void Rsm::on_snapshot_save(braft::SnapshotWriter* writer, braft::Closure* done) {
PLOG_INFO(("desc", "on_snapshot_save"));
SnapshotArg* arg = new SnapshotArg;
arg->writer = writer;
arg->done = done;
arg->container = _container;
bthread_t tid = 0;
bthread_start_urgent(&tid, nullptr, save_snapshot, arg);
}

// NOLINTNEXTLINE
int Rsm::on_snapshot_load(braft::SnapshotReader* reader) {
std::ignore = reader;
PLOG_INFO(("desc", "on_snapshot_load"));
CHECK(!is_leader()) << "Leader is not supposed to load snapshot";
auto path = reader->get_path();
auto status = _container->load_snapshot(path);
if (!status.ok()) {
PLOG_ERROR(("desc", "fail to load snapshot from ")("path", path)("error", status.error_str()));
return -1;
}
return 0;
}

Expand Down
1 change: 1 addition & 0 deletions src/deva/rsm.h
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ class Rsm : public braft::StateMachine {
struct SnapshotArg {
braft::SnapshotWriter* writer;
braft::Closure* done;
ContainerPtr container;
};

static void* save_snapshot(void* arg);
Expand Down
Loading
Loading