dcs2 0.1.0

An extensible distributed control system framework made in rust with no-std support.
Documentation
use log::{debug, error, info, warn};

use crate::communication::connection::{InMsgQueue, OutMsgQueue};
use crate::communication::messages::{
    AlertMessage, ControlMessage, GenericPackageBuilder, SimpleCodifier, Messages,
    NoCoordinatedMessage, PackageBuilder,
};
use crate::communication::router::Router;
use crate::coordination::{Stopwatch, Timer};
use crate::membership::client::{MembershipClient, MembershipManager};
use crate::membership::metadata::{CommunicationMetadata, NoCoordinatedMetadata, NodeMetadata, SerializableMetadata};
use crate::nodes::utils::{connect_with_members, SystemCommunicationService};
use crate::nodes::{CommunicationSDK, SerializationSDK, SystemClusterId, SystemNodeId};
use crate::properties::{init_buffer, EncodedMetadata, EncodedMetadataVec, SystemBufferVec, CLUSTER_NODE_COUNT};

pub trait ActuatorWorker {
    fn actuate(&mut self, alert: AlertMessage);
}

/// This is a mock implementation of an actuator. All it does in regards to acting is logging the
/// order itself. Besides that it also implements encoding/decoding, communication with the membership
/// service, and all other basic procedures for a node to exist.  
pub struct ActuatorNode<
    Clock1: Stopwatch,
    Clock2: Timer,
    Worker: ActuatorWorker,
    CommsSDK: CommunicationSDK,
    Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
    SerdeSDK: SerializationSDK<
        CommsSDK::Address,
        NoCoordinatedMessage,
        NoCoordinatedMetadata,
        CommsSDK::Metadata,
    >,
> {
    id: SystemNodeId,
    cluster: SystemClusterId,
    keep_alive: Clock1,
    idle_timer: Clock2,
    worker: Worker,
    metadata: NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>,
    membership_manager: Membership,
    communication_service: SystemCommunicationService<CommsSDK>,
    message_codifier: SerdeSDK::GeneralMessageCodifier,
}

impl<
    Clock1: Stopwatch,
    Clock2: Timer,
    Actuator: ActuatorWorker,
    Membership: MembershipManager<CLUSTER_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
    CommsSDK: CommunicationSDK,
    SerdeSDK: SerializationSDK<
        CommsSDK::Address,
        NoCoordinatedMessage,
        NoCoordinatedMetadata,
        CommsSDK::Metadata,
    >,
