dora_message/
node_to_daemon.rs1pub use crate::common::{DataMessage, LogLevel, LogMessage, SharedMemoryId, Timestamped};
2use crate::{
3 DataflowId, current_crate_version,
4 id::{DataId, NodeId},
5 metadata::Metadata,
6 versions_compatible,
7};
8
9#[allow(clippy::large_enum_variant)]
10#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
11#[non_exhaustive]
12pub enum DaemonRequest {
13 Register(NodeRegisterRequest),
14 Subscribe,
15 SendMessage {
16 output_id: DataId,
17 metadata: Metadata,
18 data: Option<DataMessage>,
19 },
20 OutputSent {
21 output_id: DataId,
22 metadata: Metadata,
23 },
24 CloseOutputs(Vec<DataId>),
25 OutputsDone,
27 NextEvent,
28 EventStreamDropped,
29 NodeConfig {
30 node_id: NodeId,
31 },
32 ExtensionStore {
40 namespace: String,
41 key: String,
42 value: Vec<u8>,
43 },
44 ExtensionLoad {
47 namespace: String,
48 key: String,
49 remove: bool,
50 },
51 ExtensionDrop {
55 namespace: String,
56 key: String,
57 },
58 ExtensionRequest {
66 namespace: String,
67 #[serde(with = "crate::bulk_bytes::vec")]
68 payload: Vec<u8>,
69 },
70}
71
72impl DaemonRequest {
73 pub fn encode_size_hint(&self) -> usize {
81 match self {
82 DaemonRequest::SendMessage { data, .. } => data.as_ref().map_or(0, DataMessage::len),
83 DaemonRequest::Register(_)
84 | DaemonRequest::Subscribe
85 | DaemonRequest::OutputSent { .. }
86 | DaemonRequest::CloseOutputs(_)
87 | DaemonRequest::OutputsDone
88 | DaemonRequest::NextEvent
89 | DaemonRequest::EventStreamDropped
90 | DaemonRequest::NodeConfig { .. } => 0,
91 DaemonRequest::ExtensionStore { value, .. } => value.len(),
93 DaemonRequest::ExtensionLoad { .. } | DaemonRequest::ExtensionDrop { .. } => 0,
94 DaemonRequest::ExtensionRequest { payload, .. } => payload.len(),
95 }
96 }
97
98 pub fn expects_tcp_binary_reply(&self) -> bool {
99 #[allow(clippy::match_like_matches_macro)]
100 match self {
101 DaemonRequest::SendMessage { .. }
102 | DaemonRequest::OutputSent { .. }
103 | DaemonRequest::NodeConfig { .. } => false,
104 DaemonRequest::Register(NodeRegisterRequest { .. })
105 | DaemonRequest::Subscribe
106 | DaemonRequest::CloseOutputs(_)
107 | DaemonRequest::OutputsDone
108 | DaemonRequest::NextEvent
109 | DaemonRequest::EventStreamDropped
110 | DaemonRequest::ExtensionRequest { .. }
111 | DaemonRequest::ExtensionStore { .. }
112 | DaemonRequest::ExtensionLoad { .. }
113 | DaemonRequest::ExtensionDrop { .. } => true,
114 }
115 }
116
117 pub fn expects_tcp_json_reply(&self) -> bool {
118 #[allow(clippy::match_like_matches_macro)]
119 match self {
120 DaemonRequest::NodeConfig { .. } => true,
121 DaemonRequest::Register(NodeRegisterRequest { .. })
122 | DaemonRequest::Subscribe
123 | DaemonRequest::CloseOutputs(_)
124 | DaemonRequest::OutputsDone
125 | DaemonRequest::NextEvent
126 | DaemonRequest::SendMessage { .. }
127 | DaemonRequest::OutputSent { .. }
128 | DaemonRequest::EventStreamDropped
129 | DaemonRequest::ExtensionRequest { .. }
130 | DaemonRequest::ExtensionStore { .. }
131 | DaemonRequest::ExtensionLoad { .. }
132 | DaemonRequest::ExtensionDrop { .. } => false,
133 }
134 }
135}
136
137#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
138pub struct NodeRegisterRequest {
139 pub dataflow_id: DataflowId,
140 pub node_id: NodeId,
141 dora_version: semver::Version,
142 metadata_version: u16,
159}
160
161impl NodeRegisterRequest {
162 pub fn new(dataflow_id: DataflowId, node_id: NodeId) -> Self {
163 Self {
164 dataflow_id,
165 node_id,
166 dora_version: semver::Version::parse(env!("CARGO_PKG_VERSION")).unwrap(),
167 metadata_version: Metadata::CURRENT_VERSION,
168 }
169 }
170
171 pub fn check_version(&self) -> Result<(), String> {
172 let crate_version = current_crate_version();
173 let specified_version = &self.dora_version;
174
175 if !versions_compatible(&crate_version, specified_version)? {
176 return Err(format!(
177 "version mismatch: message format v{} is not compatible \
178 with expected message format v{crate_version}",
179 self.dora_version
180 ));
181 }
182
183 if self.metadata_version != Metadata::CURRENT_VERSION {
188 return Err(format!(
189 "message wire-format mismatch: node speaks metadata format v{} \
190 but this daemon speaks v{}. The node and daemon were built from \
191 dora revisions with incompatible message formats; rebuild both \
192 from the same revision.",
193 self.metadata_version,
194 Metadata::CURRENT_VERSION
195 ));
196 }
197
198 Ok(())
199 }
200}
201
202#[derive(Debug, serde::Deserialize, serde::Serialize)]
203pub enum DynamicNodeEvent {
204 NodeConfig { node_id: NodeId },
205}
206
207#[cfg(test)]
208mod register_version_tests {
209 use super::*;
210
211 fn request() -> NodeRegisterRequest {
212 NodeRegisterRequest::new(uuid::Uuid::nil(), NodeId::from("test-node".to_string()))
213 }
214
215 #[test]
216 fn new_stamps_current_metadata_version() {
217 assert_eq!(request().metadata_version, Metadata::CURRENT_VERSION);
218 }
219
220 #[test]
221 fn check_version_accepts_a_matching_request() {
222 request().check_version().unwrap();
225 }
226
227 #[test]
228 fn check_version_rejects_a_wire_format_mismatch() {
229 let mut req = request();
234 req.metadata_version = Metadata::CURRENT_VERSION.wrapping_add(1);
235 let err = req
236 .check_version()
237 .expect_err("mismatched metadata wire version must be rejected");
238 assert!(
239 err.contains("wire-format") && err.contains(&Metadata::CURRENT_VERSION.to_string()),
240 "error should name the wire-format mismatch and the expected version, got: {err}"
241 );
242 }
243}