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

Method exec

flow-rs/src/node/reorder.rs:32–72  ·  view source on GitHub ↗
(&mut self, _: &Context)

Source from the content-addressed store, hash-verified

30 async fn finalize(&mut self) {}
31
32 async fn exec(&mut self, _: &Context) -> Result<()> {
33 if let Ok(msg) = self.inp.recv_any().await {
34 let id = msg
35 .info()
36 .partial_id
37 .expect("partial_id required by reorder");
38 assert!(id >= self.seq_id);
39 if id == self.seq_id {
40 self.seq_id += 1;
41 self.out.send_any(msg).await.ok();
42 } else {
43 self.cache.insert(id, msg);
44 }
45
46 let mut stop = self.seq_id;
47 for &id in self.cache.keys() {
48 assert!(id >= self.seq_id);
49 if id == self.seq_id {
50 self.seq_id += 1;
51 } else {
52 stop = id;
53 break;
54 }
55 }
56 if stop != self.seq_id {
57 let rest = if stop > self.seq_id {
58 let rest = self.cache.split_off(&stop);
59 std::mem::replace(&mut self.cache, rest)
60 } else {
61 std::mem::take(&mut self.cache)
62 };
63
64 for (_, msg) in rest {
65 self.out.send_any(msg).await.ok();
66 }
67 }
68 } else {
69 assert!(self.cache.is_empty());
70 }
71 Ok(())
72 }
73}
74
75node_register!("Reorder", Reorder);

Callers

nothing calls this directly

Calls 5

recv_anyMethod · 0.80
send_anyMethod · 0.80
insertMethod · 0.80
infoMethod · 0.45
keysMethod · 0.45

Tested by

no test coverage detected