use core::marker::PhantomData;
use heapless::LinearMap;
use log::{debug, error, info, warn};
use crate::CodificationError;
use crate::communication::connection::{InMsgQueue, OutMsgQueue};
use crate::communication::messages::{ReportedAlert, AlertMessage, ControlMessage, GenericPackage, GenericPackageBuilder, Messages, Package, PackageBuilder, MetricData, SimpleCodifier, UpdateClusterVec, ReportedMeasurement, EventMessage, SubscriptionData};
use crate::communication::router::{is_unknown_route_id, Router};
use crate::communication::service::CommunicationService;
use crate::coordination::{CoordinationService, Stopwatch, Timer};
use crate::membership::client::{MembershipClient, MembershipManager};
use crate::membership::metadata::{CommunicationMetadata, NodeMetadata};
use crate::nodes::{CommunicationSDK, CoordinationSDK, NodeRole, SerializationSDK, SystemClusterId, SystemNodeId};
use crate::nodes::utils::{connect_with_members, SystemCommunicationService};
use crate::properties::{EncodedMetadata, EncodedMetadataVec, init_buffer, SystemBufferVec, CLUSTER_NODE_COUNT};
use crate::rules::manager::{RulesEngine, SystemStatus};
use crate::rules::measurements::{Measurement, ClusterType};
use crate::rules::strategy::{Rule, RuleType};
pub trait EnvironmentModule {
fn get_measurement(&mut self) -> Option<Measurement>;
}
pub trait CommunicationServiceWrapper<CoordSDK: CoordinationSDK, CommsSDK: CommunicationSDK> {
type Service: CommunicationService<Package<SystemNodeId, CoordSDK::Message>>;
fn unwrap_comm_service(
coordinated_communication_service: Self::Service,
) -> SystemCommunicationService<CommsSDK>;
fn wrap_comm_service(
communication_service: SystemCommunicationService<CommsSDK>,
) -> Self::Service;
}
struct CoordinationHelper<
CoordSDK: CoordinationSDK,
CommsSDK: CommunicationSDK,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>>,
CommsWrapper: CommunicationServiceWrapper<CoordSDK, CommsSDK>,
> {
id: SystemNodeId,
phantom1: PhantomData<CoordSDK>,
phantom2: PhantomData<CommsSDK>,
phantom3: PhantomData<Membership>,
phantom4: PhantomData<CommsWrapper>,
}
impl<
CoordSDK: CoordinationSDK,
CommsSDK: CommunicationSDK,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>>,
CommsWrapper: CommunicationServiceWrapper<CoordSDK, CommsSDK>,
> CoordinationHelper<CoordSDK, CommsSDK, Membership, CommsWrapper>
{
pub fn new(
id: SystemNodeId,
) -> CoordinationHelper<CoordSDK, CommsSDK, Membership, CommsWrapper> {
Self {
id,
phantom1: Default::default(),
phantom2: Default::default(),
phantom3: Default::default(),
phantom4: Default::default(),
}
}
fn coordinate(
members: CoordSDK::Members,
communication_service: SystemCommunicationService<CommsSDK>,
coordination_callback: &mut dyn FnMut(&mut CommsWrapper::Service, CoordSDK::Members),
) -> SystemCommunicationService<CommsSDK> {
let mut coordinated_communication_service =
CommsWrapper::wrap_comm_service(communication_service);
coordination_callback(&mut coordinated_communication_service, members);
CommsWrapper::unwrap_comm_service(coordinated_communication_service)
}
fn members(membership_manager: &Membership) -> CoordSDK::Members {
let members = LinearMap::from_iter(
membership_manager
.members()
.iter()
.filter(|(id, metadata)| {
if let NodeRole::SENSOR(_) = metadata.node.role {
if metadata.coordination.is_none() {
error!("Coordination metadata not available for node #{}", **id);
return false;
}
return true;
};
false
})
.map(|(id, metadata)| (*id, metadata.clone().coordination.unwrap())),
);
members.into()
}
pub fn coordinate_state(
&self,
coordination_service: &mut CoordSDK::Service,
communication_service: SystemCommunicationService<CommsSDK>,
membership_manager: &mut Membership,
measurement: Measurement,
) -> SystemCommunicationService<CommsSDK> {
if coordination_service.leader().is_none() {
debug!("Measurement taken without a cluster leader. Skipping.");
return communication_service;
}
let mut callback = move |send: &mut CommsWrapper::Service, _: CoordSDK::Members| {
coordination_service.update_state(send, measurement);
debug!("State: {:?}", coordination_service.get_state());
};
let members = Self::members(membership_manager);
Self::coordinate(members, communication_service, &mut callback)
}
pub fn coordinate_message(
&self,
coordination_service: &mut CoordSDK::Service,
communication_service: SystemCommunicationService<CommsSDK>,
membership_manager: &mut Membership,
message: CoordSDK::Message,
peer_id: SystemNodeId,
) -> SystemCommunicationService<CommsSDK> {
let mut callback = |comm: &mut CommsWrapper::Service, members: CoordSDK::Members| {
let package = CoordSDK::PackageBuilder::default()
.from(peer_id)
.to(self.id)
.with_message(message.clone())
.build()
.ok();
coordination_service.process(comm, package, members);
};
let members = Self::members(membership_manager);
Self::coordinate(members, communication_service, &mut callback)
}
pub fn coordinate_leadership(
&self,
coordination_service: &mut CoordSDK::Service,
communication_service: SystemCommunicationService<CommsSDK>,
membership_manager: &mut Membership,
) -> SystemCommunicationService<CommsSDK> {
let mut callback = |comm: &mut CommsWrapper::Service, members: CoordSDK::Members| {
coordination_service.process(comm, None, members);
};
let members = Self::members(membership_manager);
Self::coordinate(members, communication_service, &mut callback)
}
pub fn coordinate_cluster(
&self,
coordination_service: &mut CoordSDK::Service,
communication_service: SystemCommunicationService<CommsSDK>,
cluster: UpdateClusterVec,
) -> SystemCommunicationService<CommsSDK> {
let mut callback = |comm: &mut CommsWrapper::Service, _: CoordSDK::Members| {
coordination_service.update_members(comm, cluster.clone());
};
Self::coordinate(LinearMap::new().into(), communication_service, &mut callback)
}
pub fn coordinate_rule(
&self,
coordination_service: &mut CoordSDK::Service,
communication_service: SystemCommunicationService<CommsSDK>,
new_rule: Rule,
) -> SystemCommunicationService<CommsSDK> {
let mut callback = |comm: &mut CommsWrapper::Service, _: CoordSDK::Members| {
coordination_service.update_rule(comm, new_rule.clone());
};
Self::coordinate(LinearMap::new().into(), communication_service, &mut callback)
}
}
pub struct SensorNode<
Clock1: Stopwatch,
Clock2: Timer,
Environment: EnvironmentModule,
CoordSDK: CoordinationSDK,
CommsSDK: CommunicationSDK,
SerdeSDK: SerializationSDK<CommsSDK::Address, CoordSDK::Message, CoordSDK::Metadata, CommsSDK::Metadata>,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>>,
WrappedComms: CommunicationServiceWrapper<CoordSDK, CommsSDK>,
> {
id: SystemNodeId,
cluster: SystemClusterId,
keep_alive: Clock1,
idle_timer: Clock2,
measurement_clock: Clock1,
metadata: NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>,
sensor_type: ClusterType,
rules_engine: RulesEngine,
message_codifier: SerdeSDK::GeneralMessageCodifier,
environment_module: Environment,
coordination_service: CoordSDK::Service,
membership_manager: Membership,
helper: CoordinationHelper<CoordSDK, CommsSDK, Membership, WrappedComms>,
}
impl<
Clock1: Stopwatch,
Clock2: Timer,
Environment: EnvironmentModule,
CoordSDK: CoordinationSDK,
CommsSDK: CommunicationSDK,
SerdeSDK: SerializationSDK<
CommsSDK::Address,
CoordSDK::Message,
CoordSDK::Metadata,
CommsSDK::Metadata,
>,
Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>>,
WrappedComms: CommunicationServiceWrapper<CoordSDK, CommsSDK>,
>
SensorNode<Clock1, Clock2, Environment, CoordSDK, CommsSDK, SerdeSDK, Membership, WrappedComms>
where
Membership: MembershipClient<Address=CommsSDK::Address>
{
#[allow(clippy::too_many_arguments)]
pub fn new(
id: SystemNodeId,
cluster: SystemClusterId,
keep_alive: Clock1,
idle_timer: Clock2,
measurement_clock: Clock1,
metadata: NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>,
sensor_type: ClusterType,
rule: RuleType,
threshold: i32,
environment_module: Environment,
coordination_service: CoordSDK::Service,
membership_manager: Membership,
) -> SensorNode<
Clock1,
Clock2,
Environment,
CoordSDK,
CommsSDK,
SerdeSDK,
Membership,
WrappedComms,
> {
if idle_timer.as_secs() > keep_alive.as_secs() {
error!("The idle-time ({}s) can't be higher than the keep-alive ({}s). Aborting.", idle_timer.as_secs(), keep_alive.as_secs());
panic!()
}
if idle_timer.as_secs() > measurement_clock.as_secs() {
error!("The idle-time ({}s) can't be higher than the loop-clock ({}s). Aborting.", idle_timer.as_secs(), measurement_clock.as_secs());
panic!()
}
let helper = CoordinationHelper::<CoordSDK, CommsSDK, Membership, WrappedComms>::new(id);
let message_codifier = SerdeSDK::GeneralMessageCodifier::default();
let rules_engine = RulesEngine::new(rule, threshold);
SensorNode {
id,
cluster,
keep_alive,
idle_timer,
measurement_clock,
metadata,
sensor_type,
rules_engine,
message_codifier,
environment_module,
coordination_service,
membership_manager,
helper,
}
}
fn encode_metadata(&self) -> EncodedMetadata {
let mut buffer = init_buffer();
let node = if let Ok(bytes) = SerdeSDK::SystemMetadataCodifier::default()
.encode(&self.metadata.node, &mut buffer)
{
EncodedMetadataVec::from_slice(&buffer[..bytes]).unwrap()
} else {
panic!("Couldn't serialize Node #{} system metadata.", self.id);
};
let mut buffer = init_buffer();
let coordination = if let Ok(bytes) = SerdeSDK::CoordinatedMetadataCodifier::default()
.encode(self.metadata.coordination.as_ref().unwrap(), &mut buffer)
{
EncodedMetadataVec::from_slice(&buffer[..bytes]).unwrap()
} else {
panic!("Couldn't serialize Node #{} system metadata.", self.id);
};
EncodedMetadata { node, coordination: Some(coordination) }
}
fn decode_metadata(
&self,
encoded_metadata: EncodedMetadata,
) -> Result<NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>, CodificationError> {
let node =
SerdeSDK::SystemMetadataCodifier::default().decode(encoded_metadata.node.as_slice())?;
let coordination = if let Some(encoded_metadata) = encoded_metadata.coordination {
Some(
SerdeSDK::CoordinatedMetadataCodifier::default()
.decode(encoded_metadata.as_slice())?,
)
} else {
None
};
Ok(NodeMetadata { node, coordination })
}
fn decode_package(
&mut self,
package: &GenericPackage,
) -> Option<(SystemNodeId, Messages<CoordSDK::Message>)> {
Some((
package.header.from,
self.message_codifier.decode(package.body.as_slice()).ok()?,
))
}
fn basic_event(&self, metric: MetricData) -> EventMessage {
EventMessage {
node_id: self.id,
cluster_id: self.cluster,
typed: self.sensor_type,
active_rule: self.rules_engine.current(),
extra: None,
metric,
}
}
fn select_actuator(&self) -> Option<SystemNodeId> {
self.membership_manager.select_member(|metadata| metadata.node.role == NodeRole::ACTUATOR)
}
fn select_translator(&self) -> Option<SystemNodeId> {
self.membership_manager.select_member(|metadata| metadata.node.role == NodeRole::TRANSLATOR)
}
fn send_message(
&self,
communication_service: &mut SystemCommunicationService<CommsSDK>,
peer_id: SystemNodeId,
message: Messages<CoordSDK::Message>,
) {
let mut buffer = init_buffer();
let codified = SerdeSDK::GeneralMessageCodifier::default()
.encode(&message, &mut buffer)
.unwrap();
if let Ok(vec) = SystemBufferVec::from_slice(&buffer[..codified]) {
let package = GenericPackageBuilder::default()
.to(peer_id)
.from(self.id)
.with_message(vec)
.build();
communication_service.push(package.unwrap());
}
}
fn evaluate_state(&mut self, communication_service: &mut SystemCommunicationService<CommsSDK>) {
let rule = self.rules_engine.current();
info!("Evaluating state using rule {} with threshold at {}.", rule.name, rule.threshold);
let state = self.coordination_service.get_state();
let system_status = self.rules_engine.evaluate(&state);
let actuator_message: Messages<CoordSDK::Message> = match system_status {
SystemStatus::DECREASE => match self.sensor_type {
ClusterType::TEMPERATURE => Messages::Actuator(AlertMessage::DecreaseTemperature),
ClusterType::HUMIDITY => Messages::Actuator(AlertMessage::DecreaseHumidity),
},
SystemStatus::INCREASE => match self.sensor_type {
ClusterType::TEMPERATURE => Messages::Actuator(AlertMessage::IncreaseTemperature),
ClusterType::HUMIDITY => Messages::Actuator(AlertMessage::IncreaseHumidity),
},
SystemStatus::OK => {
info!("System status is OK!");
return;
}
};
if let Some(actuator_id) = self.select_actuator() {
info!("Reporting alert to Actuator #{}.", actuator_id);
self.send_message(communication_service, actuator_id, actuator_message);
} else {
warn!("Couldn't find Actuator to send alert.");
debug!("Members: {:?}", self.membership_manager.members());
}
if let Some(translator_id) = self.select_translator() {
let alert = MetricData::Alert(ReportedAlert::new(system_status));
let event = self.basic_event(alert);
let translator_message = Messages::Reporter(event);
self.send_message(communication_service, translator_id, translator_message);
} else {
warn!("Couldn't find Translator to report alert.");
debug!("Members: {:?}", self.membership_manager.members());
}
}
fn store_metadata(
&mut self,
communication_service: &mut SystemCommunicationService<CommsSDK>,
previous_id: SystemNodeId,
metadata: NodeMetadata<CoordSDK::Metadata, CommsSDK::Metadata>,
) -> bool {
if !metadata.node.same_cluster(self.cluster) {
info!("Ignoring metadata from node #{} (role: {:?}) because it belongs to a different cluster/s: {}.", metadata.node.id, metadata.node.role, metadata.node.log_clusters());
debug!("Removing router entry from node #{} (temporary entry for node #{}).", previous_id, metadata.node.id);
communication_service.router.unregister(previous_id);
return false;
}
if self.id == metadata.node.id {
error!("Received metadata from a node with identical ID #{} (previously registerd as #{}). Aborting.", self.id, previous_id);
panic!();
}
info!("Updating metadata for node #{}: {}", metadata.node.id, metadata);
communication_service.update(metadata.node.id, previous_id, metadata.node.communication.address());
match metadata.node.role {
NodeRole::SENSOR(sensor_type) if sensor_type != self.sensor_type => {
warn!("Received metadata from node #{} of type {} in same cluster, ignoring.", metadata.node.id, sensor_type);
communication_service.router.unregister(metadata.node.id);
false
}
_ => {
debug!("Registering node #{} with type {:?} into Membership Manager.", metadata.node.id, metadata.node.role);
self.membership_manager.register_member(metadata.node.id, metadata);
true
}
}
}
fn is_leader(&self) -> bool {
self.coordination_service.leader() == Some(self.id)
}
fn handle_metadata_message(
&mut self,
communication_service: &mut SystemCommunicationService<CommsSDK>,
previous_id: SystemNodeId,
encoded_metadata: EncodedMetadata,
with_reply: bool,
) {
if let Ok(metadata) = self.decode_metadata(encoded_metadata) {
let should_reply = self.store_metadata(communication_service, previous_id, metadata.clone());
let response = if with_reply && should_reply {
info!("Replying to metadata request from node #{} with own node metadata.", metadata.node.id);
Messages::System(ControlMessage::Metadata(self.encode_metadata()))
} else if should_reply {
info!("Replying to metadata from node #{} with ACK message.", metadata.node.id);
Messages::System(ControlMessage::ACK)
} else {
debug!("Omitting reply to node #{}.", metadata.node.id);
return;
};
self.send_message(communication_service, metadata.node.id, response);
if metadata.node.role == NodeRole::TRANSLATOR {
let register = MetricData::Register(SubscriptionData::new(self.measurement_clock.as_secs()));
let event = self.basic_event(register);
let translator_message = Messages::Reporter(event);
info!("Registering to WebServer through Translator #{}.", metadata.node.id);
self.send_message(communication_service, metadata.node.id, translator_message);
}
}
}
fn handle_system_message(
&mut self,
communication_service: &mut SystemCommunicationService<CommsSDK>,
peer_id: SystemNodeId,
message: ControlMessage,
) {
let response: Messages<CoordSDK::Message> = match message {
ControlMessage::Start | ControlMessage::Healthcheck => {
Messages::System(ControlMessage::ACK)
}
ControlMessage::ACK | ControlMessage::Fail => {
debug!("Message '{:?}' received from node #{}", message, peer_id);
return;
}
ControlMessage::Debug => {
let state = self.coordination_service.get_state();
Messages::System(ControlMessage::DebugData(self.is_leader(), state, self.coordination_service.get_current_rule()))
}
_ => {
warn!("Unexpected message received: {:?}", message);
return;
}
};
self.send_message(communication_service, peer_id, response);
}
fn connect_with_members(
&mut self,
communication_service: &mut SystemCommunicationService<CommsSDK>,
) {
connect_with_members::<
CoordSDK::Message,
CommsSDK,
CoordSDK::Metadata,
Membership,
SerdeSDK::GeneralMessageCodifier,
CLUSTER_NODE_COUNT,
>(
self.encode_metadata(),
&mut self.membership_manager,
communication_service,
&mut self.message_codifier,
);
}
pub fn run(&mut self, mut communication_service: SystemCommunicationService<CommsSDK>) {
info!("Initializing Sensor #{}.", self.id);
self.connect_with_members(&mut communication_service);
let mut started = false;
loop {
if self.keep_alive.is_timeout() {
self.membership_manager.keep_alive();
self.keep_alive.restart();
}
let mut should_idle = false;
if let Some(package) = communication_service.pop() {
if let Some((sender_id, message)) = self.decode_package(&package) {
debug!("Message received from node #{}: {:?}", sender_id, message);
match message {
Messages::Coordination(msg) => {
if started {
if is_unknown_route_id(sender_id) {
info!("Message from an unknown node! Starting discovery workflow.");
communication_service.queue(sender_id, package);
self.connect_with_members(&mut communication_service);
} else {
communication_service = self.helper.coordinate_message(
&mut self.coordination_service,
communication_service,
&mut self.membership_manager,
msg,
sender_id,
);
}
} else {
warn!("Node #{} hasn't been started but received a Coordination message from #{}.", self.id, sender_id);
}
}
Messages::System(ControlMessage::ClusterChange(cluster)) => {
info!("Updating cluster with new members: {:?}", cluster);
communication_service = self.helper.coordinate_cluster(
&mut self.coordination_service,
communication_service,
cluster,
);
let response = Messages::System(ControlMessage::ACK);
self.send_message(&mut communication_service, sender_id, response);
}
Messages::System(ControlMessage::Exit) => {
info!("Closing Sensor #{}.", self.id);
break;
}
Messages::System(ControlMessage::Start) => {
info!("Starting Sensor #{}.", self.id);
started = true;
self.measurement_clock.restart();
self.send_message(
&mut communication_service,
sender_id,
Messages::System(ControlMessage::ACK),
);
}
Messages::System(ControlMessage::GetMetadata(bytes)) => {
self.handle_metadata_message(
&mut communication_service,
sender_id,
bytes,
true,
);
}
Messages::System(ControlMessage::Metadata(metadata)) => {
self.handle_metadata_message(
&mut communication_service,
sender_id,
metadata,
false,
);
}
Messages::System(ControlMessage::RuleChange(change)) => {
info!("Updating engine strategy to {:?} with threshold {}", change.name, change.threshold);
communication_service = self.helper.coordinate_rule(
&mut self.coordination_service,
communication_service,
change,
);
self.rules_engine.update_strategy(change);
self.send_message(& mut communication_service, sender_id, Messages::System(ControlMessage::ACK));
}
Messages::System(msg) => {
self.handle_system_message(&mut communication_service, sender_id, msg);
}
_ => error!("Unexpected message received!"),
}
} else {
error!("Couldn't decode package {:?}", package)
}
} else {
should_idle = true
}
if started {
if self.measurement_clock.is_timeout() {
self.measurement_clock.restart();
if let Some(measurement) = self.environment_module.get_measurement() {
info!("Measurement read from Environment: {:?}", measurement);
communication_service = self.helper.coordinate_state(
&mut self.coordination_service,
communication_service,
&mut self.membership_manager,
measurement,
);
if let Some(translator_id) = self.select_translator() {
let measurement = MetricData::Measurement(ReportedMeasurement::new(self.is_leader(), measurement));
let event = self.basic_event(measurement);
let translator_message = Messages::Reporter(event);
info!("Sending measurement to Translator #{}.", translator_id);
self.send_message(&mut communication_service, translator_id, translator_message);
} else {
warn!("Couldn't find Translator to report measurement.");
debug!("Members: {:?}", self.membership_manager.members());
}
}
if self.is_leader() {
info!("Evaluating state.");
if let Some(rule) = self.coordination_service.get_current_rule() {
self.rules_engine.update_strategy(rule);
}
self.evaluate_state(&mut communication_service);
}
}
communication_service = self.helper.coordinate_leadership(
&mut self.coordination_service,
communication_service,
&mut self.membership_manager,
);
}
if should_idle {
self.idle_timer.wait();
}
}
}
}