mirror of
https://gitee.com/mirrors_PX4/PX4-Autopilot.git
synced 2026-10-08 11:58:52 +08:00
Cluster manager implementation, no tests yet
This commit is contained in:
@@ -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<protocol::dynamic_node_id::server::Discovery, DiscoveryCallback> 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<protocol::dynamic_node_id::server::Discovery>& 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<uint8_t>(cluster_size_ / 2U + 1U); }
|
||||
};
|
||||
|
||||
|
||||
@@ -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<ClusterManager*>(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<uint8_t>(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<protocol::dynamic_node_id::server::Discovery>& 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<uint8_t>(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
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user