zenoh_flow/runtime/dataflow/instance/
mod.rs1pub mod builtin;
16pub mod runners;
17
18use self::runners::connector::{ZenohReceiver, ZenohSender};
19use self::runners::Runner;
20use super::DataFlow;
21use crate::io::{Inputs, Outputs};
22use crate::model::record::{LinkRecord, ZFConnectorKind};
23use crate::prelude::{Context, Node};
24use crate::runtime::InstanceContext;
25use crate::types::NodeId;
26use crate::zfresult::ErrorKind;
27use crate::Result;
28use crate::{bail, zferror};
29use std::collections::HashMap;
30use std::ops::Deref;
31use std::sync::Arc;
32use uhlc::HLC;
33
34pub struct DataFlowInstance {
40 pub(crate) _instance_context: Arc<InstanceContext>,
41 pub(crate) data_flow: DataFlow,
42 pub(crate) runners: HashMap<NodeId, Runner>,
43}
44
45impl Deref for DataFlowInstance {
46 type Target = DataFlow;
47
48 fn deref(&self) -> &Self::Target {
49 &self.data_flow
50 }
51}
52
53impl DataFlowInstance {
54 pub fn get_sinks(&self) -> Vec<NodeId> {
60 self.sink_constructors.keys().cloned().collect()
61 }
62
63 pub fn get_sources(&self) -> Vec<NodeId> {
69 self.source_constructors.keys().cloned().collect()
70 }
71
72 pub fn get_operators(&self) -> Vec<NodeId> {
78 self.operator_constructors.keys().cloned().collect()
79 }
80
81 pub fn get_connectors(&self) -> Vec<NodeId> {
87 self.connectors.keys().cloned().collect()
88 }
89
90 pub fn start_node(&mut self, node_id: &NodeId) -> Result<()> {
101 if let Some(runner) = self.runners.get_mut(node_id) {
102 runner.start();
103 return Ok(());
104 }
105
106 bail!(
107 ErrorKind::NodeNotFound(node_id.clone()),
108 "Node < {} > not found",
109 node_id
110 )
111 }
112
113 pub async fn stop_node(&mut self, node_id: &NodeId) -> Result<()> {
125 if let Some(runner) = self.runners.get_mut(node_id) {
126 return runner.stop().await;
127 }
128
129 bail!(
130 ErrorKind::NodeNotFound(node_id.clone()),
131 "Node < {} > not found",
132 node_id
133 )
134 }
135
136 pub async fn try_instantiate(data_flow: DataFlow, hlc: Arc<HLC>) -> Result<Self> {
145 let instance_context = Arc::new(InstanceContext {
146 flow_id: data_flow.flow.clone(),
147 instance_id: data_flow.uuid,
148 runtime: data_flow.context.clone(),
149 });
150
151 let mut node_ids: Vec<NodeId> = Vec::with_capacity(
152 data_flow.source_constructors.len()
153 + data_flow.operator_constructors.len()
154 + data_flow.sink_constructors.len()
155 + data_flow.connectors.len(),
156 );
157
158 node_ids.append(
159 &mut data_flow
160 .source_constructors
161 .keys()
162 .cloned()
163 .collect::<Vec<_>>(),
164 );
165 node_ids.append(
166 &mut data_flow
167 .operator_constructors
168 .keys()
169 .cloned()
170 .collect::<Vec<_>>(),
171 );
172 node_ids.append(
173 &mut data_flow
174 .sink_constructors
175 .keys()
176 .cloned()
177 .collect::<Vec<_>>(),
178 );
179 node_ids.append(&mut data_flow.connectors.keys().cloned().collect::<Vec<_>>());
180
181 let mut links = create_links(&node_ids, &data_flow.links, hlc.clone())?;
182
183 let context = Context::new(&instance_context);
184
185 let mut runners = HashMap::with_capacity(data_flow.source_constructors.len());
186 for (source_id, source_constructor) in &data_flow.source_constructors {
187 let (_, outputs) = links.remove(source_id).ok_or_else(|| {
188 zferror!(
189 ErrorKind::IOError,
190 "Links for Source < {} > were not created.",
191 &source_id
192 )
193 })?;
194
195 let source = (source_constructor.constructor)(
196 context.clone(),
197 source_constructor.configuration.clone(),
198 outputs,
199 )
200 .await?;
201
202 let runner = Runner::new(source);
203 runners.insert(source_id.clone(), runner);
204 }
205
206 for (operator_id, operator_constructor) in &data_flow.operator_constructors {
207 let (inputs, outputs) = links.remove(operator_id).ok_or_else(|| {
208 zferror!(
209 ErrorKind::IOError,
210 "Links for Operator < {} > were not created.",
211 &operator_id
212 )
213 })?;
214
215 let operator = (operator_constructor.constructor)(
216 context.clone(),
217 operator_constructor.configuration.clone(),
218 inputs,
219 outputs,
220 )
221 .await?;
222
223 let runner = Runner::new(operator);
224 runners.insert(operator_id.clone(), runner);
225 }
226
227 for (sink_id, sink_constructor) in &data_flow.sink_constructors {
228 let (inputs, _) = links.remove(sink_id).ok_or_else(|| {
229 zferror!(
230 ErrorKind::IOError,
231 "Links for Sink < {} > were not created.",
232 &sink_id
233 )
234 })?;
235
236 let sink = (sink_constructor.constructor)(
237 context.clone(),
238 sink_constructor.configuration.clone(),
239 inputs,
240 )
241 .await?;
242
243 let runner = Runner::new(sink);
244 runners.insert(sink_id.clone(), runner);
245 }
246
247 for (connector_id, connector_record) in &data_flow.connectors {
248 let node = match &connector_record.kind {
249 ZFConnectorKind::Sender => {
250 let (inputs, _) = links.remove(connector_id).ok_or_else(|| {
251 zferror!(
252 ErrorKind::IOError,
253 "Links for Sink < {} > were not created.",
254 connector_id
255 )
256 })?;
257 Arc::new(
258 ZenohSender::new(connector_record, instance_context.clone(), inputs)
259 .await?,
260 ) as Arc<dyn Node>
261 }
262 ZFConnectorKind::Receiver => {
263 let (_, outputs) = links.remove(connector_id).ok_or_else(|| {
264 zferror!(
265 ErrorKind::IOError,
266 "Links for Source < {} > were not created.",
267 &connector_id
268 )
269 })?;
270 Arc::new(
271 ZenohReceiver::new(connector_record, instance_context.clone(), outputs)
272 .await?,
273 ) as Arc<dyn Node>
274 }
275 };
276
277 let runner = Runner::new(node);
278 runners.insert(connector_id.clone(), runner);
279 }
280
281 Ok(DataFlowInstance {
282 _instance_context: instance_context,
283 data_flow,
284 runners,
285 })
286 }
287}
288
289pub(crate) fn create_links(
295 nodes: &[NodeId],
296 links: &[LinkRecord],
297 hlc: Arc<HLC>,
298) -> Result<HashMap<NodeId, (Inputs, Outputs)>> {
299 let mut io: HashMap<NodeId, (Inputs, Outputs)> = HashMap::with_capacity(nodes.len());
300
301 for link_desc in links {
302 let upstream_node = link_desc.from.node.clone();
303 let downstream_node = link_desc.to.node.clone();
304
305 if !nodes.contains(&upstream_node) || !nodes.contains(&downstream_node) {
308 continue;
309 }
310
311 let (tx, rx) = flume::unbounded();
314 let from = link_desc.from.output.clone();
315 let to = link_desc.to.input.clone();
316
317 match io.get_mut(&upstream_node) {
318 Some((_, outputs)) => outputs.insert(from.clone(), tx),
319 None => {
320 let inputs = Inputs::new();
321 let mut outputs = Outputs::new(hlc.clone());
322 outputs.insert(from.clone(), tx);
323
324 io.insert(upstream_node, (inputs, outputs));
325 }
326 }
327
328 match io.get_mut(&downstream_node) {
329 Some((inputs, _)) => inputs.insert(to.clone(), rx),
330 None => {
331 let outputs = Outputs::new(hlc.clone());
332
333 let mut inputs = Inputs::new();
334 inputs.insert(to.clone(), rx);
335
336 io.insert(downstream_node, (inputs, outputs));
337 }
338 }
339 }
340
341 Ok(io)
342}