dora-node-api 1.0.0-rc.5

`dora` goal is to be a low latency, composable, and distributed data flow.
Documentation
use std::sync::Arc;

use crate::{DaemonCommunicationWrapper, daemon_connection::DaemonChannel};
use dora_core::{
    config::{DataId, NodeId},
    uhlc::HLC,
};
use dora_message::{
    DataflowId,
    daemon_to_node::{DaemonCommunication, DaemonReply},
    metadata::Metadata,
    node_to_daemon::{DaemonRequest, DataMessage, Timestamped},
};
use eyre::{Context, bail, eyre};

pub(crate) struct ControlChannel {
    channel: DaemonChannel,
    clock: Arc<HLC>,
}

impl ControlChannel {
    #[tracing::instrument(level = "trace", skip(clock))]
    pub(crate) fn init(
        dataflow_id: DataflowId,
        node_id: &NodeId,
        daemon_communication: &DaemonCommunicationWrapper,
        clock: Arc<HLC>,
    ) -> eyre::Result<Self> {
        let channel = match daemon_communication {
            DaemonCommunicationWrapper::Standard(daemon_communication) => {
                match daemon_communication {
                    DaemonCommunication::Tcp { socket_addr } => {
                        DaemonChannel::new_tcp(*socket_addr)
                            .wrap_err("failed to connect control channel")?
                    }
                    DaemonCommunication::Interactive => {
                        DaemonChannel::Interactive(Default::default())
                    }
                }
            }
            DaemonCommunicationWrapper::Testing { channel, .. } => {
                DaemonChannel::IntegrationTestChannel(channel.clone())
            }
        };

        Self::init_on_channel(dataflow_id, node_id, channel, clock)
    }

    #[tracing::instrument(skip(channel, clock), level = "trace")]
    pub fn init_on_channel(
        dataflow_id: DataflowId,
        node_id: &NodeId,
        mut channel: DaemonChannel,
        clock: Arc<HLC>,
    ) -> eyre::Result<Self> {
        channel.register(dataflow_id, node_id.clone(), clock.new_timestamp())?;

        Ok(Self { channel, clock })
    }

    /// Drop the underlying daemon channel so this control-channel sender clone
    /// is released. The testing daemon exits after OutputsDone under shutdown
    /// independently of remaining EventStream sender clones.
    pub(crate) fn close_channel(&mut self) {
        self.channel = DaemonChannel::Interactive(Default::default());
    }

    pub fn report_outputs_done(&mut self) -> eyre::Result<()> {
        let reply = self
            .channel
            .request(&Timestamped {
                inner: DaemonRequest::OutputsDone,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to report outputs done to dora-daemon")?;
        match reply {
            DaemonReply::Result(result) => result
                .map_err(|e| eyre!(e))
                .wrap_err("failed to report outputs done event to dora-daemon")?,
            other => bail!("unexpected outputs done reply: {other:?}"),
        }
        Ok(())
    }

    pub fn report_closed_outputs(&mut self, outputs: Vec<DataId>) -> eyre::Result<()> {
        let reply = self
            .channel
            .request(&Timestamped {
                inner: DaemonRequest::CloseOutputs(outputs),
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to report closed outputs to dora-daemon")?;
        match reply {
            DaemonReply::Result(result) => result
                .map_err(|e| eyre!(e))
                .wrap_err("failed to receive closed outputs reply from dora-daemon")?,
            other => bail!("unexpected closed outputs reply: {other:?}"),
        }
        Ok(())
    }

    pub fn send_message(
        &mut self,
        output_id: DataId,
        metadata: Metadata,
        data: Option<DataMessage>,
    ) -> eyre::Result<()> {
        let request = DaemonRequest::SendMessage {
            output_id,
            metadata,
            data,
        };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send SendMessage request to dora-daemon")?;
        match reply {
            DaemonReply::Empty => Ok(()),
            other => bail!("unexpected SendMessage reply: {other:?}"),
        }
    }

    pub fn report_output_sent(
        &mut self,
        output_id: DataId,
        metadata: Metadata,
    ) -> eyre::Result<()> {
        let request = DaemonRequest::OutputSent {
            output_id,
            metadata,
        };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send OutputSent request to dora-daemon")?;
        match reply {
            DaemonReply::Empty => Ok(()),
            other => bail!("unexpected OutputSent reply: {other:?}"),
        }
    }

    pub fn extension_store(
        &mut self,
        namespace: String,
        key: String,
        value: Vec<u8>,
    ) -> eyre::Result<()> {
        let request = DaemonRequest::ExtensionStore {
            namespace,
            key,
            value,
        };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send ExtensionStore request to dora-daemon")?;
        match reply {
            DaemonReply::Result(Ok(())) => Ok(()),
            DaemonReply::Result(Err(e)) => bail!("{e}"),
            other => bail!("unexpected ExtensionStore reply: {other:?}"),
        }
    }

    pub fn extension_load(
        &mut self,
        namespace: String,
        key: String,
        remove: bool,
    ) -> eyre::Result<Option<Vec<u8>>> {
        let request = DaemonRequest::ExtensionLoad {
            namespace,
            key,
            remove,
        };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send ExtensionLoad request to dora-daemon")?;
        match reply {
            DaemonReply::ExtensionValue { value } => Ok(value),
            DaemonReply::Result(Err(e)) => bail!("{e}"),
            other => bail!("unexpected ExtensionLoad reply: {other:?}"),
        }
    }

    pub fn extension_drop(&mut self, namespace: String, key: String) -> eyre::Result<()> {
        let request = DaemonRequest::ExtensionDrop { namespace, key };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send ExtensionDrop request to dora-daemon")?;
        match reply {
            DaemonReply::Result(Ok(())) => Ok(()),
            DaemonReply::Result(Err(e)) => bail!("{e}"),
            other => bail!("unexpected ExtensionDrop reply: {other:?}"),
        }
    }

    /// Send an opaque request to the extension registered under
    /// `namespace` on this node's daemon, and return its opaque reply.
    ///
    /// dora interprets neither side; see `docs/extensions.md`.
    pub fn extension_request(
        &mut self,
        namespace: String,
        payload: Vec<u8>,
    ) -> eyre::Result<Vec<u8>> {
        let request = DaemonRequest::ExtensionRequest { namespace, payload };
        let reply = self
            .channel
            .request(&Timestamped {
                inner: request,
                timestamp: self.clock.new_timestamp(),
            })
            .wrap_err("failed to send ExtensionRequest to dora-daemon")?;
        match reply {
            DaemonReply::ExtensionReply { payload } => Ok(payload),
            DaemonReply::Result(Err(e)) => bail!("{e}"),
            other => bail!("unexpected ExtensionRequest reply: {other:?}"),
        }
    }
}