> ActuatorNode<Clock1, Clock2, Actuator, CommsSDK, Membership, SerdeSDK>
    where Membership: MembershipClient<Address=CommsSDK::Address>,
          CommsSDK::Metadata: SerializableMetadata
{
    pub fn new(
        id: SystemNodeId,
        cluster: SystemClusterId,
        keep_alive: Clock1,
        idle_timer: Clock2,
        worker: Actuator,
        metadata: NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>,
        membership_manager: Membership,
        communication_service: SystemCommunicationService<CommsSDK>,
        message_codifier: SerdeSDK::GeneralMessageCodifier,
    ) -> ActuatorNode<Clock1, Clock2, Actuator, CommsSDK, Membership, SerdeSDK> {
        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!()
        }

        ActuatorNode {
            id,
            cluster,
            idle_timer,
            keep_alive,
            worker,
            metadata,
            membership_manager,
            communication_service,
            message_codifier,
        }
    }

    pub fn connect_with_members(&mut self) {
        connect_with_members::<
            NoCoordinatedMessage,
            CommsSDK,
            NoCoordinatedMetadata,
            Membership,
            SerdeSDK::GeneralMessageCodifier,
            CLUSTER_NODE_COUNT,
        >(
            self.encode_metadata(),
            &mut self.membership_manager,
            &mut self.communication_service,
            &mut self.message_codifier,
        );
    }

    pub fn receive_message(&mut self) -> Option<(SystemNodeId, Messages<NoCoordinatedMessage>)> {
        self.communication_service.pop().and_then(|pkg| {
            Some((
                pkg.header.from,
                self.message_codifier.decode(pkg.body.as_slice()).ok()?,
            ))
        })
    }

    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);
        };

        EncodedMetadata { node, coordination: None }
    }

    fn store_metadata(
        &mut self,
        previous_id: SystemNodeId,
        encoded_metadata: EncodedMetadata,
    ) -> Option<SystemNodeId> {
        let metadata = SerdeSDK::SystemMetadataCodifier::default().decode(encoded_metadata.node.as_slice()).ok()?;

        if !metadata.same_cluster(self.cluster) {
            info!("Ignoring metadata from node #{} (role: {:?}) because it belongs belongs to a different cluster/s: {}.", metadata.id, metadata.role, metadata.log_clusters());
            debug!("Removing router entry from node #{} (temporary entry for node #{}).", previous_id, metadata.id);
            self.communication_service.router.unregister(previous_id);
            return None;
        }

        if self.id == metadata.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.id, metadata);
        self.communication_service.update(metadata.id, previous_id, metadata.communication.address());

        debug!("Registering node #{} with type {:?} into Membership Manager.", metadata.id, metadata.role);
        let member_metadata = NodeMetadata { node: metadata.clone(), coordination: None };
        self.membership_manager.register_member(metadata.id, member_metadata);
        Some(metadata.id)
    }

    pub fn send_metadata(&mut self, to: SystemNodeId) -> Option<()> {
        let metadata = self.encode_metadata();
        let message: Messages<NoCoordinatedMessage> = Messages::System(ControlMessage::Metadata(metadata));
        self.send_message(to, message)
    }

    pub fn send_message(&mut self, to: SystemNodeId, message: Messages<NoCoordinatedMessage>) -> Option<()> {
        let mut buffer = init_buffer();
        let bytes = self.message_codifier.encode(&message, &mut buffer).ok()?;
        let encoded = SystemBufferVec::from_slice(&buffer[..bytes]).ok()?;

        let package = GenericPackageBuilder::default()
            .to(to)
            .from(self.id)
            .with_message(encoded)
            .build();

        self.communication_service.push(package.ok()?);
        Some(())
    }

    pub fn run(&mut self) {
        self.connect_with_members();

        loop {
            if self.keep_alive.is_timeout() {
                self.membership_manager.keep_alive();
                self.keep_alive.restart();
            }

            if let Some((id, message)) = self.receive_message() {
                debug!("Message received from node #{}: {:?}", id, message);
                match message {
                    Messages::Actuator(alert) => {
                        info!("Actuating over new alert from node #{id}.");
                        self.worker.actuate(alert);
                    }
                    Messages::System(ControlMessage::Metadata(metadata)) => {
                        info!("Metadata from node #{id} received.");
                        if let Some(new_id) = self.store_metadata(id, metadata) {
                            info!("Replying with ACK to node #{new_id}.");
                            let message: Messages<NoCoordinatedMessage> = Messages::System(ControlMessage::ACK);
                            if self.send_message(new_id, message).is_none() {
                                error!("Couldn't send ACK to node #{new_id}.");
                            }
                        }
                    }
                    Messages::System(ControlMessage::GetMetadata(metadata)) => {
                        info!("Metadata request from node #{id} received.");
                        if let Some(new_id) = self.store_metadata(id, metadata) {
                            info!("Replying with metadata to node #{new_id}.");
                            if self.send_metadata(new_id).is_none() {
                                error!("Couldn't send metadata to node #{new_id}.");
                            }
                        }
                    }
                    Messages::System(ControlMessage::Exit) => {
                        info!("Closing Actuator #{}.", self.id);
                        break;
                    }
                    Messages::System(ControlMessage::ACK) => info!("ACK received from node #{id}."),
                    _ => warn!("Unexpected message received: {message:?}")
                }
            } else {
                self.idle_timer.wait();
            }
        }
    }
}