(&mut self, _: &Context)
| 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 | |
| 84 | node_register!("DynDemux", DynDemux); |
nothing calls this directly
no test coverage detected