From 6e3ab9eecf26e9d245ea67196ab795ae1e4ffcfe Mon Sep 17 00:00:00 2001 From: Ivan Baidakou Date: Sat, 21 Sep 2024 14:33:09 +0300 Subject: [PATCH] fix: core, introduce sync start/finish diffs + react on them in scan_actor & scheduler --- CMakeLists.txt | 2 + src/fs/scan_actor.cpp | 119 ++++++++++-------- src/fs/scan_scheduler.cpp | 23 +++- src/fs/scan_scheduler.h | 2 + src/model/diff/cluster_visitor.cpp | 12 ++ src/model/diff/cluster_visitor.h | 4 + .../diff/local/synchronization_finish.cpp | 25 ++++ src/model/diff/local/synchronization_finish.h | 19 +++ .../diff/local/synchronization_start.cpp | 25 ++++ src/model/diff/local/synchronization_start.h | 19 +++ src/model/folder.cpp | 11 +- src/model/folder.h | 3 + src/model/misc/file_iterator.cpp | 7 -- src/model/misc/file_iterator.h | 3 +- src/net/controller_actor.cpp | 23 +++- src/net/controller_actor.h | 3 + tests/050-file_iterator.cpp | 13 -- tests/075-controller.cpp | 28 +++-- tests/085-scan-scheduler.cpp | 14 ++- tests/086-scan_actor.cpp | 18 +++ tests/diff-builder.cpp | 10 ++ tests/diff-builder.h | 2 + 22 files changed, 291 insertions(+), 94 deletions(-) create mode 100644 src/model/diff/local/synchronization_finish.cpp create mode 100644 src/model/diff/local/synchronization_finish.h create mode 100644 src/model/diff/local/synchronization_start.cpp create mode 100644 src/model/diff/local/synchronization_start.h diff --git a/CMakeLists.txt b/CMakeLists.txt index 39231a47..99bcb731 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -105,6 +105,8 @@ add_library(syncspirit_lib src/model/diff/local/scan_finish.cpp src/model/diff/local/scan_request.cpp src/model/diff/local/scan_start.cpp + src/model/diff/local/synchronization_finish.cpp + src/model/diff/local/synchronization_start.cpp src/model/diff/modify/add_blocks.cpp src/model/diff/modify/add_remote_folder_infos.cpp src/model/diff/modify/add_ignored_device.cpp diff --git a/src/fs/scan_actor.cpp b/src/fs/scan_actor.cpp index 251a5dfd..a94243c6 100644 --- a/src/fs/scan_actor.cpp +++ b/src/fs/scan_actor.cpp @@ -8,6 +8,8 @@ #include "model/diff/local/update.h" #include "model/diff/local/scan_finish.h" #include "model/diff/local/scan_start.h" +#include "model/diff/local/synchronization_start.h" +#include "model/diff/local/synchronization_finish.h" #include "net/names.h" #include "utils.h" #include @@ -84,65 +86,74 @@ void scan_actor_t::on_scan(message::scan_progress_t &message) noexcept { auto folder_id = task->get_folder_id(); LOG_TRACE(log, "on_scan, folder = {}", folder_id); - auto r = task->advance(); + auto folder = cluster->get_folders().by_id(folder_id); + bool stop_processing = false; bool completed = false; - std::visit( - [&](auto &&r) { - using T = std::decay_t; - if constexpr (std::is_same_v) { - stop_processing = !r; - if (stop_processing) { - completed = true; - } - } else if constexpr (std::is_same_v) { - send(coordinator, std::move(r)); - } else if constexpr (std::is_same_v) { - auto diff = model::diff::cluster_diff_ptr_t{}; - diff = new model::diff::local::file_availability_t(r.file); - send(coordinator, std::move(diff), this); - } else if constexpr (std::is_same_v) { - on_remove(*r.file); - } else if constexpr (std::is_same_v) { - auto &file = *r.file; - auto metadata = file.as_proto(true); - auto errs = initiate_hash(task, file.get_path(), metadata); - if (errs.empty()) { - stop_processing = true; - } else { - send(coordinator, std::move(errs)); - } - } else if constexpr (std::is_same_v) { - auto errs = initiate_hash(task, r.path, r.metadata); - if (errs.empty()) { - stop_processing = true; - } else { - send(coordinator, std::move(errs)); - } - } else if constexpr (std::is_same_v) { - auto &f = r.file; - assert(f->get_source()); - stop_processing = true; - send(address, task, std::move(f), f->get_source(), std::move(r.opened_file)); - } else if constexpr (std::is_same_v) { - auto &file = *r.file; - if (file.is_locked()) { + if (folder->is_synchronizing()) { + LOG_DEBUG(log, "folder = {} synchronization in progress, cancel scanning", folder_id); + stop_processing = true; + completed = true; + } else { + auto r = task->advance(); + std::visit( + [&](auto &&r) { + using T = std::decay_t; + if constexpr (std::is_same_v) { + stop_processing = !r; + if (stop_processing) { + completed = true; + } + } else if constexpr (std::is_same_v) { + send(coordinator, std::move(r)); + } else if constexpr (std::is_same_v) { auto diff = model::diff::cluster_diff_ptr_t{}; - diff = new model::diff::modify::lock_file_t(file, false); + diff = new model::diff::local::file_availability_t(r.file); send(coordinator, std::move(diff), this); + } else if constexpr (std::is_same_v) { + on_remove(*r.file); + } else if constexpr (std::is_same_v) { + auto &file = *r.file; + auto metadata = file.as_proto(true); + auto errs = initiate_hash(task, file.get_path(), metadata); + if (errs.empty()) { + stop_processing = true; + } else { + send(coordinator, std::move(errs)); + } + } else if constexpr (std::is_same_v) { + auto errs = initiate_hash(task, r.path, r.metadata); + if (errs.empty()) { + stop_processing = true; + } else { + send(coordinator, std::move(errs)); + } + } else if constexpr (std::is_same_v) { + auto &f = r.file; + assert(f->get_source()); + stop_processing = true; + send(address, task, std::move(f), f->get_source(), + std::move(r.opened_file)); + } else if constexpr (std::is_same_v) { + auto &file = *r.file; + if (file.is_locked()) { + auto diff = model::diff::cluster_diff_ptr_t{}; + diff = new model::diff::modify::lock_file_t(file, false); + send(coordinator, std::move(diff), this); + } + } else if constexpr (std::is_same_v) { + auto diff = model::diff::cluster_diff_ptr_t{}; + diff = new model::diff::modify::lock_file_t(*r.file, false); + send(coordinator, std::move(diff), this); + auto path = make_temporal(r.file->get_path()); + scan_errors_t errors{scan_error_t{path, r.ec}}; + send(coordinator, std::move(errors)); + } else { + static_assert(always_false_v, "non-exhaustive visitor!"); } - } else if constexpr (std::is_same_v) { - auto diff = model::diff::cluster_diff_ptr_t{}; - diff = new model::diff::modify::lock_file_t(*r.file, false); - send(coordinator, std::move(diff), this); - auto path = make_temporal(r.file->get_path()); - scan_errors_t errors{scan_error_t{path, r.ec}}; - send(coordinator, std::move(errors)); - } else { - static_assert(always_false_v, "non-exhaustive visitor!"); - } - }, - r); + }, + r); + } if (!stop_processing) { send(address, std::move(task)); diff --git a/src/fs/scan_scheduler.cpp b/src/fs/scan_scheduler.cpp index 1decb283..4dc2a1a6 100644 --- a/src/fs/scan_scheduler.cpp +++ b/src/fs/scan_scheduler.cpp @@ -7,6 +7,7 @@ #include "model/diff/local/scan_finish.h" #include "model/diff/local/scan_request.h" #include "model/diff/local/scan_start.h" +#include "model/diff/local/synchronization_finish.h" using namespace syncspirit::fs; @@ -72,6 +73,14 @@ auto scan_scheduler_t::operator()(const model::diff::local::scan_finish_t &diff, return diff.visit_next(*this, custom); } +auto scan_scheduler_t::operator()(const model::diff::local::synchronization_finish_t &diff, void *custom) noexcept + -> outcome::result { + if (!scan_in_progress) { + scan_next_or_schedule(); + } + return diff.visit_next(*this, custom); +} + void scan_scheduler_t::scan_next_or_schedule() noexcept { auto next = scan_next(); if (!next) { @@ -80,7 +89,7 @@ void scan_scheduler_t::scan_next_or_schedule() noexcept { bool do_start_timer = false; if (timer_id) { if (schedule_option->at > next.value().at) { - LOG_TRACE(log, "cancellling previous schedule"); + LOG_TRACE(log, "cancelling previous schedule"); cancel_timer(*timer_id); do_start_timer = true; } @@ -98,7 +107,8 @@ auto scan_scheduler_t::scan_next() noexcept -> schedule_option_t { while (!scan_queue.empty()) { auto folder_id = scan_queue.front(); scan_queue.pop_front(); - if (!cluster->get_folders().by_id(folder_id)) { + auto folder = cluster->get_folders().by_id(folder_id); + if (!folder || folder->is_synchronizing()) { continue; } for (auto it = scan_queue.begin(); it != scan_queue.end();) { @@ -116,18 +126,19 @@ auto scan_scheduler_t::scan_next() noexcept -> schedule_option_t { auto deadline = r::pt::ptime{}; auto now = r::pt::second_clock::local_time(); for (auto it : cluster->get_folders()) { - auto interval = it.item->get_rescan_interval(); - if (!interval || it.item->is_scanning()) { + auto &f = it.item; + auto interval = f->get_rescan_interval(); + if (!interval || f->is_scanning() || f->is_synchronizing()) { continue; } auto interval_s = r::pt::seconds{interval}; - auto prev_scan = it.item->get_scan_finish(); + auto prev_scan = f->get_scan_finish(); auto it_deadline = prev_scan.is_not_a_date_time() ? now : prev_scan + interval_s; auto eq = it_deadline == deadline; auto select_it = !folder || it_deadline < deadline || ((it_deadline == deadline) && (folder->get_rescan_interval() > interval)); if (select_it) { - folder = it.item; + folder = f; deadline = it_deadline; } } diff --git a/src/fs/scan_scheduler.h b/src/fs/scan_scheduler.h index ad8eb7fe..1cab626d 100644 --- a/src/fs/scan_scheduler.h +++ b/src/fs/scan_scheduler.h @@ -60,6 +60,8 @@ struct SYNCSPIRIT_API scan_scheduler_t : public r::actor_base_t, private model:: outcome::result operator()(const model::diff::modify::upsert_folder_t &, void *custom) noexcept override; outcome::result operator()(const model::diff::local::scan_request_t &, void *custom) noexcept override; outcome::result operator()(const model::diff::local::scan_finish_t &, void *custom) noexcept override; + outcome::result operator()(const model::diff::local::synchronization_finish_t &, + void *custom) noexcept override; model::cluster_ptr_t cluster; scan_queue_t scan_queue; diff --git a/src/model/diff/cluster_visitor.cpp b/src/model/diff/cluster_visitor.cpp index a4903100..f188be62 100644 --- a/src/model/diff/cluster_visitor.cpp +++ b/src/model/diff/cluster_visitor.cpp @@ -19,6 +19,8 @@ #include "local/scan_finish.h" #include "local/scan_request.h" #include "local/scan_start.h" +#include "local/synchronization_finish.h" +#include "local/synchronization_start.h" #include "local/update.h" #include "modify/add_blocks.h" #include "modify/add_ignored_device.h" @@ -125,6 +127,16 @@ auto cluster_visitor_t::operator()(const local::scan_start_t &diff, void *custom return diff.visit_next(*this, custom); } +auto cluster_visitor_t::operator()(const local::synchronization_start_t &diff, void *custom) noexcept + -> outcome::result { + return diff.visit_next(*this, custom); +} + +auto cluster_visitor_t::operator()(const local::synchronization_finish_t &diff, void *custom) noexcept + -> outcome::result { + return diff.visit_next(*this, custom); +} + auto cluster_visitor_t::operator()(const peer::cluster_update_t &diff, void *custom) noexcept -> outcome::result { return diff.visit_next(*this, custom); } diff --git a/src/model/diff/cluster_visitor.h b/src/model/diff/cluster_visitor.h index aeac0d0d..177b0f49 100644 --- a/src/model/diff/cluster_visitor.h +++ b/src/model/diff/cluster_visitor.h @@ -33,6 +33,8 @@ struct file_availability_t; struct scan_finish_t; struct scan_request_t; struct scan_start_t; +struct synchronization_start_t; +struct synchronization_finish_t; struct update_t; } // namespace local @@ -93,6 +95,8 @@ template <> struct SYNCSPIRIT_API generic_visitor_t operator()(const local::scan_finish_t &, void *custom) noexcept; virtual outcome::result operator()(const local::scan_request_t &, void *custom) noexcept; virtual outcome::result operator()(const local::scan_start_t &, void *custom) noexcept; + virtual outcome::result operator()(const local::synchronization_start_t &, void *custom) noexcept; + virtual outcome::result operator()(const local::synchronization_finish_t &, void *custom) noexcept; virtual outcome::result operator()(const peer::cluster_update_t &, void *custom) noexcept; virtual outcome::result operator()(const peer::update_folder_t &, void *custom) noexcept; diff --git a/src/model/diff/local/synchronization_finish.cpp b/src/model/diff/local/synchronization_finish.cpp new file mode 100644 index 00000000..a5b3099c --- /dev/null +++ b/src/model/diff/local/synchronization_finish.cpp @@ -0,0 +1,25 @@ +// SPDX-License-Identifier: GPL-3.0-or-later +// SPDX-FileCopyrightText: 2024 Ivan Baidakou + +#include "synchronization_finish.h" +#include "model/cluster.h" +#include "model/diff/cluster_visitor.h" + +using namespace syncspirit::model::diff::local; + +synchronization_finish_t::synchronization_finish_t(std::string_view folder_id_) : folder_id{std::move(folder_id_)} {} + +auto synchronization_finish_t::apply_impl(cluster_t &cluster) const noexcept -> outcome::result { + auto folder = cluster.get_folders().by_id(folder_id); + folder->set_synchronizing(false); + auto r = applicator_t::apply_sibling(cluster); + if (auto aug = folder->get_augmentation(); aug) { + aug->on_update(); + } + return r; +} + +auto synchronization_finish_t::visit(cluster_visitor_t &visitor, void *custom) const noexcept -> outcome::result { + LOG_TRACE(log, "visiting synchronization_finish_t"); + return visitor(*this, custom); +} diff --git a/src/model/diff/local/synchronization_finish.h b/src/model/diff/local/synchronization_finish.h new file mode 100644 index 00000000..0b9340ed --- /dev/null +++ b/src/model/diff/local/synchronization_finish.h @@ -0,0 +1,19 @@ +// SPDX-License-Identifier: GPL-3.0-or-later +// SPDX-FileCopyrightText: 2024 Ivan Baidakou + +#pragma once + +#include +#include "../cluster_diff.h" + +namespace syncspirit::model::diff::local { + +struct SYNCSPIRIT_API synchronization_finish_t final : cluster_diff_t { + synchronization_finish_t(std::string_view folder_id); + outcome::result apply_impl(cluster_t &) const noexcept override; + outcome::result visit(cluster_visitor_t &, void *) const noexcept override; + + std::string folder_id; +}; + +} // namespace syncspirit::model::diff::local diff --git a/src/model/diff/local/synchronization_start.cpp b/src/model/diff/local/synchronization_start.cpp new file mode 100644 index 00000000..41433548 --- /dev/null +++ b/src/model/diff/local/synchronization_start.cpp @@ -0,0 +1,25 @@ +// SPDX-License-Identifier: GPL-3.0-or-later +// SPDX-FileCopyrightText: 2024 Ivan Baidakou + +#include "synchronization_start.h" +#include "model/cluster.h" +#include "model/diff/cluster_visitor.h" + +using namespace syncspirit::model::diff::local; + +synchronization_start_t::synchronization_start_t(std::string_view folder_id_) : folder_id{std::move(folder_id_)} {} + +auto synchronization_start_t::apply_impl(cluster_t &cluster) const noexcept -> outcome::result { + auto folder = cluster.get_folders().by_id(folder_id); + folder->set_synchronizing(true); + auto r = applicator_t::apply_sibling(cluster); + if (auto aug = folder->get_augmentation(); aug) { + aug->on_update(); + } + return r; +} + +auto synchronization_start_t::visit(cluster_visitor_t &visitor, void *custom) const noexcept -> outcome::result { + LOG_TRACE(log, "visiting synchronization_start_t"); + return visitor(*this, custom); +} diff --git a/src/model/diff/local/synchronization_start.h b/src/model/diff/local/synchronization_start.h new file mode 100644 index 00000000..c6e6f07d --- /dev/null +++ b/src/model/diff/local/synchronization_start.h @@ -0,0 +1,19 @@ +// SPDX-License-Identifier: GPL-3.0-or-later +// SPDX-FileCopyrightText: 2024 Ivan Baidakou + +#pragma once + +#include +#include "../cluster_diff.h" + +namespace syncspirit::model::diff::local { + +struct SYNCSPIRIT_API synchronization_start_t final : cluster_diff_t { + synchronization_start_t(std::string_view folder_id); + outcome::result apply_impl(cluster_t &) const noexcept override; + outcome::result visit(cluster_visitor_t &, void *) const noexcept override; + + std::string folder_id; +}; + +} // namespace syncspirit::model::diff::local diff --git a/src/model/folder.cpp b/src/model/folder.cpp index 35808ac3..0b3e86a5 100644 --- a/src/model/folder.cpp +++ b/src/model/folder.cpp @@ -38,9 +38,9 @@ outcome::result folder_t::create(const uuid_t &uuid, const db::Fol return outcome::success(std::move(ptr)); } -folder_t::folder_t(std::string_view key_) noexcept { std::copy(key_.begin(), key_.end(), key); } +folder_t::folder_t(std::string_view key_) noexcept : synchronizing{false} { std::copy(key_.begin(), key_.end(), key); } -folder_t::folder_t(const uuid_t &uuid) noexcept { +folder_t::folder_t(const uuid_t &uuid) noexcept : synchronizing{false} { key[0] = prefix; std::copy(uuid.begin(), uuid.end(), key + 1); } @@ -114,6 +114,13 @@ const bool folder_t::is_scanning() const noexcept { return scan_start > scan_finish; } +const bool folder_t::is_synchronizing() const noexcept { return synchronizing; } + +void folder_t::set_synchronizing(bool value) noexcept { + assert(synchronizing != value); + synchronizing = value; +} + template <> SYNCSPIRIT_API std::string_view get_index<0>(const folder_ptr_t &item) noexcept { return item->get_key(); } template <> SYNCSPIRIT_API std::string_view get_index<1>(const folder_ptr_t &item) noexcept { return item->get_id(); } diff --git a/src/model/folder.h b/src/model/folder.h index 6698ca8c..8e0e529a 100644 --- a/src/model/folder.h +++ b/src/model/folder.h @@ -55,6 +55,8 @@ struct SYNCSPIRIT_API folder_t final : augmentable_t, folder_data_t { const pt::ptime &get_scan_finish() noexcept; void set_scan_finish(const pt::ptime &value) noexcept; const bool is_scanning() const noexcept; + const bool is_synchronizing() const noexcept; + void set_synchronizing(bool value) noexcept; using folder_data_t::get_path; using folder_data_t::set_path; @@ -76,6 +78,7 @@ struct SYNCSPIRIT_API folder_t final : augmentable_t, folder_data_t { folder_infos_map_t folder_infos; cluster_t *cluster = nullptr; char key[data_length]; + bool synchronizing; }; struct SYNCSPIRIT_API folders_map_t : generic_map_t { diff --git a/src/model/misc/file_iterator.cpp b/src/model/misc/file_iterator.cpp index d5dbcf32..0a75af2d 100644 --- a/src/model/misc/file_iterator.cpp +++ b/src/model/misc/file_iterator.cpp @@ -67,13 +67,6 @@ void file_iterator_t::reset() noexcept { } } -void file_iterator_t::renew(file_info_t &f) noexcept { - append(f); - if (!file) { - prepare(); - } -} - void file_iterator_t::prepare() noexcept { if (!incomplete.empty()) { file = incomplete.front(); diff --git a/src/model/misc/file_iterator.h b/src/model/misc/file_iterator.h index f4417264..ad6421fd 100644 --- a/src/model/misc/file_iterator.h +++ b/src/model/misc/file_iterator.h @@ -1,5 +1,5 @@ // SPDX-License-Identifier: GPL-3.0-or-later -// SPDX-FileCopyrightText: 2019-2022 Ivan Baidakou +// SPDX-FileCopyrightText: 2019-2024 Ivan Baidakou #pragma once @@ -23,7 +23,6 @@ struct SYNCSPIRIT_API file_iterator_t : arc_base_t { file_info_ptr_t next() noexcept; void reset() noexcept; - void renew(file_info_t &file) noexcept; private: using queue_t = std::deque; diff --git a/src/net/controller_actor.cpp b/src/net/controller_actor.cpp index 9b2caa20..cd0d3d78 100644 --- a/src/net/controller_actor.cpp +++ b/src/net/controller_actor.cpp @@ -5,6 +5,8 @@ #include "names.h" #include "constants.h" #include "model/diff/local/update.h" +#include "model/diff/local/synchronization_finish.h" +#include "model/diff/local/synchronization_start.h" #include "model/diff/modify/append_block.h" #include "model/diff/modify/block_ack.h" #include "model/diff/modify/block_rej.h" @@ -113,6 +115,12 @@ void controller_actor_t::shutdown_finish() noexcept { ++diff_counter; }; + for (auto &[folder, counter] : synchronizing_folders) { + if (counter) { + assign(new model::diff::local::synchronization_finish_t(folder->get_id())); + } + } + for (auto &file : locked_files) { if (!file->is_unlocking()) { LOG_TRACE(log, "going to unlock {} ({}); is_unlocking {}", file->get_full_name(), (void *)file.get(), @@ -208,7 +216,7 @@ void controller_actor_t::push_pending() noexcept { model::file_info_ptr_t controller_actor_t::next_file(bool reset) noexcept { if (reset) { - file_iterator = new model::file_iterator_t(*cluster, peer); + file_iterator.reset(new model::file_iterator_t(*cluster, peer)); } if (file_iterator && *file_iterator) { return file_iterator->next(); @@ -245,6 +253,12 @@ void controller_actor_t::on_pull_ready(message::pull_signal_t &) noexcept { if (reset_block) { auto diff = model::diff::cluster_diff_ptr_t{}; diff = new model::diff::modify::lock_file_t(*file, true); + auto folder = model::folder_ptr_t{file->get_folder_info()->get_folder()}; + auto &counter = ++synchronizing_folders[folder]; + if (counter == 1) { + auto raw = new model::diff::local::synchronization_start_t(folder->get_id()); + diff->sibling.reset(raw); + } send(coordinator, std::move(diff), this); } else { preprocess_block(block); @@ -384,6 +398,13 @@ auto controller_actor_t::operator()(const model::diff::modify::lock_file_t &diff auto it = locked_files.find(file); assert(it != locked_files.end()); locked_files.erase(it); + auto &counter = --synchronizing_folders[folder]; + if (counter == 0) { + auto diff = model::diff::cluster_diff_ptr_t{}; + diff = new model::diff::local::synchronization_finish_t(folder->get_id()); + send(coordinator, std::move(diff), this); + synchronizing_folders.erase(folder); + } } return diff.visit_next(*this, custom); } diff --git a/src/net/controller_actor.h b/src/net/controller_actor.h index 8d6b85c9..d994586f 100644 --- a/src/net/controller_actor.h +++ b/src/net/controller_actor.h @@ -17,6 +17,7 @@ #include "fs/messages.h" #include +#include #include namespace syncspirit { @@ -117,6 +118,7 @@ struct SYNCSPIRIT_API controller_actor_t : public r::actor_base_t, private model using unlink_requests_t = std::vector; using block_write_queue_t = std::deque; using dispose_callback_t = model::diff::modify::block_transaction_t::dispose_callback_t; + using synchronizing_folders_t = std::unordered_map; void on_termination(message::termination_signal_t &message) noexcept; void on_forward(message::forwarded_message_t &message) noexcept; @@ -187,6 +189,7 @@ struct SYNCSPIRIT_API controller_actor_t : public r::actor_base_t, private model int substate = substate_t::none; locked_files_t locked_files; locked_files_t locally_locked_files; + synchronizing_folders_t synchronizing_folders; block_write_queue_t block_write_queue; }; diff --git a/tests/050-file_iterator.cpp b/tests/050-file_iterator.cpp index 664f23c4..033fa110 100644 --- a/tests/050-file_iterator.cpp +++ b/tests/050-file_iterator.cpp @@ -125,19 +125,6 @@ TEST_CASE("file iterator", "[model]") { REQUIRE(!next()); } - SECTION("appending already visited file") { - auto f1 = next(true); - REQUIRE(f1); - CHECK(f1->get_name() == "a.txt"); - file_iterator->renew(*f1); - - auto f2 = next(); - REQUIRE(f2); - CHECK(f2->get_name() == "b.txt"); - - REQUIRE(!next()); - } - SECTION("one file is already exists on my side") { auto my_folder = folder_infos.by_device(*my_device); auto pr_file = proto::FileInfo(); diff --git a/tests/075-controller.cpp b/tests/075-controller.cpp index 7cd52310..354a8cc9 100644 --- a/tests/075-controller.cpp +++ b/tests/075-controller.cpp @@ -129,14 +129,7 @@ struct sample_peer_t : r::actor_base_t { messages.emplace_back(fwd_msg); } - void on_block_request(net::message::block_request_t &req) noexcept { - block_requests.push_front(&req); - ++blocks_requested; - log->debug("{}, requesting block # {}", identity, - block_requests.front()->payload.request_payload.block.block_index()); - if (block_responses.size()) { - log->debug("{}, top response block # {}", identity, block_responses.front().block_index); - } + void process_block_requests() noexcept { auto condition = [&]() -> bool { return block_requests.size() && block_responses.size() && block_requests.front()->payload.request_payload.block.block_index() == @@ -156,6 +149,17 @@ struct sample_peer_t : r::actor_base_t { } } + void on_block_request(net::message::block_request_t &req) noexcept { + block_requests.push_front(&req); + ++blocks_requested; + log->debug("{}, requesting block # {}", identity, + block_requests.front()->payload.request_payload.block.block_index()); + if (block_responses.size()) { + log->debug("{}, top response block # {}", identity, block_responses.front().block_index); + } + process_block_requests(); + } + void forward(net::payload::forwarded_message_t payload) noexcept { send(controller, std::move(payload)); } @@ -561,11 +565,17 @@ void test_downloading() { auto folder_my = folder_infos.by_device(*my_device); CHECK(folder_my->get_max_sequence() == 0ul); + CHECK(!folder_my->get_folder()->is_synchronizing()); peer_actor->forward(proto::message::Index(new proto::Index(index))); + sup->do_process(); + CHECK(folder_my->get_folder()->is_synchronizing()); + peer_actor->push_block("12345", 0); + peer_actor->process_block_requests(); sup->do_process(); + CHECK(!folder_my->get_folder()->is_synchronizing()); REQUIRE(folder_my); CHECK(folder_my->get_max_sequence() == 1ul); REQUIRE(folder_my->get_file_infos().size() == 1); @@ -928,6 +938,8 @@ void test_downloading_errors() { CHECK(!lf->is_locally_locked()); CHECK(!lf->is_locked()); + CHECK(!folder_my->get_folder()->is_synchronizing()); + sup->do_process(); } }; diff --git a/tests/085-scan-scheduler.cpp b/tests/085-scan-scheduler.cpp index fe1a45a5..ca3d40ae 100644 --- a/tests/085-scan-scheduler.cpp +++ b/tests/085-scan-scheduler.cpp @@ -93,7 +93,7 @@ void test_1_folder() { CHECK(folder->is_scanning()); SECTION("scan start/finish") { - builder.scan_finish(folder_id).upsert_folder(db_folder).apply(*sup); + builder.scan_finish(folder_id).apply(*sup); REQUIRE(!folder->is_scanning()); REQUIRE(sup->timers.size() == 1); @@ -101,6 +101,18 @@ void test_1_folder() { sup->do_process(); REQUIRE(folder->is_scanning()); } + + SECTION("synchronization start/finish") { + builder.scan_finish(folder_id).synchronization_start(folder_id).scan_request(folder_id).apply(*sup); + REQUIRE(!folder->is_scanning()); + REQUIRE(sup->timers.size() == 0); + + builder.synchronization_finish(folder_id).apply(*sup); + REQUIRE(sup->timers.size() == 1); + sup->do_invoke_timer((*sup->timers.begin())->request_id); + sup->do_process(); + REQUIRE(folder->is_scanning()); + } } } }; diff --git a/tests/086-scan_actor.cpp b/tests/086-scan_actor.cpp index 15020a8e..135b13c6 100644 --- a/tests/086-scan_actor.cpp +++ b/tests/086-scan_actor.cpp @@ -589,10 +589,28 @@ void test_remove_file() { F().run(); }; +void test_synchronization() { + struct F : fixture_t { + void main() noexcept override { + sys::error_code ec; + auto &blocks = cluster->get_blocks(); + + auto file_path = root_path / "file.ext"; + write_file(file_path, "12345"); + builder->scan_start(folder->get_id()).synchronization_start(folder->get_id()).apply(*sup); + REQUIRE(!folder->get_scan_finish().is_not_a_date_time()); + REQUIRE(!folder->is_scanning()); + REQUIRE(files->size() == 0); + } + }; + F().run(); +}; + int _init() { REGISTER_TEST_CASE(test_meta_changes, "test_meta_changes", "[fs]"); REGISTER_TEST_CASE(test_new_files, "test_new_files", "[fs]"); REGISTER_TEST_CASE(test_remove_file, "test_remove_file", "[fs]"); + REGISTER_TEST_CASE(test_synchronization, "test_synchronization", "[fs]"); return 1; } diff --git a/tests/diff-builder.cpp b/tests/diff-builder.cpp index 4877eda3..05d7754b 100644 --- a/tests/diff-builder.cpp +++ b/tests/diff-builder.cpp @@ -7,6 +7,8 @@ #include "model/diff/local/scan_finish.h" #include "model/diff/local/scan_request.h" #include "model/diff/local/scan_start.h" +#include "model/diff/local/synchronization_finish.h" +#include "model/diff/local/synchronization_start.h" #include "model/diff/modify/add_ignored_device.h" #include "model/diff/modify/add_pending_device.h" #include "model/diff/modify/append_block.h" @@ -254,6 +256,14 @@ diff_builder_t &diff_builder_t::scan_request(std::string_view id) noexcept { return assign(new model::diff::local::scan_request_t(std::string(id))); } +diff_builder_t &diff_builder_t::synchronization_start(std::string_view id) noexcept { + return assign(new model::diff::local::synchronization_start_t(std::string(id))); +} + +diff_builder_t &diff_builder_t::synchronization_finish(std::string_view id) noexcept { + return assign(new model::diff::local::synchronization_finish_t(std::string(id))); +} + template static void generic_assign(Holder *holder, Diff *diff) noexcept { if (!(*holder)) { holder->reset(diff); diff --git a/tests/diff-builder.h b/tests/diff-builder.h index dc21a721..2dee7d6c 100644 --- a/tests/diff-builder.h +++ b/tests/diff-builder.h @@ -85,6 +85,8 @@ struct SYNCSPIRIT_TEST_API diff_builder_t { diff_builder_t &scan_start(std::string_view id, const r::pt::ptime & = {}) noexcept; diff_builder_t &scan_finish(std::string_view id, const r::pt::ptime & = {}) noexcept; diff_builder_t &scan_request(std::string_view id) noexcept; + diff_builder_t &synchronization_start(std::string_view id) noexcept; + diff_builder_t &synchronization_finish(std::string_view id) noexcept; model::sequencer_t &get_sequencer() noexcept;