diff --git a/docs/tutorials/plugin-system.rst b/docs/tutorials/plugin-system.rst index 622b50e17..d1640f62a 100644 --- a/docs/tutorials/plugin-system.rst +++ b/docs/tutorials/plugin-system.rst @@ -876,8 +876,9 @@ Multiple plugins can be loaded simultaneously: - **DataProvider / OperationProvider / FaultProvider**: These use per-entity routing based on entity ownership. Entities created by a plugin's IntrospectionProvider are automatically routed to that same plugin's DataProvider, OperationProvider, and FaultProvider. Multiple - plugins can each serve different entities concurrently - there is no "first wins" conflict - because each plugin only handles requests for its own entities. + plugins can each serve different entities concurrently, because each plugin only handles + requests for the entities it owns. An entity ID that several plugins publish has one owner + (see `Entity Ownership`_). - **Custom routes**: All plugins can register endpoints (use unique path prefixes) Entity Ownership @@ -894,6 +895,17 @@ route requests to the correct plugin. - When a data, operation, or fault request arrives for an entity, the handler looks up the owning plugin and delegates to its corresponding provider. Entities not owned by any plugin fall through to the default gateway behavior. +- When several plugins publish the same entity ID, for example a parent area that each of + them creates by default, the plugin loaded last among them owns it. Each refresh registers + the plugins in load order and every registration takes the ID over, so during a refresh + the ID passes from plugin to plugin and ends with the last one. Outside hybrid discovery + the entity cache also keeps the copy from the plugin loaded last, so requests go to the + plugin whose copy of the entity is served. +- The gateway logs a warning the first time ownership of an ID passes between two plugins + and does not repeat it on later refreshes. Once a refresh ends with no plugin owning the ID, + the gateway forgets which pairs it reported, so a conflict that comes back is logged again. +- When the owner stops publishing the ID, another plugin that still publishes it owns it by + the end of the same refresh. This model allows multiple plugins to coexist without conflict - each plugin manages its own entities independently. diff --git a/src/ros2_medkit_gateway/include/ros2_medkit_gateway/core/plugins/plugin_manager.hpp b/src/ros2_medkit_gateway/include/ros2_medkit_gateway/core/plugins/plugin_manager.hpp index 95b5534ac..676c39f11 100644 --- a/src/ros2_medkit_gateway/include/ros2_medkit_gateway/core/plugins/plugin_manager.hpp +++ b/src/ros2_medkit_gateway/include/ros2_medkit_gateway/core/plugins/plugin_manager.hpp @@ -34,8 +34,10 @@ #include #include #include +#include #include #include +#include #include #include @@ -188,12 +190,28 @@ class PluginManager : public LogProviderRegistry { /// Called after IntrospectionProvider::introspect() returns new entities. /// Maps entity IDs to the plugin that created them, enabling per-entity /// provider routing in handlers. + /// + /// An ID another plugin owns is transferred to this plugin, so when several + /// plugins publish one ID the plugin registered last owns it. The refresh + /// registers plugins in load order. Outside hybrid discovery it also adds + /// their entities to the cache in that order, and the cache keeps the last + /// copy of a duplicated entity, so the owner is the plugin whose copy is + /// served. Each refresh transfers such an ID again, so a transfer is logged + /// once per ID and pair of plugins, until a refresh ends with the ID + /// unowned (see finish_ownership_refresh()). void register_entity_ownership(const std::string & plugin_name, const std::vector & entity_ids); /// Clear all entity ownership entries for a given plugin. /// Called before re-registering during refresh to remove stale entries. void clear_entity_ownership(const std::string & plugin_name); + /// Close an entity refresh: forget the logged ownership transfers of IDs + /// that no plugin owns any more. Called once after every plugin's ownership + /// was registered. An ID that several plugins keep publishing stays owned, + /// so its transfers stay logged once, and the record holds no more than the + /// conflicts over IDs that are currently owned. + void finish_ownership_refresh(); + /// Get DataProvider for a specific entity (if plugin-owned) /// @return Non-owning pointer, or nullptr if entity is not plugin-owned /// or owning plugin doesn't implement DataProvider @@ -306,6 +324,10 @@ class PluginManager : public LogProviderRegistry { TransportRegistry * transport_registry_ = nullptr; /// Entity ID -> plugin name mapping (populated from IntrospectionProvider results) std::unordered_map entity_ownership_; + /// Ownership transfers already logged, as (entity ID, plugin, plugin) with + /// the two plugin names sorted. finish_ownership_refresh() drops the + /// entries of IDs nobody owns. + std::set> reported_ownership_conflicts_; bool shutdown_called_ = false; }; diff --git a/src/ros2_medkit_gateway/src/gateway_node.cpp b/src/ros2_medkit_gateway/src/gateway_node.cpp index b17256644..467692274 100644 --- a/src/ros2_medkit_gateway/src/gateway_node.cpp +++ b/src/ros2_medkit_gateway/src/gateway_node.cpp @@ -2461,6 +2461,7 @@ void GatewayNode::refresh_cache() { plugin_mgr_->clear_entity_ownership(name); plugin_mgr_->register_entity_ownership(name, entity_ids); } + plugin_mgr_->finish_ownership_refresh(); } // Filter ROS 2 internal nodes (underscore prefix convention). diff --git a/src/ros2_medkit_gateway/src/plugins/plugin_manager.cpp b/src/ros2_medkit_gateway/src/plugins/plugin_manager.cpp index 47b8e7036..c389f6a47 100644 --- a/src/ros2_medkit_gateway/src/plugins/plugin_manager.cpp +++ b/src/ros2_medkit_gateway/src/plugins/plugin_manager.cpp @@ -17,6 +17,7 @@ #include #include +#include #include #include #include @@ -614,13 +615,32 @@ void PluginManager::register_entity_ownership(const std::string & plugin_name, for (const auto & eid : entity_ids) { auto it = entity_ownership_.find(eid); if (it != entity_ownership_.end() && it->second != plugin_name) { - RCLCPP_WARN(logger(), "Entity '%s' ownership transferred from plugin '%s' to '%s'", eid.c_str(), - it->second.c_str(), plugin_name.c_str()); + // Plugins that all publish this ID pass it along on every refresh, so + // log each pair of them once. + const auto & [first, second] = std::minmax(it->second, plugin_name); + if (reported_ownership_conflicts_.emplace(eid, first, second).second) { + RCLCPP_WARN(logger(), + "Entity '%s' ownership transferred from plugin '%s' to '%s'. When several plugins publish an " + "entity, the one registered last in a refresh owns it. Not logged again for these two plugins " + "while the entity has an owner", + eid.c_str(), it->second.c_str(), plugin_name.c_str()); + } } entity_ownership_[eid] = plugin_name; } } +void PluginManager::finish_ownership_refresh() { + std::unique_lock lock(plugins_mutex_); + for (auto it = reported_ownership_conflicts_.begin(); it != reported_ownership_conflicts_.end();) { + if (entity_ownership_.count(std::get<0>(*it)) == 0) { + it = reported_ownership_conflicts_.erase(it); + } else { + ++it; + } + } +} + void PluginManager::clear_entity_ownership(const std::string & plugin_name) { std::unique_lock lock(plugins_mutex_); for (auto it = entity_ownership_.begin(); it != entity_ownership_.end();) { diff --git a/src/ros2_medkit_gateway/test/log_capture.hpp b/src/ros2_medkit_gateway/test/log_capture.hpp new file mode 100644 index 000000000..fa65bf079 --- /dev/null +++ b/src/ros2_medkit_gateway/test/log_capture.hpp @@ -0,0 +1,97 @@ +// Copyright 2026 bburda +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#pragma once + +// LogCapture - records rcutils log output for tests that assert what was +// logged, or how often. + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "rcutils/logging.h" + +namespace ros2_medkit_gateway::test { + +/// Captures every log line while alive and puts the previous output handler +/// back on every exit path, so a failing case cannot swallow the output of +/// later ones or leave rclcpp's handler replaced. +/// +/// The handler is a plain C function pointer with no user-data slot, so the +/// live capture is reached through a static. The pointer is atomic because +/// the handler is process-global and other threads log too. +class LogCapture { + public: + LogCapture() : previous_(rcutils_logging_get_output_handler()) { + active().store(this); + rcutils_logging_set_output_handler(&LogCapture::handler); + } + ~LogCapture() { + rcutils_logging_set_output_handler(previous_); + active().store(nullptr); + } + LogCapture(const LogCapture &) = delete; + LogCapture & operator=(const LogCapture &) = delete; + LogCapture(LogCapture &&) = delete; + LogCapture & operator=(LogCapture &&) = delete; + + /// Captured lines, each as ": ", that contain `needle`. + std::vector matching(const std::string & needle) const { + std::lock_guard lk(mutex_); + std::vector out; + std::copy_if(lines_.begin(), lines_.end(), std::back_inserter(out), [&needle](const std::string & line) { + return line.find(needle) != std::string::npos; + }); + return out; + } + + private: + static std::atomic & active() { + static std::atomic current{nullptr}; + return current; + } + + static void handler(const rcutils_log_location_t * /*location*/, int /*severity*/, const char * name, + rcutils_time_point_value_t /*timestamp*/, const char * format, va_list * args) { + char buf[1024]; + va_list copy; + va_copy(copy, *args); + // The format string arrives through the handler signature, so there is no + // literal to check. GCC exempts va_list formatters from + // -Wformat-nonliteral, clang does not. Scoped to the single call. +#pragma GCC diagnostic push +#pragma GCC diagnostic ignored "-Wformat-nonliteral" + vsnprintf(buf, sizeof(buf), format, copy); +#pragma GCC diagnostic pop + va_end(copy); + LogCapture * capture = active().load(); + if (capture == nullptr) { + return; + } + std::lock_guard lk(capture->mutex_); + capture->lines_.push_back(std::string(name != nullptr ? name : "") + ": " + buf); + } + + rcutils_logging_output_handler_t previous_; + mutable std::mutex mutex_; + std::vector lines_; +}; + +} // namespace ros2_medkit_gateway::test diff --git a/src/ros2_medkit_gateway/test/test_plugin_entity_routing.cpp b/src/ros2_medkit_gateway/test/test_plugin_entity_routing.cpp index 260fad967..0e1641059 100644 --- a/src/ros2_medkit_gateway/test/test_plugin_entity_routing.cpp +++ b/src/ros2_medkit_gateway/test/test_plugin_entity_routing.cpp @@ -14,9 +14,18 @@ #include +#include +#include +#include +#include +#include +#include + +#include "log_capture.hpp" #include "ros2_medkit_gateway/core/plugins/plugin_manager.hpp" #include "ros2_medkit_gateway/core/providers/data_provider.hpp" #include "ros2_medkit_gateway/core/providers/fault_provider.hpp" +#include "ros2_medkit_gateway/core/providers/introspection_provider.hpp" #include "ros2_medkit_gateway/core/providers/operation_provider.hpp" #include "ros2_medkit_gateway/dto/data.hpp" #include "ros2_medkit_gateway/dto/faults.hpp" @@ -479,6 +488,248 @@ TEST(PluginEntityRouting, ClearEntityOwnership) { EXPECT_EQ(*mgr.get_entity_owner("ent3"), "plugin_b"); } +// ============================================================================= +// Entity Published By Several Plugins +// ============================================================================= + +namespace { + +using ros2_medkit_gateway::test::LogCapture; + +/// Publishes a configurable set of areas through introspect() and serves data +/// for every entity it owns, so a test can see which plugin a request for a +/// shared id reaches. +class MockAreaPublisher : public GatewayPlugin, public IntrospectionProvider, public DataProvider { + public: + explicit MockAreaPublisher(std::string name, std::vector ids = {}) + : name_(std::move(name)), area_ids(std::move(ids)) { + } + std::string name() const override { + return name_; + } + void configure(const json & /*config*/) override { + } + void shutdown() override { + } + + IntrospectionResult introspect(const IntrospectionInput & /*input*/) override { + IntrospectionResult result; + for (const auto & id : area_ids) { + Area area; + area.id = id; + area.name = id; + result.new_entities.areas.push_back(std::move(area)); + } + return result; + } + + tl::expected list_data(const std::string & /*entity_id*/) override { + return dto::DataListResult{json{{"items", json::array()}}}; + } + tl::expected read_data(const std::string & /*entity_id*/, + const std::string & /*resource*/) override { + return dto::DataValue{json{{"plugin", name_}}}; + } + tl::expected + write_data(const std::string & /*entity_id*/, const std::string & /*resource*/, const json & /*payload*/) override { + return dto::DataWriteResult{json{{"status", "ok"}}}; + } + + std::string name_; + std::vector area_ids; +}; + +/// Runs one entity refresh the way the gateway's cache refresh does: every +/// introspecting plugin, in load order, is introspected and then has its +/// ownership cleared and registered again with the ids it returned, and the +/// refresh is closed with finish_ownership_refresh(). Returns the owner of +/// `watched_id` as a request would see it after each plugin's step. +std::vector> refresh_entities(PluginManager & mgr, const std::string & watched_id) { + std::vector> owners_seen; + for (const auto & [name, provider] : mgr.get_named_introspection_providers()) { + const auto result = provider->introspect(IntrospectionInput{}); + std::vector ids; + ids.reserve(result.new_entities.areas.size()); + for (const auto & area : result.new_entities.areas) { + ids.push_back(area.id); + } + mgr.clear_entity_ownership(name); + mgr.register_entity_ownership(name, ids); + owners_seen.push_back(mgr.get_entity_owner(watched_id)); + } + mgr.finish_ownership_refresh(); + return owners_seen; +} + +MockAreaPublisher * add_publisher(PluginManager & mgr, const std::string & name, + std::vector ids = {"shared_area"}) { + auto plugin = std::make_unique(name, std::move(ids)); + auto * raw = plugin.get(); + mgr.add_plugin(std::move(plugin)); + return raw; +} + +} // namespace + +// Two plugins publish one area id, as with a parent area several plugins +// default to. The plugin loaded last owns it after every refresh, requests +// reach its providers, and the conflict is logged once, not per refresh. +TEST(PluginEntityRouting, SharedEntityOwnedByLastPublisherAcrossRefreshes) { + PluginManager mgr; + add_publisher(mgr, "plugin_a"); + auto * last = add_publisher(mgr, "plugin_b"); + + const LogCapture log; + for (int refresh = 1; refresh <= 5; ++refresh) { + refresh_entities(mgr, "shared_area"); + EXPECT_EQ(mgr.get_entity_owner("shared_area"), std::optional("plugin_b")) << "refresh " << refresh; + } + EXPECT_EQ(mgr.get_data_provider_for_entity("shared_area"), static_cast(last)); + + const auto conflict_lines = log.matching("'shared_area'"); + ASSERT_EQ(conflict_lines.size(), 1u) << "the conflict must be logged once, not on every refresh"; + EXPECT_NE(conflict_lines[0].find("'plugin_a'"), std::string::npos) << conflict_lines[0]; + EXPECT_NE(conflict_lines[0].find("'plugin_b'"), std::string::npos) << conflict_lines[0]; +} + +// Four plugins publish one area id. Ownership passes from plugin to plugin +// inside every refresh, and each pair of plugins it passes between is logged +// once. After that, refreshes log nothing. +TEST(PluginEntityRouting, SharedEntityAmongFourPublishersStopsLogging) { + PluginManager mgr; + for (const char * name : {"plugin_a", "plugin_b", "plugin_c", "plugin_d"}) { + add_publisher(mgr, name); + } + + const LogCapture log; + refresh_entities(mgr, "shared_area"); + refresh_entities(mgr, "shared_area"); + const auto first_two = log.matching("'shared_area'"); + for (int refresh = 3; refresh <= 6; ++refresh) { + refresh_entities(mgr, "shared_area"); + EXPECT_EQ(mgr.get_entity_owner("shared_area"), std::optional("plugin_d")) << "refresh " << refresh; + } + EXPECT_EQ(log.matching("'shared_area'").size(), first_two.size()) << "no new line once every pair was reported"; + + // Pairs in the order ownership moves: a-b, b-c, c-d in the first refresh, + // d-a in the second. + ASSERT_EQ(first_two.size(), 4u); + EXPECT_EQ(std::set(first_two.begin(), first_two.end()).size(), 4u); +} + +enum class StoppingPublisher { kLoadedLater, kLoadedEarlier }; + +class SharedEntityHandover : public ::testing::TestWithParam {}; + +// One of two publishers of an id stops publishing it: the owner, which is the +// one loaded later, or the one loaded earlier. Either way the remaining +// publisher owns the id at the end of that same refresh, and no step of the +// refresh leaves the id without an owner. +TEST_P(SharedEntityHandover, RemainingPublisherOwnsItInTheSameRefresh) { + PluginManager mgr; + auto * early = add_publisher(mgr, "plugin_a", {}); + auto * late = add_publisher(mgr, "plugin_b"); + + // The plugin loaded later owns the id first, then the earlier one starts + // publishing it too. + refresh_entities(mgr, "shared_area"); + early->area_ids = {"shared_area"}; + refresh_entities(mgr, "shared_area"); + ASSERT_EQ(mgr.get_entity_owner("shared_area"), std::optional("plugin_b")); + + const bool later_stops = GetParam() == StoppingPublisher::kLoadedLater; + MockAreaPublisher * stopping = later_stops ? late : early; + MockAreaPublisher * remaining = later_stops ? early : late; + + const LogCapture log; + stopping->area_ids.clear(); + const auto owners_seen = refresh_entities(mgr, "shared_area"); + for (size_t step = 0; step < owners_seen.size(); ++step) { + EXPECT_TRUE(owners_seen[step].has_value()) << "step " << step << " left the id without an owner"; + } + EXPECT_EQ(mgr.get_entity_owner("shared_area"), std::optional(remaining->name())); + EXPECT_EQ(mgr.get_data_provider_for_entity("shared_area"), static_cast(remaining)); + + refresh_entities(mgr, "shared_area"); + EXPECT_EQ(mgr.get_entity_owner("shared_area"), std::optional(remaining->name())); + EXPECT_TRUE(log.matching("'shared_area'").empty()) << "the conflict of these two plugins was already logged"; +} + +INSTANTIATE_TEST_SUITE_P(PluginEntityRouting, SharedEntityHandover, + ::testing::Values(StoppingPublisher::kLoadedLater, StoppingPublisher::kLoadedEarlier), + [](const ::testing::TestParamInfo & param_info) { + return param_info.param == StoppingPublisher::kLoadedLater + ? std::string("OwnerLoadedLaterStops") + : std::string("PluginLoadedEarlierStops"); + }); + +// A reported conflict is remembered while some plugin owns the id, so a +// steady conflict is not logged again. Once no plugin owns the id at the end +// of a refresh it is forgotten, which keeps the record set bounded by the +// owned ids, and the conflict is logged again if it comes back. +TEST(PluginEntityRouting, SharedEntityConflictForgottenOnceTheIdIsUnowned) { + PluginManager mgr; + auto * first = add_publisher(mgr, "plugin_a"); + auto * second = add_publisher(mgr, "plugin_b"); + + const LogCapture log; + const auto conflicts = [&log]() { + return log.matching("'shared_area'").size(); + }; + for (int refresh = 1; refresh <= 3; ++refresh) { + refresh_entities(mgr, "shared_area"); + } + EXPECT_EQ(conflicts(), 1u); + + // One publisher pauses and resumes. Some plugin owns the id throughout. + first->area_ids.clear(); + refresh_entities(mgr, "shared_area"); + first->area_ids = {"shared_area"}; + refresh_entities(mgr, "shared_area"); + refresh_entities(mgr, "shared_area"); + EXPECT_EQ(conflicts(), 1u) << "the id stayed owned, so the conflict is still the one already reported"; + + first->area_ids.clear(); + second->area_ids.clear(); + refresh_entities(mgr, "shared_area"); + ASSERT_FALSE(mgr.get_entity_owner("shared_area").has_value()); + + first->area_ids = {"shared_area"}; + second->area_ids = {"shared_area"}; + refresh_entities(mgr, "shared_area"); + refresh_entities(mgr, "shared_area"); + EXPECT_EQ(conflicts(), 2u) << "the conflict was forgotten with the id and is reported once more"; +} + +// Positive control: a single plugin keeps what it publishes, releases what it +// stops publishing and logs nothing about ownership. The second plugin at the +// end shows that this capture does see an ownership conflict. +TEST(PluginEntityRouting, SinglePublisherOwnershipUnchanged) { + PluginManager mgr; + auto * solo = add_publisher(mgr, "plugin_a", {"area_1", "area_2"}); + + const LogCapture log; + for (int refresh = 1; refresh <= 3; ++refresh) { + refresh_entities(mgr, "area_1"); + EXPECT_EQ(mgr.get_entity_owner("area_1"), std::optional("plugin_a")) << "refresh " << refresh; + EXPECT_EQ(mgr.get_entity_owner("area_2"), std::optional("plugin_a")) << "refresh " << refresh; + } + + solo->area_ids = {"area_1", "area_3"}; + refresh_entities(mgr, "area_1"); + EXPECT_EQ(mgr.get_entity_owner("area_1"), std::optional("plugin_a")); + EXPECT_FALSE(mgr.get_entity_owner("area_2").has_value()); + EXPECT_EQ(mgr.get_entity_owner("area_3"), std::optional("plugin_a")); + EXPECT_EQ(mgr.get_data_provider_for_entity("area_3"), static_cast(solo)); + EXPECT_TRUE(log.matching("plugin_manager: ").empty()); + + add_publisher(mgr, "plugin_b", {"area_1"}); + refresh_entities(mgr, "area_1"); + EXPECT_EQ(mgr.get_entity_owner("area_1"), std::optional("plugin_b")); + EXPECT_EQ(log.matching("plugin_manager: ").size(), 1u); + EXPECT_EQ(log.matching("'area_1'").size(), 1u); +} + // ============================================================================= // OperationProvider list+filter contract (used by handle_get_operation) // ============================================================================= diff --git a/src/ros2_medkit_gateway/test/test_plugin_notify_integration.cpp b/src/ros2_medkit_gateway/test/test_plugin_notify_integration.cpp index 04bbf2b79..7c23157a4 100644 --- a/src/ros2_medkit_gateway/test/test_plugin_notify_integration.cpp +++ b/src/ros2_medkit_gateway/test/test_plugin_notify_integration.cpp @@ -24,6 +24,9 @@ // 3. assert the new app is visible via the ManifestManager // 4. remove the fragment + notify again // 5. assert the app is gone +// It also checks which plugin owns, and which plugin's copy the cache serves +// for, an area id that two plugins publish, over the same notify-driven +// refresh passes. #include #include @@ -38,10 +41,15 @@ #include #include #include +#include #include #include +#include +#include "log_capture.hpp" #include "ros2_medkit_gateway/core/plugins/entity_change_scope.hpp" +#include "ros2_medkit_gateway/core/plugins/gateway_plugin.hpp" +#include "ros2_medkit_gateway/core/providers/introspection_provider.hpp" #include "ros2_medkit_gateway/discovery/discovery_manager.hpp" #include "ros2_medkit_gateway/discovery/manifest/manifest_manager.hpp" #include "ros2_medkit_gateway/gateway_node.hpp" @@ -130,9 +138,16 @@ class NotifyIntegrationTest : public ::testing::Test { manifest_path = work_dir / "manifest.yaml"; std::ofstream(manifest_path) << kBaseManifest; - // Start the GatewayNode with the base manifest + fragments_dir wired up. - // discovery.mode=hybrid to exercise the manifest-load path. runtime - // discovery is disabled so the only source of apps is manifest + fragments. + // discovery.mode=hybrid to exercise the manifest-load path. + start_node("hybrid"); + } + + /// (Re)start the GatewayNode with the base manifest + fragments_dir wired up + /// in the given discovery mode. + void start_node(const std::string & discovery_mode) { + node.reset(); + // Runtime discovery is disabled so the only source of apps is manifest + + // fragments. // Reserve a free loopback port per test instance - GatewayNode starts its // REST server unconditionally, so we cannot share :8080 with parallel // gtest suites. `server.enabled` is not a real parameter; override @@ -142,7 +157,7 @@ class NotifyIntegrationTest : public ::testing::Test { ASSERT_GT(server_port, 0); rclcpp::NodeOptions opts; opts.parameter_overrides({ - {"discovery.mode", "hybrid"}, + {"discovery.mode", discovery_mode}, {"discovery.manifest_path", manifest_path.string()}, {"discovery.manifest.enabled", true}, {"discovery.manifest_strict_validation", false}, @@ -183,6 +198,35 @@ class NotifyIntegrationTest : public ::testing::Test { } }; +/// Plugin whose introspect() publishes the area `shared_area`, described as +/// coming from this plugin, while `publishing` is set. +class SharedAreaPublisher : public ros2_medkit_gateway::GatewayPlugin, + public ros2_medkit_gateway::IntrospectionProvider { + public: + SharedAreaPublisher(std::string name, bool publish) : name_(std::move(name)), publishing(publish) { + } + std::string name() const override { + return name_; + } + void configure(const nlohmann::json & /*config*/) override { + } + ros2_medkit_gateway::IntrospectionResult + introspect(const ros2_medkit_gateway::IntrospectionInput & /*input*/) override { + ros2_medkit_gateway::IntrospectionResult result; + if (publishing) { + ros2_medkit_gateway::Area area; + area.id = "shared_area"; + area.name = "Shared area"; + area.description = "published by " + name_; + result.new_entities.areas.push_back(std::move(area)); + } + return result; + } + + std::string name_; + bool publishing; +}; + } // namespace TEST_F(NotifyIntegrationTest, FragmentAddedAfterStartupBecomesVisibleOnNotify) { @@ -315,3 +359,77 @@ TEST_F(NotifyIntegrationTest, NotifyWithoutAnyFragmentIsANoOp) { }); EXPECT_NE(it, comps.end()) << "base manifest entity lost after notify"; } + +TEST_F(NotifyIntegrationTest, SharedPluginAreaServedAndOwnedByTheSamePlugin) { + // Outside hybrid mode the refresh adds every plugin's copy of the area to + // the entity cache. Requests must go to the plugin whose copy the cache + // serves, through every change of who publishes it. + start_node("manifest_only"); + ASSERT_NE(node, nullptr); + node->stop_discovery_refresh_for_testing(); + auto * pm = node->get_plugin_manager(); + ASSERT_NE(pm, nullptr); + auto first = std::make_unique("plugin_a", false); + auto * first_raw = first.get(); + pm->add_plugin(std::move(first)); + auto second = std::make_unique("plugin_b", true); + auto * second_raw = second.get(); + pm->add_plugin(std::move(second)); + auto ctx = ros2_medkit_gateway::make_gateway_plugin_context(node.get(), node->get_fault_manager(), nullptr); + const auto refresh_and_expect = [&](const std::string & plugin, const std::string & when) { + ctx->notify_entities_changed(ros2_medkit_gateway::EntityChangeScope::full_refresh()); + EXPECT_EQ(pm->get_entity_owner("shared_area"), std::optional(plugin)) << when; + auto area = node->get_thread_safe_cache().get_area("shared_area"); + ASSERT_TRUE(area.has_value()) << when; + EXPECT_EQ(area->description, "published by " + plugin) << when; + }; + + refresh_and_expect("plugin_b", "only plugin_b publishes"); + first_raw->publishing = true; + for (int pass = 1; pass <= 3; ++pass) { + refresh_and_expect("plugin_b", "both publish, pass " + std::to_string(pass)); + } + second_raw->publishing = false; + refresh_and_expect("plugin_a", "the refresh in which plugin_b stops"); + second_raw->publishing = true; + refresh_and_expect("plugin_b", "plugin_b publishes again"); + refresh_and_expect("plugin_b", "plugin_b publishes again, next pass"); +} + +TEST_F(NotifyIntegrationTest, SharedPluginAreaConflictLoggedOnceWhileOwned) { + // Two plugins publish one area id on every refresh pass. The conflict is + // logged once, and once more only after a pass in which no plugin owned + // the id. + node->stop_discovery_refresh_for_testing(); + auto * pm = node->get_plugin_manager(); + ASSERT_NE(pm, nullptr); + auto first = std::make_unique("plugin_a", true); + auto * first_raw = first.get(); + pm->add_plugin(std::move(first)); + auto second = std::make_unique("plugin_b", true); + auto * second_raw = second.get(); + pm->add_plugin(std::move(second)); + auto ctx = ros2_medkit_gateway::make_gateway_plugin_context(node.get(), node->get_fault_manager(), nullptr); + const auto refresh = [&ctx]() { + ctx->notify_entities_changed(ros2_medkit_gateway::EntityChangeScope::full_refresh()); + }; + + const ros2_medkit_gateway::test::LogCapture log; + for (int pass = 1; pass <= 3; ++pass) { + refresh(); + } + EXPECT_EQ(pm->get_entity_owner("shared_area"), std::optional("plugin_b")); + EXPECT_EQ(log.matching("'shared_area'").size(), 1u) << "the conflict must be logged once, not on every pass"; + + first_raw->publishing = false; + second_raw->publishing = false; + refresh(); + ASSERT_FALSE(pm->get_entity_owner("shared_area").has_value()); + + first_raw->publishing = true; + second_raw->publishing = true; + refresh(); + refresh(); + EXPECT_EQ(log.matching("'shared_area'").size(), 2u) + << "the refresh forgets the conflict of an id nobody owns, so it is reported once more"; +}