Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 14 additions & 2 deletions docs/tutorials/plugin-system.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

"Outside hybrid discovery" hides that hybrid does the opposite: PluginLayer defaults to MergePolicy::ENRICHMENT, so the merge pipeline keeps the field values of the plugin loaded first and only fills gaps from later plugins, while routing goes to the plugin loaded last. State that here, since this paragraph justifies last-wins by the served copy.

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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,10 @@
#include <memory>
#include <nlohmann/json.hpp>
#include <optional>
#include <set>
#include <shared_mutex>
#include <string>
#include <tuple>
#include <unordered_map>
#include <vector>

Expand Down Expand Up @@ -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<std::string> & 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
Expand Down Expand Up @@ -306,6 +324,10 @@ class PluginManager : public LogProviderRegistry {
TransportRegistry * transport_registry_ = nullptr;
/// Entity ID -> plugin name mapping (populated from IntrospectionProvider results)
std::unordered_map<std::string, std::string> 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<std::tuple<std::string, std::string, std::string>> reported_ownership_conflicts_;
bool shutdown_called_ = false;
};

Expand Down
1 change: 1 addition & 0 deletions src/ros2_medkit_gateway/src/gateway_node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
24 changes: 22 additions & 2 deletions src/ros2_medkit_gateway/src/plugins/plugin_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
#include <dlfcn.h>
#include <httplib.h>

#include <algorithm>
#include <rclcpp/rclcpp.hpp>
#include <regex>
#include <unordered_set>
Expand Down Expand Up @@ -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(),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When an id already has an owner from the previous refresh, the first transfer seen is "previous owner -> first publisher in load order", the reverse of where the id ends up, and the sorted pair then suppresses the real direction for good: with plugin_b owning shared_area and plugin_a starting to publish it, the only line ever logged is "transferred from plugin_b to plugin_a" while plugin_b keeps the id and its copy is served. With four publishers the same mechanism adds a "plugin_d to plugin_a" line one refresh after the first three. Log where the outcome is known instead: collect each id's publishers during the refresh and let finish_ownership_refresh() emit one line per id naming the publishers and the final owner, deduplicated on the id plus its publisher set.

"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<std::shared_mutex> 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<std::shared_mutex> lock(plugins_mutex_);
for (auto it = entity_ownership_.begin(); it != entity_ownership_.end();) {
Expand Down
97 changes: 97 additions & 0 deletions src/ros2_medkit_gateway/test/log_capture.hpp
Original file line number Diff line number Diff line change
@@ -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 <algorithm>
#include <atomic>
#include <cstdarg>
#include <cstdio>
#include <iterator>
#include <mutex>
#include <string>
#include <vector>

#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 "<logger name>: <message>", that contain `needle`.
std::vector<std::string> matching(const std::string & needle) const {
std::lock_guard<std::mutex> lk(mutex_);
std::vector<std::string> 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<LogCapture *> & active() {
static std::atomic<LogCapture *> 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<std::mutex> lk(capture->mutex_);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Between active().load() and this lock nothing stops ~LogCapture() on the test thread from restoring the handler and destroying the object, so a WARN from another thread (the REST server thread or an executor callback in the notify test) locks a freed mutex. Guard the pointer and the push with one static mutex and take that same mutex in the destructor around active().store(nullptr), which also makes mutex_ redundant.

capture->lines_.push_back(std::string(name != nullptr ? name : "") + ": " + buf);
}

rcutils_logging_output_handler_t previous_;
mutable std::mutex mutex_;
std::vector<std::string> lines_;
};

} // namespace ros2_medkit_gateway::test
Loading
Loading