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,
metadata_version: u16,
}
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(),
metadata_version: Metadata::CURRENT_VERSION,
}
}
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)? {
return Err(format!(
"version mismatch: message format v{} is not compatible \
with expected message format v{crate_version}",
self.dora_version
));
}
if self.metadata_version != Metadata::CURRENT_VERSION {
return Err(format!(
"message wire-format mismatch: node speaks metadata format v{} \
but this daemon speaks v{}. The node and daemon were built from \
dora revisions with incompatible message formats; rebuild both \
from the same revision.",
self.metadata_version,
Metadata::CURRENT_VERSION
));
}
Ok(())
}
}
#[derive(Debug, serde::Deserialize, serde::Serialize)]
pub enum DynamicNodeEvent {
NodeConfig { node_id: NodeId },
}
#[cfg(test)]
mod register_version_tests {
use super::*;
fn request() -> NodeRegisterRequest {
NodeRegisterRequest::new(uuid::Uuid::nil(), NodeId::from("test-node".to_string()))
}
#[test]
fn new_stamps_current_metadata_version() {
assert_eq!(request().metadata_version, Metadata::CURRENT_VERSION);
}
#[test]
fn check_version_accepts_a_matching_request() {
request().check_version().unwrap();
}
#[test]
fn check_version_rejects_a_wire_format_mismatch() {
let mut req = request();
req.metadata_version = Metadata::CURRENT_VERSION.wrapping_add(1);
let err = req
.check_version()
.expect_err("mismatched metadata wire version must be rejected");
assert!(
err.contains("wire-format") && err.contains(&Metadata::CURRENT_VERSION.to_string()),
"error should name the wire-format mismatch and the expected version, got: {err}"
);
}
}