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

Method exec

flow-plugins/src/video_server.rs:114–201  ·  view source on GitHub ↗
(&mut self, _: &flow_rs::graph::Context)

Source from the content-addressed store, hash-verified

112 }
113 async fn finalize(&mut self) {}
114 async fn exec(&mut self, _: &flow_rs::graph::Context) -> Result<()> {
115 let (s, r) = flow_rs::rt::channel::unbounded();
116 let state = State::new(s);
117 let mapping = state.mapping.clone();
118 let messages = Messages::new().with_lock();
119 let messages_cloned = messages.clone();
120
121 let (spec, filter) = openapi::spec().build(move || {
122 start(state.clone())
123 .or(stop(state.clone()))
124 .or(list(state))
125 .or(get_msgs(messages_cloned))
126 });
127
128 let mut recv_msgs = FuturesUnordered::new();
129 let mut recv_conns = FuturesUnordered::new();
130 let mut spawn_decode = FuturesUnordered::new();
131 let listen = serve(filter.or(openapi_docs(spec)))
132 .run(([0, 0, 0, 0], self.port))
133 .fuse();
134 recv_conns.push(self.inp.fetch());
135 spawn_decode.push(r.recv());
136
137 pin_mut!(listen);
138
139 loop {
140 select! {
141 // server listen
142 _ = listen => break,
143 // spawn task to wait message from subgraph
144 conns = recv_conns.select_next_some() => {
145 if let Ok((id, port)) = conns {
146 let messages = messages.clone();
147 recv_msgs.push(async move {
148 while let Ok(mut msg) = port.recv().await {
149 let messages_map = &mut messages.lock().await.mapping;
150 let messages = messages_map.entry(id).or_default();
151 let msg = Python::with_gil(|py| -> PyResult<_> {
152 let msg: PyObject = msg.unpack();
153 let msg = msg.as_ref(py).extract()?;
154 Ok(msg)
155 }).expect("plugin[VideoServier] only support python string as input");
156 messages.push(msg);
157 }
158 });
159 recv_conns.push(self.inp.fetch());
160 }
161 },
162 // wait message from subgraph
163 _ = recv_msgs.select_next_some() => {},
164 // spawn subgraph
165 ret = spawn_decode.select_next_some() => {
166 if let Ok((id, url, waker)) = ret {
167 self.out.create(id, self.resources.clone().unwrap(), Default::default()).await.expect("broker has closed");
168 let (_, port) = self.out.fetch().await.expect("broker has closed");
169 let url = urlencoding::decode(&url).unwrap().into_owned();
170 let url_cloned = url.clone();
171 let port_cloned = port.clone();

Callers

nothing calls this directly

Calls 12

unboundedFunction · 0.85
stopFunction · 0.85
listFunction · 0.85
get_msgsFunction · 0.85
with_lockMethod · 0.80
pushMethod · 0.80
startFunction · 0.70
cloneMethod · 0.45
buildMethod · 0.45
runMethod · 0.45
fetchMethod · 0.45
recvMethod · 0.45

Tested by

no test coverage detected