diff --git a/src/common/rocksdb_store.cc b/src/common/rocksdb_store.cc index b87316f..df0f650 100644 --- a/src/common/rocksdb_store.cc +++ b/src/common/rocksdb_store.cc @@ -119,14 +119,23 @@ Status RocksdbStore::check_point(const char* to, std::vector* files } std::unique_ptr 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 dir_reader(fs->directory_reader(to)); if (dir_reader == nullptr) { @@ -143,7 +152,7 @@ Status RocksdbStore::check_point(const char* to, std::vector* files std::vector 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); @@ -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(); diff --git a/src/common/rocksdb_util.h b/src/common/rocksdb_util.h index b5849e7..9a5e39a 100644 --- a/src/common/rocksdb_util.h +++ b/src/common/rocksdb_util.h @@ -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())); } } diff --git a/src/deva/deva.cc b/src/deva/deva.cc index 91eeba5..08a54a9 100644 --- a/src/deva/deva.cc +++ b/src/deva/deva.cc @@ -150,14 +150,11 @@ DEVA_METHOD(SealAndNewChunk) { } Status Deva::save_snapshot(std::string_view path, std::vector* 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 diff --git a/src/deva/mock/deva_machine.h b/src/deva/mock/deva_machine.h index cb704bf..387cec3 100644 --- a/src/deva/mock/deva_machine.h +++ b/src/deva/mock/deva_machine.h @@ -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); @@ -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"); diff --git a/src/deva/mock/mock_deva.cc b/src/deva/mock/mock_deva.cc index ebe6c20..5f3c35a 100644 --- a/src/deva/mock/mock_deva.cc +++ b/src/deva/mock/mock_deva.cc @@ -1,4 +1,5 @@ #include "deva/mock/mock_deva.h" +#include #include #include #include @@ -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)); @@ -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()) { @@ -43,6 +48,7 @@ Status MockDeva::start() { } Status MockDeva::start(int index) { + PLOG_INFO(("desc", "start")("index", index)); _deva_machine[index] = std::make_unique( _data_paths[index].c_str(), _group.c_str(), _node_addrs[index].c_str(), _node_conf.c_str()); return _deva_machine[index]->start(); @@ -50,18 +56,21 @@ Status MockDeva::start(int index) { 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); @@ -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 diff --git a/src/deva/mock/mock_deva.h b/src/deva/mock/mock_deva.h index 60ea852..bdef7be 100644 --- a/src/deva/mock/mock_deva.h +++ b/src/deva/mock/mock_deva.h @@ -29,6 +29,13 @@ 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; @@ -36,6 +43,7 @@ class MockDeva { std::vector _data_paths; std::vector _node_addrs; std::vector> _deva_machine; + bool _do_not_remove_data_path = false; }; } // namespace pain::deva::mock diff --git a/src/deva/namespace.cc b/src/deva/namespace.cc index 25515ba..ca65fef 100644 --- a/src/deva/namespace.cc +++ b/src/deva/namespace.cc @@ -156,7 +156,7 @@ void Namespace::list(const UUID& parent, std::list* entries) const { Status Namespace::parse_path(const char* path, std::list* components) const { const char* p = path; if (*p != '/') { - return Status(EINVAL, "Invalid path"); + return Status(EINVAL, fmt::format("Invalid path:{}", path)); } while (*p != '\0') { diff --git a/src/deva/op.cc b/src/deva/op.cc index 373a8e4..84b50e8 100644 --- a/src/deva/op.cc +++ b/src/deva/op.cc @@ -1,4 +1,5 @@ #include "deva/op.h" +#include #include #include "butil/iobuf.h" #include "butil/time.h" @@ -21,11 +22,15 @@ void encode(int32_t version, OpPtr op, IOBuf* buf) { 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; + 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; } diff --git a/src/deva/rsm.cc b/src/deva/rsm.cc index ba7428c..b5cd4a6 100644 --- a/src/deva/rsm.cc +++ b/src/deva/rsm.cc @@ -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"; @@ -108,25 +114,46 @@ struct SnapshotArg { }; void* Rsm::save_snapshot(void* arg) { + PLOG_INFO(("desc", "save_snapshot")); SnapshotArg* sa = (SnapshotArg*)arg; std::unique_ptr 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 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; } diff --git a/src/deva/rsm.h b/src/deva/rsm.h index 9a6d7ce..9879978 100644 --- a/src/deva/rsm.h +++ b/src/deva/rsm.h @@ -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); diff --git a/src/deva/test/test_deva.cc b/src/deva/test/test_deva.cc index 78e252b..a381565 100644 --- a/src/deva/test/test_deva.cc +++ b/src/deva/test/test_deva.cc @@ -32,6 +32,12 @@ class TestDeva : public ::testing::Test { _mock_deva.group().c_str(), &pain::proto::deva::DevaService::ReadDir, &request, response); } + void TearDown() override { + if (::testing::Test::HasFailure()) { + _mock_deva.do_not_remove_data_path(); + } + } + protected: pain::deva::mock::MockDeva _mock_deva; }; @@ -153,4 +159,47 @@ TEST_F(TestDeva, ReadDir) { } } +TEST_F(TestDeva, Snapshot) { + _mock_deva.start(); + SCOPE_EXIT { + _mock_deva.stop(); + }; + std::string leader; + auto status = _mock_deva.wait_for_leader(&leader); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "leader: " << leader << std::endl; + + pain::proto::deva::MkdirResponse mkdir_response; + status = mkdir("/test", &mkdir_response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "mkdir response: " << mkdir_response.DebugString() << std::endl; + + pain::proto::deva::OpenFileResponse response; + status = open("/test/test.txt", pain::proto::deva::OpenFlag::OPEN_CREATE, &response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "response: " << response.DebugString() << std::endl; + + status = _mock_deva.snapshot(); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "snapshot response: " << status.error_str() << "(" << status.error_code() << ")" << std::endl; + + // restart the deva + _mock_deva.stop(); + status = _mock_deva.start(); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "start response: " << status.error_str() << "(" << status.error_code() << ")" << std::endl; + + status = _mock_deva.wait_for_leader(&leader); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "leader: " << leader << std::endl; + + pain::proto::deva::ReadDirResponse readdir_response; + status = readdir("/test", &readdir_response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "readdir response: " << readdir_response.DebugString() << std::endl; + ASSERT_EQ(readdir_response.entries_size(), 1); + EXPECT_EQ(readdir_response.entries(0).name(), "test.txt"); + EXPECT_EQ(readdir_response.entries(0).type(), pain::proto::FileType::FILE_TYPE_FILE); +} + } // namespace diff --git a/src/deva/test/test_main.cc b/src/deva/test/test_main.cc index b5027bb..3e46897 100644 --- a/src/deva/test/test_main.cc +++ b/src/deva/test/test_main.cc @@ -26,5 +26,7 @@ int main(int argc, char** argv) { pain::cleanup_tracer(); }); - return RUN_ALL_TESTS(); + auto ret = RUN_ALL_TESTS(); + PLOG_INFO(("desc", "test finished")("ret", ret)); + return ret; }