(&mut self, _: &flow_rs::graph::Context)
| 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(); |
nothing calls this directly
no test coverage detected