From 51cd8404b133fe34f09d1f5516389e73aba08a9f Mon Sep 17 00:00:00 2001 From: Pavel Kirienko Date: Mon, 4 May 2015 19:00:39 +0300 Subject: [PATCH] Cluster manager implementation, no tests yet --- .../dynamic_node_id_allocation_server.hpp | 38 ++- .../uc_dynamic_node_id_allocation_server.cpp | 285 ++++++++++++++++++ 2 files changed, 317 insertions(+), 6 deletions(-) diff --git a/libuavcan/include/uavcan/protocol/dynamic_node_id_allocation_server.hpp b/libuavcan/include/uavcan/protocol/dynamic_node_id_allocation_server.hpp index cdd9d6c1bb..e37d3a5c87 100644 --- a/libuavcan/include/uavcan/protocol/dynamic_node_id_allocation_server.hpp +++ b/libuavcan/include/uavcan/protocol/dynamic_node_id_allocation_server.hpp @@ -239,7 +239,7 @@ class ClusterManager : private TimerBase struct Server { - const NodeID node_id; + NodeID node_id; Log::Index next_index; Log::Index match_index; @@ -251,7 +251,7 @@ class ClusterManager : private TimerBase enum { MaxServers = protocol::dynamic_node_id::server::Discovery::FieldTypes::known_nodes::MaxSize }; - const IDynamicNodeIDStorageBackend& storage_; + IDynamicNodeIDStorageBackend& storage_; const Log& log_; Subscriber discovery_sub_; @@ -262,11 +262,23 @@ class ClusterManager : private TimerBase uint8_t cluster_size_; uint8_t num_known_servers_; + bool had_discovery_activity_; + + static IDynamicNodeIDStorageBackend::String getStorageKeyForClusterSize() { return "cluster_size"; } + + INode& getNode() { return discovery_sub_.getNode(); } + const INode& getNode() const { return discovery_sub_.getNode(); } + + Server* findServer(NodeID node_id); + const Server* findServer(NodeID node_id) const; + bool isKnownServer(NodeID node_id) const; + void addServer(NodeID node_id); + virtual void handleTimerEvent(const TimerEvent&); void handleDiscovery(const ReceivedDataStructure& msg); - void publishDiscovery() const; + void publishDiscovery(); public: enum { ClusterSizeUnknown = 0 }; @@ -276,7 +288,7 @@ public: * @param storage Needed to read the cluster size parameter from the storage * @param log Needed to initialize nextIndex[] values after elections */ - ClusterManager(INode& node, const IDynamicNodeIDStorageBackend& storage, const Log& log) + ClusterManager(INode& node, IDynamicNodeIDStorageBackend& storage, const Log& log) : TimerBase(node) , storage_(storage) , log_(log) @@ -284,6 +296,7 @@ public: , discovery_pub_(node) , cluster_size_(0) , num_known_servers_(0) + , had_discovery_activity_(false) { } /** @@ -291,7 +304,7 @@ public: * storage backend using key 'cluster_size'. * Returns negative error code. */ - int init(uint8_t cluster_size = ClusterSizeUnknown); + int init(uint8_t init_cluster_size = ClusterSizeUnknown); /** * An invalid node ID will be returned if there's no such server. @@ -317,8 +330,21 @@ public: */ void resetAllServerIndices(); + /** + * This method returns true if there was at least one Discovery message received since last call. + */ + bool hadDiscoveryActivity() + { + if (had_discovery_activity_) + { + had_discovery_activity_ = false; + return true; + } + return false; + } + uint8_t getNumKnownServers() const { return num_known_servers_; } - uint8_t getConfiguredClusterSize() const { return cluster_size_; } + uint8_t getClusterSize() const { return cluster_size_; } uint8_t getQuorumSize() const { return static_cast(cluster_size_ / 2U + 1U); } }; diff --git a/libuavcan/src/protocol/uc_dynamic_node_id_allocation_server.cpp b/libuavcan/src/protocol/uc_dynamic_node_id_allocation_server.cpp index 68be3636f4..a69ca1bfb7 100644 --- a/libuavcan/src/protocol/uc_dynamic_node_id_allocation_server.cpp +++ b/libuavcan/src/protocol/uc_dynamic_node_id_allocation_server.cpp @@ -530,6 +530,291 @@ int PersistentState::setVotedFor(const NodeID node_id) return 0; } +/* + * ClusterManager + */ +ClusterManager::Server* ClusterManager::findServer(NodeID node_id) +{ + for (uint8_t i = 0; i < num_known_servers_; i++) + { + UAVCAN_ASSERT(servers_[i].node_id.isUnicast()); + if (servers_[i].node_id == node_id) + { + return &servers_[i]; + } + } + return NULL; +} + +const ClusterManager::Server* ClusterManager::findServer(NodeID node_id) const +{ + return const_cast(this)->findServer(node_id); +} + +bool ClusterManager::isKnownServer(NodeID node_id) const +{ + if (node_id == getNode().getNodeID()) + { + return true; + } + for (uint8_t i = 0; i < num_known_servers_; i++) + { + UAVCAN_ASSERT(servers_[i].node_id.isUnicast()); + UAVCAN_ASSERT(servers_[i].node_id != getNode().getNodeID()); + if (servers_[i].node_id == node_id) + { + return true; + } + } + return false; +} + +void ClusterManager::addServer(NodeID node_id) +{ + UAVCAN_ASSERT((num_known_servers_ + 1) < (MaxServers - 2)); + if (!isKnownServer(node_id) && node_id.isUnicast()) + { + servers_[num_known_servers_].node_id = node_id; + num_known_servers_ = static_cast(num_known_servers_ + 1U); + } + else + { + UAVCAN_ASSERT(0); + } +} + +void ClusterManager::handleTimerEvent(const TimerEvent&) +{ + UAVCAN_ASSERT(num_known_servers_ < cluster_size_); + if (num_known_servers_ < (cluster_size_ - 1)) + { + publishDiscovery(); + } + else + { + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", "Cluster is fully discovered, no more broadcasts"); + stop(); + } +} + +void ClusterManager::handleDiscovery(const ReceivedDataStructure& msg) +{ + /* + * Validating cluster configuration + * If there's a case of misconfiguration, the message will be ignored. + */ + if (msg.configured_cluster_size != cluster_size_) + { + getNode().registerInternalFailure("Bad Raft cluster size"); + return; + } + + had_discovery_activity_ = true; + + /* + * Updating the set of known servers + */ + for (uint8_t i = 0; i < msg.known_nodes.size(); i++) + { + if (num_known_servers_ >= (cluster_size_ - 1)) + { + break; + } + + const NodeID node_id(msg.known_nodes[i]); + if (node_id.isUnicast() && !isKnownServer(node_id)) + { + addServer(node_id); + } + } + + /* + * Publishing a new Discovery request if the timer is stopped already and the publishing server needs to + * learn about more servers. + */ + if ((msg.configured_cluster_size > msg.known_nodes.size()) && !isRunning()) + { + publishDiscovery(); + } +} + +void ClusterManager::publishDiscovery() +{ + protocol::dynamic_node_id::server::Discovery msg; + + msg.configured_cluster_size = cluster_size_; + + for (uint8_t i = 0; i < num_known_servers_; i++) + { + UAVCAN_ASSERT(servers_[i].node_id.isUnicast()); + msg.known_nodes.push_back(servers_[i].node_id.get()); + } + + UAVCAN_ASSERT(msg.known_nodes.size() == num_known_servers_); + + msg.known_nodes.push_back(getNode().getNodeID().get()); + + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", "Broadcasting Discovery message; known nodes: %d of %d", + int(msg.known_nodes.size()), int(cluster_size_)); + + const int res = discovery_pub_.broadcast(msg); + if (res < 0) + { + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", "Discovery broadcst failed: %d", res); + getNode().registerInternalFailure("Raft discovery broadcast"); + } +} + +int ClusterManager::init(const uint8_t init_cluster_size) +{ + /* + * Figuring out the cluster size + */ + if (init_cluster_size == ClusterSizeUnknown) + { + // Reading from the storage + MarshallingStorageDecorator io(storage_); + uint32_t value = 0; + int res = io.get(getStorageKeyForClusterSize(), value); + if (res < 0) + { + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", + "Cluster size is neither configured nor stored in the storage"); + return res; + } + if ((value == 0) || (value > MaxServers)) + { + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", "Cluster size is invalid"); + return -ErrFailure; + } + cluster_size_ = static_cast(value); + } + else + { + if ((init_cluster_size == 0) || (init_cluster_size > MaxServers)) + { + return -ErrInvalidParam; + } + cluster_size_ = init_cluster_size; + + // Writing the storage + MarshallingStorageDecorator io(storage_); + uint32_t value = init_cluster_size; + int res = io.setAndGetBack(getStorageKeyForClusterSize(), value); + if ((res < 0) || (value != init_cluster_size)) + { + UAVCAN_TRACE("dynamic_node_id_server_impl::ClusterManager", "Failed to store cluster size"); + return -ErrFailure; + } + } + + UAVCAN_ASSERT(cluster_size_ > 0); + UAVCAN_ASSERT(cluster_size_ <= MaxServers); + + /* + * Initializing pub/sub and timer + */ + int res = discovery_pub_.init(); + if (res < 0) + { + return res; + } + + res = discovery_sub_.start(DiscoveryCallback(this, &ClusterManager::handleDiscovery)); + if (res < 0) + { + return res; + } + + startPeriodic(MonotonicDuration::fromMSec(protocol::dynamic_node_id::server::Discovery::BROADCASTING_INTERVAL_MS)); + + /* + * Misc + */ + resetAllServerIndices(); + return 0; +} + +NodeID ClusterManager::getRemoteServerNodeIDAtIndex(uint8_t index) const +{ + if (index < num_known_servers_) + { + return servers_[index].node_id; + } + return NodeID(); +} + +Log::Index ClusterManager::getServerNextIndex(NodeID server_node_id) const +{ + const Server* const s = findServer(server_node_id); + if (s != NULL) + { + return s->next_index; + } + UAVCAN_ASSERT(0); + return 0; +} + +void ClusterManager::incrementServerNextIndexBy(NodeID server_node_id, Log::Index increment) +{ + Server* const s = findServer(server_node_id); + if (s != NULL) + { + s->next_index = Log::Index(s->next_index + increment); + } + else + { + UAVCAN_ASSERT(0); + } +} + +void ClusterManager::decrementServerNextIndex(NodeID server_node_id) +{ + Server* const s = findServer(server_node_id); + if (s != NULL) + { + s->next_index--; + } + else + { + UAVCAN_ASSERT(0); + } +} + +Log::Index ClusterManager::getServerMatchIndex(NodeID server_node_id) const +{ + const Server* const s = findServer(server_node_id); + if (s != NULL) + { + return s->match_index; + } + UAVCAN_ASSERT(0); + return 0; +} + +void ClusterManager::setServerMatchIndex(NodeID server_node_id, Log::Index match_index) +{ + Server* const s = findServer(server_node_id); + if (s != NULL) + { + s->match_index = match_index; + } + else + { + UAVCAN_ASSERT(0); + } +} + +void ClusterManager::resetAllServerIndices() +{ + for (uint8_t i = 0; i < num_known_servers_; i++) + { + UAVCAN_ASSERT(servers_[i].node_id.isUnicast()); + servers_[i].next_index = Log::Index(log_.getLastIndex() + 1U); + servers_[i].match_index = 0; + } +} + } // dynamic_node_id_server_impl }