dcs2 0.1.0

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

use crate::communication::connection::{InMsgQueue, OutMsgQueue};
use crate::communication::messages::{ControlMessage, GenericPackageBuilder, SimpleCodifier, Messages, NoCoordinatedMessage, PackageBuilder, EventMessage};
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, CLUSTERS_PER_TRANSLATOR, EncodedMetadata, SystemBufferVec, CLUSTER_NODE_COUNT};

const TRANSLATOR_NODE_COUNT: usize = CLUSTER_NODE_COUNT * CLUSTERS_PER_TRANSLATOR;

pub trait ServerModule {
    fn inbound(&mut self) -> Option<(SystemNodeId, Messages<NoCoordinatedMessage>)>;
    fn outbound(&mut self, message: EventMessage);
}

pub struct TranslatorNode<
    Clock1: Stopwatch,
    Clock2: Timer,
    CommsSDK: CommunicationSDK,
    SerdeSDK: SerializationSDK<
        CommsSDK::Address,
        NoCoordinatedMessage,
        NoCoordinatedMetadata,
        CommsSDK::Metadata,
    >,
    Membership: MembershipManager<TRANSLATOR_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
    ServerHandler: ServerModule,
> {
    id: SystemNodeId,
    clusters: Vec<SystemClusterId, CLUSTERS_PER_TRANSLATOR>,
    keep_alive: Clock1,
    idle_timer: Clock2,
    metadata: EncodedMetadata,
    node_codifier: SerdeSDK::SystemMetadataCodifier,
    message_codifier: SerdeSDK::GeneralMessageCodifier,
    membership_manager: Membership,
    communication_service: SystemCommunicationService<CommsSDK>,
    server_handler: ServerHandler,
}

impl<
    Clock1: Stopwatch,
    Clock2: Timer,
    CommsSDK: CommunicationSDK,
    SerdeSDK: SerializationSDK<
        CommsSDK::Address,
        NoCoordinatedMessage,
        NoCoordinatedMetadata,
        CommsSDK::Metadata,
    >,
    Membership: MembershipManager<TRANSLATOR_NODE_COUNT, NodeMetadata=NodeMetadata<NoCoordinatedMetadata, CommsSDK::Metadata>>,
    ServerHandler: ServerModule,
> TranslatorNode<Clock1, Clock2, CommsSDK, SerdeSDK, Membership, ServerHandler>
    where Membership: MembershipClient<Address=CommsSDK::Address>,
          CommsSDK::Metadata: SerializableMetadata,
{
    #[allow(clippy::too_many_arguments)]
    pub fn new(
        id: SystemNodeId,
        clusters: Vec<SystemClusterId, CLUSTERS_PER_TRANSLATOR>,
        keep_alive: Clock1,
        idle_timer: Clock2,
        metadata: EncodedMetadata,
        node_codifier: SerdeSDK::SystemMetadataCodifier,
        message_codifier: SerdeSDK::GeneralMessageCodifier,
        membership_manager: Membership,
        communication_service: SystemCommunicationService<CommsSDK>,
        server_handler: ServerHandler,
    ) -> TranslatorNode<Clock1, Clock2, CommsSDK, SerdeSDK, Membership, ServerHandler> {
        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!()
        }

        TranslatorNode {
            id,
            clusters,
            keep_alive,
            idle_timer,
            metadata,
            node_codifier,
            message_codifier,
            membership_manager,
            communication_service,
            server_handler,
        }
    }

    fn store_metadata(
        &mut self,
        previous_id: SystemNodeId,
        encoded_metadata: EncodedMetadata,
    ) -> Option<SystemNodeId>  {
        let metadata = self.node_codifier.decode(encoded_metadata.node.as_slice()).ok()?;
        if !self.clusters.iter().any(|cluster| metadata.same_cluster(*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)
    }

    fn send_message(&mut self, peer_id: SystemNodeId, message: Messages<NoCoordinatedMessage>) {
        let mut buffer = init_buffer();

        if let Ok(codified) = self.message_codifier.encode(&message, &mut buffer) {
            if let Ok(vec) = SystemBufferVec::from_slice(&buffer[..codified]) {
                let package = GenericPackageBuilder::default()
                    .to(peer_id)
                    .from(self.id)
                    .with_message(vec)
                    .build();

                self.communication_service.push(package.unwrap());
            }
        }
    }

    fn check_for_server_messages(&mut self) -> bool {
        return if let Some((id, message)) = self.server_handler.inbound() {
            info!("About to send a message to node #{}.", id);
            debug!("Message: {:?}", message);
            self.send_message(id, message);
            true
        } else {
            false
        }
    }

    fn check_for_nodes_messages(&mut self) -> bool {
        if let Some(package) = self.communication_service.pop() {
            if let Some(message) = self.message_codifier.decode(&package.body).ok() {
                let sender_id = package.header.from;
                match message {
                    Messages::System(ControlMessage::Metadata(encoded_metadata)) => {
                        info!("Metadata from node #{sender_id} received.");
                        if let Some(new_id) = self.store_metadata(sender_id, encoded_metadata) {
                            info!("Replying with ACK to node #{new_id}.");
                            let response = Messages::System(ControlMessage::ACK);
                            self.send_message(new_id, response);
                        }
                    }
                    Messages::System(ControlMessage::GetMetadata(encoded_metadata)) => {
                        info!("Metadata request from node #{sender_id} received.");
                        if let Some(new_id) = self.store_metadata(sender_id, encoded_metadata) {
                            info!("Replying with metadata to node #{new_id}.");
                            let response = Messages::System(ControlMessage::Metadata(self.metadata.clone()));
                            self.send_message(new_id, response);
                        }
                    }
                    Messages::Reporter(message) => {
                        info!("Received event from node #{sender_id}.");
                        self.server_handler.outbound(message);
                    }
                    Messages::System(ControlMessage::ACK) => info!("ACK received from node #{sender_id}."),
                    _ => warn!("Unexpected message received: {message:?}"),
                }
            } else {
                warn!("Couldn't decode package: {package:?}");
            }

            true
        } else {
            false
        }
    }

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

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

        info!("Starting to process messages.");
        loop {
            if self.keep_alive.is_timeout() {
                self.membership_manager.keep_alive();
                self.keep_alive.restart();
            }

            let received1 = self.check_for_nodes_messages();
            let received2 = self.check_for_server_messages();

            if !received1 && !received2 {
                self.idle_timer.wait();
            }
        }
    }
}