// SPDX-License-Identifier: GPL-3.0-or-later // SPDX-FileCopyrightText: 2019-2026 Ivan Baidakou #include #include "db/error_code.h" #include "model/diff/iterative_controller.h" #include "test-utils.h" #include "diff-builder.h" #include "model/diff/load/commit.h" #include "model/diff/peer/cluster_update.h" #include "model/diff/contact/ignored_connected.h" #include "model/diff/contact/unknown_connected.h" #include "proto/proto-helpers-impl.hpp" #include "test_supervisor.h" #include "access.h" #include "model/cluster.h" #include "db/utils.h" #include "net/db_actor.h" #include "net/names.h" #include #include using namespace syncspirit; using namespace syncspirit::db; using namespace syncspirit::test; using namespace syncspirit::model; using namespace syncspirit::net; namespace fs = std::filesystem; namespace { struct env {}; } // namespace namespace syncspirit::net { template <> inline auto &db_actor_t::access() noexcept { return env; } } // namespace syncspirit::net namespace { struct db_supervisor_t : supervisor_t { using parent_t = supervisor_t; using parent_t::parent_t; void on_package(bouncer::message::package_t &msg) noexcept { LOG_TRACE(log, "on package"); ++packages; put(std::move(msg.payload)); } int packages = 0; }; struct fixture_t { using load_cluster_success_t = net::message::load_cluster_success_t; using load_cluster_fail_t = net::message::load_cluster_fail_t; using stats_msg_t = net::message::db_info_response_t; using stats_msg_ptr_t = r::intrusive_ptr_t; using supervisor_ptr_t = r::intrusive_ptr_t; fixture_t() noexcept : root_path{unique_path()}, path_quard{root_path} { test::init_logging(); bfs::create_directory(root_path); } fixture_t(fixture_t &source) = delete; fixture_t(fixture_t &&source) noexcept : root_path(std::move(source.root_path)), path_quard(std::move(source.path_quard)) {} virtual configure_callback_t configure() noexcept { return [&](r::plugin::plugin_base_t &plugin) { plugin.template with_casted([&](auto &p) { p.subscribe_actor(r::lambda([&](load_cluster_success_t &msg) { load_diff = msg.payload.diff; spdlog::info("received load cluster diff"); sup->send(db_addr, load_diff, nullptr); })); p.subscribe_actor(r::lambda([&](load_cluster_fail_t &msg) { ee = msg.payload.ee; spdlog::warn("fail loading cluster: {}", ee->message({})); })); p.subscribe_actor(r::lambda([&](stats_msg_t &msg) { stats_reply = &msg; })); }); }; } cluster_ptr_t make_cluster(bool add_peer = true) noexcept { auto my_id = "KHQNO2S-5QSILRK-YX4JZZ4-7L77APM-QNVGZJT-EKU7IFI-PNEPBMY-4MXFMQD"; auto my_device_id = device_id_t::from_string(my_id).value(); my_device = device_t::create(my_device_id, "my-device").value(); auto cluster = cluster_ptr_t(new cluster_t(my_device, 1)); cluster->get_devices().put(my_device); auto peer_id = "VUV42CZ-IQD5A37-RPEBPM4-VVQK6E4-6WSKC7B-PVJQHHD-4PZD44V-ENC6WAZ"; auto peer_device_id = device_id_t::from_string(peer_id).value(); peer_device = device_t::create(peer_device_id, "peer-device").value(); if (add_peer) { cluster->get_devices().put(peer_device); } return cluster; } virtual config::db_config_t make_config() noexcept { return config::db_config_t{1024 * 1024, 0, 64}; } virtual supervisor_ptr_t make_supervisor(r::system_context_t &ctx) noexcept { return ctx.create_supervisor().timeout(timeout).create_registry().finish(); } virtual fixture_t &run() noexcept { cluster = make_cluster(); r::system_context_t ctx; sup = make_supervisor(ctx); sup->cluster = cluster; sup->configure_callback = configure(); sup->start(); sup->do_process(); CHECK(static_cast(sup.get())->access() == r::state_t::OPERATIONAL); launch_db(); main(); load_diff.reset(); CHECK(get_reading_txn() == 0); sup->shutdown(); sup->do_process(); CHECK(static_cast(sup.get())->access() == r::state_t::SHUT_DOWN); return *this; } int get_reading_txn() { int counter = 0; auto state = static_cast(db_actor.get())->access(); if (state != r::state_t::SHUT_DOWN) { auto &db_env = db_actor->access(); auto reader = [](void *ctx, int, int, mdbx_pid_t, mdbx_tid_t, uint64_t, uint64_t, size_t, size_t) noexcept -> int { ++(*reinterpret_cast(ctx)); return 0; }; auto r = mdbx_reader_list(db_env, reader, &counter); CHECK(((r == MDBX_RESULT_TRUE) || (r == 0))); if (r > 0) { auto ec = db::make_error_code(r); auto log = utils::get_logger(""); LOG_INFO(log, "mdbx_reader_list, error: {}", ec.message()); } } return counter; } virtual std::uint32_t get_max_files_per_diff() const { return 32; } virtual void launch_db() { db_actor = sup->create_actor() .bouncer_address(sup->get_address()) .cluster(cluster) .db_dir(root_path) .db_config(make_config()) .timeout(timeout) .max_files_per_diff(get_max_files_per_diff()) .finish(); db_addr = db_actor->get_address(); sup->do_process(); CHECK(static_cast(db_actor.get())->access() == r::state_t::OPERATIONAL); } virtual void main() noexcept {} r::address_ptr_t db_addr; r::pt::time_duration timeout = r::pt::millisec{10}; cluster_ptr_t cluster; device_ptr_t peer_device; device_ptr_t my_device; supervisor_ptr_t sup; r::intrusive_ptr_t db_actor; bfs::path root_path; test::path_guard_t path_quard; r::system_context_t ctx; model::diff::cluster_diff_ptr_t load_diff; r::extended_error_ptr_t ee; stats_msg_ptr_t stats_reply; }; } // namespace void test_db_population() { struct F : fixture_t { void main() noexcept override { auto &db_env = db_actor->access(); auto txn_opt = db::make_transaction(db::transaction_type_t::RW, db_env); REQUIRE(txn_opt); auto &txn = txn_opt.value(); auto load_opt = db::load(db::prefix::device, txn); REQUIRE(load_opt); auto &values = load_opt.value(); REQUIRE(values.size() == 1); } }; F().run(); } void test_loading_empty_db() { struct F : fixture_t { void main() noexcept override { CHECK(get_reading_txn() == 1); REQUIRE(load_diff); REQUIRE(load_diff->apply(*sup, {})); sup->do_process(); load_diff.reset(); CHECK(get_reading_txn() == 0); auto devices = cluster->get_devices(); REQUIRE(devices.size() == 2); REQUIRE(devices.by_sha256(my_device->device_id().get_sha256())); REQUIRE(devices.by_sha256(peer_device->device_id().get_sha256())); sup->request(sup->get_address()).send(timeout); sup->do_process(); REQUIRE(stats_reply); auto &stats = stats_reply->payload.res; CHECK(stats.entries > 0); CHECK(stats.page_size > 0); } }; F().run(); } void test_forget_to_commit_other_thread() { struct F : fixture_t { void main() noexcept override { using namespace model::diff; struct V : diff::cluster_visitor_t { MDBX_txn *txn; outcome::result operator()(const load::commit_t &diff, void *custom) noexcept { auto msg = static_cast(diff.commit_message.get()); txn = msg->payload.txn.txn; return outcome::success(); } }; CHECK(get_reading_txn() == 1); REQUIRE(load_diff); auto visitor = V(); REQUIRE(load_diff->visit(visitor, {})); REQUIRE(visitor.txn); auto diff = std::move(load_diff); std::thread([diff = std::move(diff)]() mutable { diff.reset(); }).join(); CHECK(get_reading_txn() == 1); mdbx_txn_commit(visitor.txn); sup->do_process(); } }; F().run(); } void test_forget_to_commit_own_thread() { struct F : fixture_t { void main() noexcept override { CHECK(get_reading_txn() == 1); load_diff.reset(); CHECK(get_reading_txn() == 0); sup->do_process(); } }; F().run(); } void test_folder_upserting() { struct F : fixture_t { void main() noexcept override { auto folder_id = "1234-5678"; auto builder = diff_builder_t(*cluster); builder.upsert_folder(folder_id, "/my/path", "my-label").apply(*sup); auto folder = cluster->get_folders().by_id(folder_id); REQUIRE(folder); REQUIRE(folder->get_folder_infos().by_device(*cluster->get_device())); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder_clone = cluster_clone->get_folders().by_id(folder->get_id()); REQUIRE(folder_clone); REQUIRE(folder.get() != folder_clone.get()); REQUIRE(folder_clone->get_label() == "my-label"); REQUIRE(folder_clone->get_path().string() == "/my/path"); REQUIRE(folder_clone->get_folder_infos().size() == 1); REQUIRE(folder_clone->get_folder_infos().by_device(*cluster->get_device())); } SECTION("upserting") { builder.upsert_folder(folder_id, "/my/path", "my-label-2").apply(*sup); REQUIRE(folder->get_label() == "my-label-2"); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder_clone = cluster_clone->get_folders().by_id(folder->get_id()); REQUIRE(folder_clone->get_label() == "my-label-2"); } } }; F().run(); } void test_unknown_and_ignored_devices_1() { struct F : fixture_t { void main() noexcept override { auto d_id1 = device_id_t::from_string("LYXKCHX-VI3NYZR-ALCJBHF-WMZYSPK-QG6QJA3-MPFYMSO-U56GTUK-NA2MIAW").value(); auto d_id2 = device_id_t::from_string("XBOWTOU-Y7H6RM6-D7WT3UB-7P2DZ5G-R6GNZG6-T5CCG54-SGVF3U5-LBM7RQB").value(); db::SomeDevice sd_1; db::set_name(sd_1, "x1"); auto unknown_device = pending_device_t::create(d_id1, sd_1).value(); db::SomeDevice sd_2; db::set_name(sd_2, "x2"); auto ignored_device = ignored_device_t::create(d_id2, sd_2).value(); auto builder = diff_builder_t(*cluster); builder.add_unknown_device(d_id1, sd_1).add_ignored_device(d_id2, sd_2).apply(*sup); REQUIRE(cluster->get_pending_devices().size() == 1); REQUIRE(cluster->get_ignored_devices().size() == 1); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); CHECK(cluster_clone->get_pending_devices().by_sha256(d_id1.get_sha256())); CHECK(cluster_clone->get_ignored_devices().by_sha256(d_id2.get_sha256())); } db::set_name(sd_1, "x1_2"); db::set_name(sd_2, "x2_2"); auto diff = model::diff::cluster_diff_ptr_t{}; diff = new model::diff::contact::unknown_connected_t(*cluster, d_id1, sd_1); sup->send(sup->get_address(), std::move(diff), nullptr); diff = new model::diff::contact::ignored_connected_t(*cluster, d_id2, sd_2); sup->send(sup->get_address(), std::move(diff), nullptr); sup->do_process(); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_pending_devices().size() == 1); REQUIRE(cluster_clone->get_ignored_devices().size() == 1); auto unknown = cluster_clone->get_pending_devices().by_sha256(d_id1.get_sha256()); auto ignored = cluster_clone->get_ignored_devices().by_sha256(d_id2.get_sha256()); REQUIRE(unknown); REQUIRE(ignored); CHECK(unknown->get_name() == "x1_2"); CHECK(ignored->get_name() == "x2_2"); } builder.remove_unknown_device(*unknown_device).remove_ignored_device(*ignored_device).apply(*sup); REQUIRE(cluster->get_pending_devices().size() == 0); REQUIRE(cluster->get_ignored_devices().size() == 0); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_pending_devices().size() == 0); REQUIRE(cluster_clone->get_ignored_devices().size() == 0); } } }; F().run(); } void test_unknown_and_ignored_devices_2() { struct F : fixture_t { void main() noexcept override { auto d_id = device_id_t::from_string("LYXKCHX-VI3NYZR-ALCJBHF-WMZYSPK-QG6QJA3-MPFYMSO-U56GTUK-NA2MIAW").value(); db::SomeDevice sd; db::set_name(sd, "x1"); auto builder = diff_builder_t(*cluster); builder.add_unknown_device(d_id, sd).apply(*sup); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); CHECK(cluster_clone->get_pending_devices().by_sha256(d_id.get_sha256())); CHECK(!cluster_clone->get_ignored_devices().by_sha256(d_id.get_sha256())); } builder.add_ignored_device(d_id, sd).apply(*sup); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); CHECK(!cluster_clone->get_pending_devices().by_sha256(d_id.get_sha256())); CHECK(cluster_clone->get_ignored_devices().by_sha256(d_id.get_sha256())); } } }; F().run(); } void test_peer_updating() { struct F : fixture_t { void main() noexcept override { auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id(), "some_name", "some-cn", true).apply(*sup); auto sha256 = peer_device->device_id().get_sha256(); auto device = cluster->get_devices().by_sha256(sha256); REQUIRE(device); CHECK(device->get_name() == "some_name"); CHECK(device->get_cert_name() == "some-cn"); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_devices().size() == 2); auto device_clone = cluster_clone->get_devices().by_sha256(sha256); REQUIRE(device_clone); REQUIRE(device.get() != device_clone.get()); CHECK(device_clone->get_name() == "some_name"); CHECK(device_clone->get_cert_name() == "some-cn"); } }; F().run(); } void test_folder_sharing() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto folder_id = "1234-5678"; auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .apply(*sup) .configure_cluster(sha256) .add(sha256, folder_id, 5, 4) .finish() .apply(*sup); REQUIRE(cluster->get_pending_folders().size() == 1); builder.share_folder(sha256, folder_id).apply(*sup); REQUIRE(cluster->get_pending_folders().size() == 0); CHECK(static_cast(db_actor.get())->access() == r::state_t::OPERATIONAL); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto peer_device = cluster_clone->get_devices().by_sha256(sha256); REQUIRE(peer_device); auto folder = cluster_clone->get_folders().by_id(folder_id); REQUIRE(folder); REQUIRE(folder->get_folder_infos().size() == 2); auto fi = folder->get_folder_infos().by_device(*peer_device); REQUIRE(fi); CHECK(fi->get_index() == 5); CHECK(fi->get_max_sequence() == 0); REQUIRE(cluster_clone->get_pending_folders().size() == 0); } }; F().run(); } void test_cluster_update_and_remove() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto folder_id = "1234-5678"; auto unknown_folder_id = "5678-999"; auto file = proto::FileInfo(); proto::set_name(file, "a.txt"); proto::set_size(file, 5ul); proto::set_block_size(file, 5ul); proto::set_sequence(file, 6ul); auto b_hash = utils::sha256_digest(as_bytes("12345")).value(); auto &b = proto::add_blocks(file); proto::set_size(b, 5ul); proto::set_hash(b, b_hash); auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .configure_cluster(sha256) .add(sha256, folder_id, 5, proto::get_sequence(file)) .add(sha256, unknown_folder_id, 5, 5) .finish() .apply(*sup) .share_folder(sha256, folder_id) .apply(*sup) .make_index(sha256, folder_id) .add(file, peer_device) .finish() .apply(*sup); REQUIRE(cluster->get_blocks().size() == 1); auto block = cluster->get_blocks().by_hash(b_hash); REQUIRE(block); auto folder = cluster->get_folders().by_id(folder_id); auto peer_folder_info = folder->get_folder_infos().by_device(*peer_device); REQUIRE(peer_folder_info); CHECK(peer_folder_info->get_max_sequence() == 6ul); REQUIRE(peer_folder_info->get_file_infos().size() == 1); auto peer_file = peer_folder_info->get_file_infos().by_name("a.txt"); REQUIRE(peer_file); auto &unknown_folders = cluster->get_pending_folders(); CHECK(std::distance(unknown_folders.begin(), unknown_folders.end()) == 1); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_blocks().size() == 1); CHECK(cluster_clone->get_blocks().by_hash(b_hash)); auto folder = cluster_clone->get_folders().by_id(folder_id); auto peer_folder_info = folder->get_folder_infos().by_device(*peer_device); REQUIRE(peer_folder_info); REQUIRE(peer_folder_info->get_file_infos().size() == 1); REQUIRE(peer_folder_info->get_file_infos().by_name("a.txt")); REQUIRE(cluster_clone->get_pending_folders().size() == 1); } auto pr_msg = proto::ClusterConfig(); auto &pr_f = proto::add_folders(pr_msg); proto::set_id(pr_f, folder_id); auto &pr_device = proto::add_devices(pr_f); proto::set_id(pr_device, peer_device->device_id().get_sha256()); proto::set_max_sequence(pr_device, 1); proto::set_index_id(pr_device, peer_folder_info->get_index() + 1); auto diff = diff::peer::cluster_update_t::create({}, *cluster, *sup->sequencer, *peer_device, pr_msg).value(); sup->send(sup->get_address(), diff, nullptr); sup->do_process(); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_blocks().size() == 0); auto &fis = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(fis.size() == 2); auto folder_info = fis.by_device(*peer_device); REQUIRE(folder_info); REQUIRE(fis.by_device(*cluster->get_device())); REQUIRE(cluster_clone->get_pending_folders().size() == 0); } } }; F().run(); } void test_unshare_and_remove_folder() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto folder_id = "1234-5678"; auto file = proto::FileInfo(); proto::set_name(file, "a.txt"); proto::set_size(file, 5ul); proto::set_block_size(file, 5ul); proto::set_sequence(file, 6ul); auto b_hash = utils::sha256_digest(as_bytes("12345")).value(); auto &b = proto::add_blocks(file); proto::set_size(b, 5ul); proto::set_hash(b, b_hash); auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .apply(*sup) .configure_cluster(sha256) .add(sha256, folder_id, 5, proto::get_sequence(file)) .finish() .share_folder(sha256, folder_id) .apply(*sup) .make_index(sha256, folder_id) .add(file, peer_device) .finish() .apply(*sup); REQUIRE(cluster->get_blocks().size() == 1); auto block = cluster->get_blocks().by_hash(b_hash); REQUIRE(block); auto folder = cluster->get_folders().by_id(folder_id); auto peer_folder_info = folder->get_folder_infos().by_device(*peer_device); REQUIRE(peer_folder_info); CHECK(peer_folder_info->get_max_sequence() == 6ul); REQUIRE(peer_folder_info->get_file_infos().size() == 1); auto peer_file = peer_folder_info->get_file_infos().by_name("a.txt"); REQUIRE(peer_file); SECTION("unshare") { builder.unshare_folder(*peer_folder_info).apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto &fis = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(fis.size() == 1); REQUIRE(!fis.by_device(*peer_device)); REQUIRE(fis.by_device(*cluster->get_device())); REQUIRE(cluster_clone->get_blocks().size() == 0); } } SECTION("remove") { builder.remove_folder(*folder).apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_folders().size() == 0); REQUIRE(cluster_clone->get_blocks().size() == 0); } } } }; F().run(); } void test_remote_copy() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto folder_id = "1234-5678"; auto file = proto::FileInfo(); proto::set_name(file, "a.txt"); proto::set_sequence(file, 6ul); auto &v = proto::get_version(file); proto::add_counters(v, proto::Counter(peer_device->device_id().get_uint(), 1)); auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .configure_cluster(sha256) .add(sha256, folder_id, 5, proto::get_sequence(file)) .finish() .apply(*sup) .share_folder(sha256, folder_id) .apply(*sup); auto folder = cluster->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto folder_peer = folder->get_folder_infos().by_device(*peer_device); SECTION("file without blocks") { builder.make_index(sha256, folder_id).add(file, peer_device).finish().apply(*sup); auto file_peer = folder_peer->get_file_infos().by_name(proto::get_name(file)); REQUIRE(file_peer); builder.remote_copy(*file_peer, *folder_peer).apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_blocks().size() == 0); auto &fis = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(fis.size() == 2); auto folder_info_clone = fis.by_device(*cluster_clone->get_device()); auto file_clone = folder_info_clone->get_file_infos().by_name(proto::get_name(file)); REQUIRE(file_clone); REQUIRE(file_clone->get_name()->get_full_name() == proto::get_name(file)); REQUIRE(file_clone->iterate_blocks().get_total() == 0); REQUIRE(file_clone->get_sequence() == 1); REQUIRE(folder_info_clone->get_max_sequence() == 1); } } SECTION("file with blocks") { proto::set_size(file, 5ul); proto::set_block_size(file, 5ul); auto b_hash = utils::sha256_digest(as_bytes("12345")).value(); auto &b = proto::add_blocks(file); proto::set_size(b, 5ul); proto::set_hash(b, b_hash); builder.make_index(sha256, folder_id).add(file, peer_device).finish().apply(*sup); auto folder = cluster->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto folder_peer = folder->get_folder_infos().by_device(*peer_device); auto file_peer = folder_peer->get_file_infos().by_name(proto::get_name(file)); REQUIRE(file_peer); builder.remote_copy(*file_peer, *folder_peer).apply(*sup); REQUIRE(folder_my->get_max_sequence() == 1); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_blocks().size() == 1); auto &fis = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(fis.size() == 2); auto folder_info_clone = fis.by_device(*cluster_clone->get_device()); auto file_clone = folder_info_clone->get_file_infos().by_name(proto::get_name(file)); REQUIRE(file_clone); REQUIRE(file_clone->get_name()->get_full_name() == proto::get_name(file)); REQUIRE(file_clone->iterate_blocks().get_total() == 1); REQUIRE(file_clone->get_sequence() == 1); REQUIRE(folder_info_clone->get_max_sequence() == 1); } file_peer = folder_peer->get_file_infos().by_name(proto::get_name(file)); file_peer->mark_local_available(0); REQUIRE(file_peer->is_locally_available()); builder.remote_copy(*file_peer, *folder_peer).apply(*sup); { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); REQUIRE(cluster_clone->get_blocks().size() == 1); auto &fis = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(fis.size() == 2); auto folder_info_clone = fis.by_device(*cluster_clone->get_device()); auto file_clone = folder_info_clone->get_file_infos().by_name(proto::get_name(file)); REQUIRE(file_clone); REQUIRE(file_clone->get_name()->get_full_name() == proto::get_name(file)); REQUIRE(file_clone->iterate_blocks().get_total() == 1); REQUIRE(file_clone->iterate_blocks().next()); REQUIRE(file_clone->get_sequence() == 2); REQUIRE(folder_info_clone->get_max_sequence() == 2); } } } }; F().run(); } void test_local_update() { struct F : fixture_t { void main() noexcept override { auto folder_id = "1234-5678"; auto pr_file = proto::FileInfo(); proto::set_name(pr_file, "a.txt"); proto::set_size(pr_file, 5ul); auto hash = utils::sha256_digest(as_bytes("12345")).value(); auto &pr_block = proto::add_blocks(pr_file); proto::set_size(pr_block, 5ul); proto::set_hash(pr_block, hash); auto builder = diff_builder_t(*cluster); builder.upsert_folder(folder_id, "/my/path").apply(*sup).local_update(folder_id, pr_file).apply(*sup); SECTION("check saved file with new blocks") { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder = cluster_clone->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto file = folder_my->get_file_infos().by_name("a.txt"); REQUIRE(file); CHECK(cluster_clone->get_blocks().size() == 1); CHECK(file->iterate_blocks().get_total() == 1); } proto::set_deleted(pr_file, true); proto::set_size(pr_file, 0); proto::clear_blocks(pr_file); builder.local_update(folder_id, pr_file).apply(*sup); SECTION("check deleted blocks") { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder = cluster_clone->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto file = folder_my->get_file_infos().by_name("a.txt"); REQUIRE(file); CHECK(file->is_deleted()); CHECK(cluster_clone->get_blocks().size() == 0); CHECK(file->iterate_blocks().get_total() == 0); } } }; F().run(); }; void test_peer_going_offline() { struct F : fixture_t { void main() noexcept override { auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()).apply(*sup); auto sha256 = peer_device->device_id().get_sha256(); db::Device db_peer; auto peer = cluster->get_devices().by_sha256(sha256); REQUIRE(db::get_last_seen(db_peer) == 0); peer->update_state(peer->get_state().connecting().connected().online("tcp://1.1.1.1:1")); builder.update_state(*peer, {}, peer->get_state().offline()).apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto peer_clone = cluster_clone->get_devices().by_sha256(sha256); db::Device db_peer_clone; peer_clone->serialize(db_peer_clone); CHECK((peer->get_last_seen() - peer_clone->get_last_seen()).total_seconds() < 2); REQUIRE(db::get_last_seen(db_peer_clone) != 0); } }; F().run(); }; void test_remove_peer() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto folder_id = "1234-5678"; auto unknown_folder_id = "5678-999"; auto file = proto::FileInfo(); proto::set_name(file, "a.txt"); proto::set_size(file, 5ul); proto::set_block_size(file, 5ul); proto::set_sequence(file, 6ul); auto b_hash = utils::sha256_digest(as_bytes("12345")).value(); auto &b = proto::add_blocks(file); proto::set_size(b, 5ul); proto::set_hash(b, b_hash); auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .configure_cluster(sha256) .add(sha256, folder_id, 5, proto::get_sequence(file)) .add(sha256, unknown_folder_id, 5, 5) .finish() .apply(*sup) .share_folder(sha256, folder_id) .apply(*sup) .make_index(sha256, folder_id) .add(file, peer_device) .finish() .apply(*sup); CHECK(cluster->get_pending_folders().size() == 1); REQUIRE(cluster->get_blocks().size() == 1); auto block = cluster->get_blocks().by_hash(b_hash); REQUIRE(block); auto folder = cluster->get_folders().by_id(folder_id); auto peer_folder_info = folder->get_folder_infos().by_device(*peer_device); REQUIRE(peer_folder_info); CHECK(peer_folder_info->get_max_sequence() == 6ul); REQUIRE(peer_folder_info->get_file_infos().size() == 1); auto peer_file = peer_folder_info->get_file_infos().by_name("a.txt"); REQUIRE(peer_file); builder.remove_peer(*peer_device).apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); { auto cluster_clone = make_cluster(false); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); CHECK(cluster_clone->get_pending_folders().size() == 0); CHECK(cluster_clone->get_devices().size() == 1); REQUIRE(cluster_clone->get_blocks().size() == 0); } } }; F().run(); } void test_update_peer() { struct F : fixture_t { void main() noexcept override { auto builder = diff_builder_t(*cluster); db::SomeDevice db; db::set_name(db, "x1"); SECTION("unknown device is removed") { builder.add_unknown_device(peer_device->device_id(), db) .apply(*sup) .update_peer(peer_device->device_id(), "p1") .apply(*sup); } SECTION("ignored device is removed") { builder.add_ignored_device(peer_device->device_id(), db) .apply(*sup) .update_peer(peer_device->device_id(), "p1") .apply(*sup); } { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); CHECK(cluster_clone->get_pending_devices().size() == 0); CHECK(cluster_clone->get_ignored_devices().size() == 0); CHECK(cluster_clone->get_devices().size() == 2); } } }; F().run(); } void test_peer_3_folders_6_files() { struct F : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto f1_id = "123"; auto f2_id = "356"; auto f3_id = "789"; auto next_sequence = 6ul; auto make_file = [&](std::string name) { auto file = proto::FileInfo(); proto::set_name(file, name); proto::set_sequence(file, ++next_sequence); auto &v = proto::get_version(file); proto::add_counters(v, proto::Counter(peer_device->device_id().get_uint(), 1)); return file; }; auto builder = diff_builder_t(*cluster); // clang-format off builder .update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(f1_id, "/my/path1", "my-label1") .upsert_folder(f2_id, "/my/path2", "my-label2") .upsert_folder(f3_id, "/my/path3", "my-label3") .configure_cluster(sha256) .add(sha256, f1_id, 10, 5) .add(sha256, f2_id, 11, 5) .add(sha256, f3_id, 12, 5) .finish() .apply(*sup) .share_folder(sha256, f1_id) .share_folder(sha256, f2_id) .share_folder(sha256, f3_id) .apply(*sup) .make_index(sha256, f1_id).add(make_file("f1.1"), peer_device).add(make_file("f1.2"), peer_device).finish() .make_index(sha256, f2_id).add(make_file("f2.1"), peer_device).add(make_file("f2.2"), peer_device).finish() .make_index(sha256, f3_id).add(make_file("f3.1"), peer_device).add(make_file("f3.2"), peer_device).finish() .apply(*sup); // clang-format on { auto get_peer_file = [&](std::string_view folder_id, std::string_view name) -> std::pair { auto folder = cluster->get_folders().by_id(folder_id); auto folder_info = folder->get_folder_infos().by_device(*peer_device); auto file = folder_info->get_file_infos().by_name(name); return {file, folder_info.get()}; }; auto [file_11, fi_11] = get_peer_file(f1_id, "f1.1"); auto [file_12, fi_12] = get_peer_file(f1_id, "f1.2"); auto [file_21, fi_21] = get_peer_file(f2_id, "f2.1"); auto [file_22, fi_22] = get_peer_file(f2_id, "f2.2"); auto [file_31, fi_31] = get_peer_file(f3_id, "f3.1"); auto [file_32, fi_32] = get_peer_file(f3_id, "f3.2"); REQUIRE(file_11); REQUIRE(file_12); REQUIRE(file_21); REQUIRE(file_22); REQUIRE(file_31); REQUIRE(file_32); // clang-format off builder .remote_copy(*file_11, *fi_11).remote_copy(*file_12, *fi_12) .remote_copy(*file_21, *fi_21).remote_copy(*file_22, *fi_22) .remote_copy(*file_31, *fi_31).remote_copy(*file_32, *fi_32) .apply(*sup); // clang-format on } { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto get_my_file = [&](std::string_view folder_id, std::string_view name) { auto folder = cluster->get_folders().by_id(folder_id); auto folder_info = folder->get_folder_infos().by_device(*my_device); auto file = folder_info->get_file_infos().by_name(name); return file; }; auto file_11 = get_my_file(f1_id, "f1.1"); auto file_12 = get_my_file(f1_id, "f1.2"); auto file_21 = get_my_file(f2_id, "f2.1"); auto file_22 = get_my_file(f2_id, "f2.2"); auto file_31 = get_my_file(f3_id, "f3.1"); auto file_32 = get_my_file(f3_id, "f3.2"); REQUIRE(file_11); REQUIRE(file_12); REQUIRE(file_21); REQUIRE(file_22); REQUIRE(file_31); REQUIRE(file_32); } } }; F().run(); } void test_db_migration_1_2() { static constexpr auto folder_id = "1234-5678"; struct F1 : fixture_t { void main() noexcept override { auto sha256 = peer_device->device_id().get_sha256(); auto builder = diff_builder_t(*cluster); builder.update_peer(peer_device->device_id()) .apply(*sup) .upsert_folder(folder_id, "/my/path") .apply(*sup) .share_folder(sha256, folder_id) .apply(*sup); auto &folder_infos = cluster->get_folders().by_id(folder_id)->get_folder_infos(); auto fi_my = folder_infos.by_device(*my_device); auto fi_peer = folder_infos.by_device(*peer_device); REQUIRE(fi_my->is_introduced_by(my_device->device_id())); REQUIRE(fi_peer->is_introduced_by(my_device->device_id())); auto &db_env = db_actor->access(); auto txn_opt = db::make_transaction(db::transaction_type_t::RW, db_env); REQUIRE(txn_opt); auto &txn = txn_opt.value(); REQUIRE(db::save_version(1, txn)); auto save = [&](const model::folder_info_t &fi) -> outcome::result { auto fi_db = db::FolderInfo(); auto fi_key = fi.get_key(); fi.serialize(fi_db); db::set_introducer_device_key(fi_db, {}); auto fi_data = db::encode(fi_db); return db::save({fi_key, fi_data}, txn); }; REQUIRE(save(*fi_my)); REQUIRE(save(*fi_peer)); REQUIRE(txn.commit()); } }; struct F2 : fixture_t { F2(fixture_t &other) : fixture_t(std::move(other)) {} void main() noexcept override { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto f_cloned = cluster_clone->get_folders().by_id(folder_id); auto &folder_infos = f_cloned->get_folder_infos(); auto fi_my = folder_infos.by_device(*my_device); auto fi_peer = folder_infos.by_device(*peer_device); CHECK(fi_my->is_introduced_by(my_device->device_id())); CHECK(fi_peer->is_introduced_by(my_device->device_id())); } }; F2(F1().run()).run(); } void test_db_migration_2_3() { static constexpr auto folder_id = "1234-5678"; struct F1 : fixture_t { void main() noexcept override { // clang-format off using BlockInfoEx = pp::message< pp::uint32_field <"weak_hash", 1>, pp::int32_field <"size", 2> >; // clang-format on auto bi = BlockInfoEx{1, 0x1234}; auto new_value = syncspirit::details::generic_encode(bi); unsigned char key[model::block_info_t::digest_length + 1] = {0}; key[0] = db::prefix::block_info; auto &db_env = db_actor->access(); auto txn_opt = db::make_transaction(db::transaction_type_t::RW, db_env); REQUIRE(txn_opt); auto &txn = txn_opt.value(); REQUIRE(db::save_version(2, txn)); REQUIRE(db::save({key, new_value}, txn)); REQUIRE(txn.commit()); } }; struct F2 : fixture_t { F2(fixture_t &other) : fixture_t(std::move(other)) {} void main() noexcept override { load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto &blocks = cluster_clone->get_blocks(); REQUIRE(blocks.size() == 1); auto b = *blocks.begin(); CHECK(b->get_size() == 0x1234); } }; F2(F1().run()).run(); } void test_corrupted_file() { struct F : fixture_t { void main() noexcept override { auto folder_id = "1234-5678"; auto builder = diff_builder_t(*cluster); builder.upsert_folder(folder_id, "/my/path").apply(*sup); CHECK(static_cast(db_actor.get())->access() == r::state_t::OPERATIONAL); auto pr_file = proto::FileInfo(); proto::set_name(pr_file, "a.txt"); proto::set_size(pr_file, 5ul); auto hash = utils::sha256_digest(as_bytes("12345")).value(); auto &pr_block = proto::add_blocks(pr_file); proto::set_size(pr_block, 5ul); proto::set_hash(pr_block, hash); builder.local_update(folder_id, pr_file).apply(*sup); auto block = *cluster->get_blocks().begin(); auto folder = cluster->get_folders().by_id(folder_id); auto fi = folder->get_folder_infos().by_device(*my_device); auto file = fi->get_file_infos().by_name("a.txt"); REQUIRE(file); auto &db_env = db_actor->access(); auto txn_opt = db::make_transaction(db::transaction_type_t::RW, db_env); REQUIRE(txn_opt); auto &txn = txn_opt.value(); auto diff = model::diff::cluster_diff_ptr_t(); SECTION("missing block") { auto key = test::make_key(block); REQUIRE(db::remove(key, txn)); REQUIRE(!txn.commit().has_error()); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(false); auto controller = make_apply_controller(cluster_clone); diff = load_diff; REQUIRE(diff->apply(*controller, {})); auto &folder_infos = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); auto folder_info = folder_infos.by_device(*my_device); CHECK(folder_info->get_file_infos().size() == 0); auto &blocks = cluster_clone->get_blocks(); CHECK(blocks.size() == 0); } SECTION("missing folder") { auto key = fi->get_key(); REQUIRE(db::remove(key, txn)); REQUIRE(!txn.commit().has_error()); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(false); auto controller = make_apply_controller(cluster_clone); diff = load_diff; REQUIRE(diff->apply(*controller, {})); auto &folder_infos = cluster_clone->get_folders().by_id(folder_id)->get_folder_infos(); REQUIRE(folder_infos.size() == 0); auto &blocks = cluster_clone->get_blocks(); CHECK(blocks.size() == 1); } REQUIRE(diff); sup->send(sup->get_address(), std::move(diff)); sup->do_process(); txn_opt = db::make_transaction(db::transaction_type_t::RW, db_env); REQUIRE(txn_opt); auto files_opt = db::load(db::prefix::file_info, txn); REQUIRE(files_opt); CHECK(files_opt.value().size() == 0); } }; F().run(); } void test_flush_on_shutdown() { struct F : fixture_t { config::db_config_t make_config() noexcept override { return config::db_config_t{1024 * 1024, 100, 64}; } void main() noexcept override { auto folder_id = "1234-5678"; auto pr_file = proto::FileInfo(); proto::set_name(pr_file, "a.txt"); proto::set_size(pr_file, 5ul); auto hash = utils::sha256_digest(as_bytes("12345")).value(); auto &pr_block = proto::add_blocks(pr_file); proto::set_size(pr_block, 5ul); proto::set_hash(pr_block, hash); auto builder = diff_builder_t(*cluster); builder.upsert_folder(folder_id, "/my/path").apply(*sup).local_update(folder_id, pr_file).apply(*sup); load_diff = {}; db_actor->do_shutdown(); sup->do_process(); launch_db(); SECTION("loaded data") { sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder = cluster_clone->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto file = folder_my->get_file_infos().by_name("a.txt"); REQUIRE(file); CHECK(file->iterate_blocks().get_total() == 1); CHECK(cluster_clone->get_blocks().size() == 1); } } }; F().run(); }; struct min_concurrencty_fixture_t : fixture_t { using parent_t = fixture_t; using parent_t::parent_t; std::uint32_t get_max_files_per_diff() const override { return 1; } }; void test_iterative_loading() { struct F : min_concurrencty_fixture_t { config::db_config_t make_config() noexcept override { return config::db_config_t{1024 * 1024, 0, 1}; } void main() noexcept override { auto builder = diff_builder_t(*cluster); auto folder_id = "1234-5678"; builder.upsert_folder(folder_id, "/my/path").apply(*sup); for (int i = 0; i < 20; ++i) { auto pr_file = proto::FileInfo(); proto::set_name(pr_file, fmt::format("file-{:03}.bin", i)); builder.local_update(folder_id, pr_file); } builder.apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = make_apply_controller(cluster_clone); REQUIRE(load_diff->apply(*controller, {})); auto folder = cluster_clone->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); CHECK(folder_my->get_file_infos().size() == 20); CHECK(sup->packages >= 21); } }; F().run(); }; void test_iterative_loading_interrupt() { struct F : min_concurrencty_fixture_t { config::db_config_t make_config() noexcept override { return config::db_config_t{1024 * 1024, 0, 1}; } supervisor_ptr_t make_supervisor(r::system_context_t &ctx) noexcept override { struct sup_t : db_supervisor_t { using parent_t = db_supervisor_t; using parent_t::parent_t; void on_package(bouncer::message::package_t &msg) noexcept { parent_t::on_package(msg); if (packages == 5) { fixture->db_actor->do_shutdown(); } } fixture_t *fixture = {}; }; auto sup = ctx.create_supervisor().timeout(timeout).create_registry().finish(); sup->fixture = this; return sup; } void main() noexcept override { auto builder = diff_builder_t(*cluster); auto folder_id = "1234-5678"; builder.upsert_folder(folder_id, "/my/path").apply(*sup); for (int i = 0; i < 20; ++i) { auto pr_file = proto::FileInfo(); proto::set_name(pr_file, fmt::format("file-{:03}.bin", i)); builder.local_update(folder_id, pr_file); } builder.apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(!load_diff); REQUIRE(ee); CHECK(sup->packages < 10); } }; F().run(); }; template using my_controller_base_t = model::diff::iterative_controller_t; void test_iterative_application_interrupt() { struct F : min_concurrencty_fixture_t { config::db_config_t make_config() noexcept override { return config::db_config_t{1024 * 1024, 0, 1}; } void main() noexcept override { struct my_controller_t : my_controller_base_t { using parent_t = my_controller_base_t; using parent_t::bouncer; using parent_t::cluster; using parent_t::resources; my_controller_t(r::actor_config_t &cfg) noexcept : parent_t(this, 0, cfg) {} void configure(r::plugin::plugin_base_t &plugin) noexcept override { plugin.with_casted([&](auto &p) { p.set_identity("my.controller", false); log = utils::get_logger(identity); }); plugin.with_casted([&](auto &p) { p.discover_name(net::names::coordinator, coordinator, true) .link(false) .callback([&](auto phase, auto &ee) { if (!ee && phase == r::plugin::registry_plugin_t::phase_t::linking) { auto p = get_plugin(r::plugin::starter_plugin_t::class_identity); auto plugin = static_cast(p); plugin->subscribe_actor(&my_controller_t::on_model_update, address); } }); }); plugin.with_casted( [&](auto &p) { p.subscribe_actor(&my_controller_t::on_model_interrupt); }, r::plugin::config_phase_t::PREINIT); } void process(model::diff::cluster_diff_t &diff, model::payload::apply_context_t &context) noexcept override { parent_t::process(diff, context); auto &folder = cluster->get_folders().begin()->item; auto &fi = folder->get_folder_infos().begin()->item; if (fi->get_file_infos().size() == 10) { do_shutdown(); resources->acquire(0); } } }; auto builder = diff_builder_t(*cluster); auto folder_id = "1234-5678"; builder.upsert_folder(folder_id, "/my/path").apply(*sup); for (int i = 0; i < 20; ++i) { auto pr_file = proto::FileInfo(); proto::set_name(pr_file, fmt::format("file-{:03}.bin", i)); builder.local_update(folder_id, pr_file); } builder.apply(*sup); load_diff = {}; sup->send(db_addr); sup->do_process(); REQUIRE(load_diff); auto cluster_clone = make_cluster(); auto controller = sup->create_actor().timeout(timeout).finish(); controller->bouncer = sup->get_address(); auto addr = controller->get_address(); sup->do_process(); controller->cluster = cluster_clone; sup->send(addr, std::move(load_diff)); sup->do_process(); auto folder = cluster_clone->get_folders().by_id(folder_id); auto folder_my = folder->get_folder_infos().by_device(*my_device); auto files_sz = folder_my->get_file_infos().size(); CHECK(files_sz >= 10); CHECK(files_sz <= 12); CHECK(static_cast(controller.get())->access() == r::state_t::SHUTTING_DOWN); controller->resources->release(0); sup->do_process(); CHECK(static_cast(controller.get())->access() == r::state_t::SHUT_DOWN); CHECK(get_reading_txn() == 0); } }; F().run(); }; int _init() { REGISTER_TEST_CASE(test_db_population, "test_db_population", "[db]"); REGISTER_TEST_CASE(test_loading_empty_db, "test_loading_empty_db", "[db]"); REGISTER_TEST_CASE(test_forget_to_commit_other_thread, "test_forget_to_commit_other_thread", "[db]"); REGISTER_TEST_CASE(test_forget_to_commit_own_thread, "test_forget_to_commit_own_thread", "[db]"); REGISTER_TEST_CASE(test_unknown_and_ignored_devices_1, "test_unknown_and_ignored_devices_1", "[db]"); REGISTER_TEST_CASE(test_unknown_and_ignored_devices_2, "test_unknown_and_ignored_devices_2", "[db]"); REGISTER_TEST_CASE(test_folder_upserting, "test_folder_upserting", "[db]"); REGISTER_TEST_CASE(test_peer_updating, "test_peer_updating", "[db]"); REGISTER_TEST_CASE(test_folder_sharing, "test_folder_sharing", "[db]"); REGISTER_TEST_CASE(test_cluster_update_and_remove, "test_cluster_update_and_remove", "[db]"); REGISTER_TEST_CASE(test_unshare_and_remove_folder, "test_unshare_and_remove_folder", "[db]"); REGISTER_TEST_CASE(test_remote_copy, "test_remote_copy", "[db]"); REGISTER_TEST_CASE(test_local_update, "test_local_update", "[db]"); REGISTER_TEST_CASE(test_peer_going_offline, "test_peer_going_offline", "[db]"); REGISTER_TEST_CASE(test_remove_peer, "test_remove_peer", "[db]"); REGISTER_TEST_CASE(test_update_peer, "test_update_peer", "[db]"); REGISTER_TEST_CASE(test_peer_3_folders_6_files, "test_peer_3_folders_6_files", "[db]"); REGISTER_TEST_CASE(test_db_migration_1_2, "test_db_migration_1_2", "[db]"); REGISTER_TEST_CASE(test_db_migration_2_3, "test_db_migration_2_3", "[db]"); REGISTER_TEST_CASE(test_corrupted_file, "test_corrupted_file", "[db]"); REGISTER_TEST_CASE(test_flush_on_shutdown, "test_flush_on_shutdown", "[db]"); REGISTER_TEST_CASE(test_iterative_loading, "test_iterative_loading", "[db]"); REGISTER_TEST_CASE(test_iterative_loading_interrupt, "test_iterative_loading_interrupt", "[db]"); REGISTER_TEST_CASE(test_iterative_application_interrupt, "test_iterative_application_interrupt", "[db]"); return 1; } static int v = _init();