pub use crate::common::{DataMessage, LogLevel, LogMessage, SharedMemoryId, Timestamped};
use crate::{
DataflowId, current_crate_version,
id::{DataId, NodeId},
metadata::Metadata,
versions_compatible,
};
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub enum DaemonRequest {
Register(NodeRegisterRequest),
Subscribe,
SendMessage {
output_id: DataId,
metadata: Metadata,
data: Option<DataMessage>,
},
OutputSent {
output_id: DataId,
metadata: Metadata,
},
CloseOutputs(Vec<DataId>),
OutputsDone,
NextEvent,
EventStreamDropped,
NodeConfig {
node_id: NodeId,
},
RegisterPinnedMemory {
shared_memory_id: String,
metadata: Metadata,
},
ReadPinnedMemory {
shared_memory_id: String,
free: bool,
},
FreePinnedMemory {
shared_memory_id: String,
},
}
impl DaemonRequest {
pub fn expects_tcp_bincode_reply(&self) -> bool {
#[allow(clippy::match_like_matches_macro)]
match self {
DaemonRequest::SendMessage { .. }
| DaemonRequest::OutputSent { .. }
| DaemonRequest::NodeConfig { .. } => false,
DaemonRequest::Register(NodeRegisterRequest { .. })
| DaemonRequest::Subscribe
| DaemonRequest::CloseOutputs(_)
| DaemonRequest::OutputsDone
| DaemonRequest::NextEvent
| DaemonRequest::EventStreamDropped
| DaemonRequest::RegisterPinnedMemory { .. }
| DaemonRequest::ReadPinnedMemory { .. }
| DaemonRequest::FreePinnedMemory { .. } => true,
}
}
pub fn expects_tcp_json_reply(&self) -> bool {
#[allow(clippy::match_like_matches_macro)]
match self {
DaemonRequest::NodeConfig { .. } => true,
DaemonRequest::Register(NodeRegisterRequest { .. })
| DaemonRequest::Subscribe
| DaemonRequest::CloseOutputs(_)
| DaemonRequest::OutputsDone
| DaemonRequest::NextEvent
| DaemonRequest::SendMessage { .. }
| DaemonRequest::OutputSent { .. }
| DaemonRequest::EventStreamDropped
| DaemonRequest::RegisterPinnedMemory { .. }
| DaemonRequest::ReadPinnedMemory { .. }
| DaemonRequest::FreePinnedMemory { .. } => false,
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct NodeRegisterRequest {
pub dataflow_id: DataflowId,
pub node_id: NodeId,
dora_version: semver::Version,
}
impl NodeRegisterRequest {
pub fn new(dataflow_id: DataflowId, node_id: NodeId) -> Self {
Self {
dataflow_id,
node_id,
dora_version: semver::Version::parse(env!("CARGO_PKG_VERSION")).unwrap(),
}
}
pub fn check_version(&self) -> Result<(), String> {
let crate_version = current_crate_version();
let specified_version = &self.dora_version;
if versions_compatible(&crate_version, specified_version)? {
Ok(())
} else {
Err(format!(
"version mismatch: message format v{} is not compatible \
with expected message format v{crate_version}",
self.dora_version
))
}
}
}
#[derive(Debug, serde::Deserialize, serde::Serialize)]
pub enum DynamicNodeEvent {
NodeConfig { node_id: NodeId },
}