// SPDX-License-Identifier: GPL-3.0-or-later // SPDX-FileCopyrightText: 2025 Ivan Baidakou #include "test_peer.h" #include "net/names.h" #include "model/diff/contact/peer_state.h" #include "utils/log.h" #include "proto/bep_support.h" #include "proto/proto-helpers-bep.h" #include "model/messages.h" using namespace syncspirit; using namespace syncspirit::test; test_peer_t::test_peer_t(config_t &config) : r::actor_base_t{config}, auto_share(config.auto_share), coordinator(config.coordinator), peer_device{config.peer_device}, cluster{config.cluster}, url(config.url), peer_state{model::device_state_t::make_offline()} { auto id = fmt::format("test.peer.{}", url.c_str()); log = utils::get_logger(id); assert(cluster); assert(peer_device); assert(!url.empty()); } void test_peer_t::configure(r::plugin::plugin_base_t &plugin) noexcept { r::actor_base_t::configure(plugin); plugin.with_casted([&](auto &p) { auto id = fmt::format("test.peer.{}", url.c_str()); p.set_identity(id, false); }); plugin.with_casted([&](auto &p) { p.discover_name(net::names::coordinator, coordinator, false).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(&test_peer_t::on_controller_up, coordinator); plugin->subscribe_actor(&test_peer_t::on_controller_predown, coordinator); } }); }); plugin.with_casted([&](auto &p) { p.subscribe_actor(&test_peer_t::on_transfer); }); } void test_peer_t::on_start() noexcept { r::actor_base_t::on_start(); peer_state = peer_state.connecting().connected().online(url); auto diff = model::diff::contact::peer_state_t::create(*cluster, peer_device->device_id().get_sha256(), get_address(), peer_state); assert(diff); auto sup_addr = supervisor->get_address(); send(sup_addr, std::move(diff), nullptr); } void test_peer_t::shutdown_start() noexcept { LOG_TRACE(log, "shutdown_start"); if (controller) { send(controller, address, shutdown_reason); } if (shutdown_start_callback) { shutdown_start_callback(); } r::actor_base_t::shutdown_start(); } void test_peer_t::shutdown_finish() noexcept { r::actor_base_t::shutdown_finish(); LOG_TRACE(log, "shutdown_finish, blocks requested = {}", blocks_requested); if (controller) { send(controller, address, shutdown_reason); } auto sha256 = peer_device->device_id().get_sha256(); auto diff = model::diff::contact::peer_state_t::create(*cluster, sha256, address, peer_state.offline()); if (diff) { send(coordinator, std::move(diff)); } } void test_peer_t::on_controller_up(net::message::controller_up_t &msg) noexcept { auto &p = msg.payload; if (p.peer == peer_device->device_id() && p.url->c_str() == url) { LOG_TRACE(log, "on_controller_up"); controller = msg.payload.controller; reading = true; } } void test_peer_t::on_controller_predown(net::message::controller_predown_t &msg) noexcept { auto for_me = msg.payload.peer == address; LOG_TRACE(log, "on_controller_predown, for_me = {}", (for_me ? "yes" : "no")); if (for_me && !shutdown_reason) { auto &ee = msg.payload.ee; auto reason = ee->message(); LOG_TRACE(log, "on_termination: {}", reason); do_shutdown(ee); } } void test_peer_t::on_transfer(net::message::transfer_data_t &message) noexcept { auto &data = message.payload.data; auto buff = utils::bytes_view_t(data); auto messages = net::payload::forwarded_messages_t{}; while (buff.size()) { auto result = proto::parse_bep(buff); auto &value = result.value(); auto sz = value.consumed; buff = buff.subspan(sz); auto orig = std::move(value).message; auto type = proto::MessageType::UNKNOWN; auto variant = net::payload::forwarded_message_t(); bool has_variant = false; std::visit( [&](auto &msg) { using T = std::decay_t; using V = net::payload::forwarded_message_t; type = proto::message::get_bep_type(); if constexpr (std::is_same_v) { in_requests_copy.push_back(msg); in_requests.push_back(std::move(msg)); ++blocks_requested; process_block_requests(); } else if constexpr (std::is_same_v) { in_responses.push_back(std::move(msg)); } else if constexpr (std::is_constructible_v) { variant = std::move(msg); has_variant = true; } }, orig); LOG_TRACE(log, "on_transfer, bytes = {}, type = {}", sz, (int)type); if (has_variant) { messages.emplace_back(variant); bep_messages.emplace_back(variant); } for (auto &msg : bep_messages) { if (auto m = std::get_if(&msg); m) { auto folder = proto::get_folder(*m); allowed_index_updates.emplace(std::move(folder)); } if (auto m = std::get_if(&msg); m) { auto folder = std::string(proto::get_folder(*m)); if ((allowed_index_updates.count(folder) == 0) && !auto_share) { LOG_WARN(log, "IndexUpdate w/o previously recevied index"); std::abort(); } } } } if (messages.size()) { auto fwd_msg = new net::message::forwarded_messages_t(address, std::move(messages)); raw_messages.emplace_back(fwd_msg); } } void test_peer_t::process_block_requests() noexcept { auto condition = [&]() -> bool { if (in_requests.size() && out_responses.size()) { auto &req = in_requests.front(); auto &res = out_responses.front(); auto req_id = proto::get_id(req); auto res_id = proto::get_id(res); if (req_id == res_id) { return true; } } return false; }; auto replies = net::payload::forwarded_messages_t(); while (condition()) { auto &req = in_requests.front(); auto &res = out_responses.front(); auto res_id = proto::get_id(res); auto code = proto::get_code(res); log->debug("request & responce match by id '{}', replying..., code = {}", res_id, (int)code); replies.emplace_back(std::move(res)); in_requests.pop_front(); out_responses.pop_front(); } if (replies.size()) { send(controller, std::move(replies)); } } void test_peer_t::process_any_block_requests() noexcept { auto replies = net::payload::forwarded_messages_t(); for (auto it = in_requests.begin(); it != in_requests.end();) { auto &req = *it; auto req_id = proto::get_id(req); auto found = false; for (auto j = out_responses.begin(); j != out_responses.end(); ++j) { auto &res = *j; auto res_id = proto::get_id(res); if (req_id == res_id) { replies.emplace_back(std::move(res)); found = true; log->debug("request & responce match by id '{}', replying...", res_id); out_responses.erase(j); it = in_requests.erase(it); break; } } if (!found) { ++it; } } if (replies.size()) { send(controller, std::move(replies)); } } void test_peer_t::push_response(proto::ErrorCode code, std::int32_t request_id) noexcept { assert((int)code); proto::Response res; proto::set_id(res, request_id); proto::set_code(res, code); out_responses.emplace_back(std::move(res)); } void test_peer_t::push_response(utils::bytes_view_t data, std::int32_t request_id) noexcept { proto::Response res; proto::set_id(res, request_id); proto::set_data(res, data); out_responses.emplace_back(std::move(res)); } void test_peer_t::forward(net::payload::forwarded_message_t payload) noexcept { auto msgs = net::payload::forwarded_messages_t{std::move(payload)}; send(controller, std::move(msgs)); }