MCPcopy Create free account
hub / github.com/MegEngine/MegFlow / exec

Method exec

flow-rs/src/node/demux.rs:39–81  ·  view source on GitHub ↗
(&mut self, _: &Context)

Source from the content-addressed store, hash-verified

37 }
38
39 async fn exec(&mut self, _: &Context) -> Result<()> {
40 if let Ok(msg) = self.inp.recv_any().await {
41 let id = msg
42 .info()
43 .to_addr
44 .expect("the envelope has no destination address");
45 let tag = msg.info().tag.as_ref();
46 if msg.is_none() {
47 self.out.fetch_with_cache().await.remove(&id);
48 if let Some(task) = self.tasks.remove(&id) {
49 task.await.ok();
50 }
51 } else {
52 if !self.out.cache().contains_key(&id) {
53 if let Some(tag) = tag {
54 self.tasks.insert(
55 id,
56 self.out
57 .create_spec(
58 id,
59 tag,
60 self.resources.clone().unwrap(),
61 Default::default(),
62 )
63 .await
64 .expect("create subgraph fault"),
65 );
66 } else {
67 self.tasks.insert(
68 id,
69 self.out
70 .create(id, self.resources.clone().unwrap(), Default::default())
71 .await
72 .expect("create subgraph fault"),
73 );
74 }
75 }
76 let out = self.out.fetch_with_cache().await.get(&id).unwrap();
77 out.send_any(msg).await.ok();
78 }
79 }
80 Ok(())
81 }
82}
83
84node_register!("DynDemux", DynDemux);

Callers

nothing calls this directly

Calls 13

recv_anyMethod · 0.80
as_refMethod · 0.80
removeMethod · 0.80
fetch_with_cacheMethod · 0.80
cacheMethod · 0.80
insertMethod · 0.80
create_specMethod · 0.80
createMethod · 0.80
send_anyMethod · 0.80
infoMethod · 0.45
is_noneMethod · 0.45
cloneMethod · 0.45

Tested by

no test coverage detected