(
mut tx: Tx,
mut rx: Rx,
server_dir: ServerDir,
state_ref: StateRef,
senders: &Senders,
end_flag: Arc<Notify>,
)
| 190 | } |
| 191 | |
| 192 | pub async fn client_rpc_loop< |
| 193 | Tx: Sink<ToClientMessage, Error = tako::Error> + Unpin + 'static, |
| 194 | Rx: Stream<Item = tako::Result<FromClientMessage>> + Unpin, |
| 195 | >( |
| 196 | mut tx: Tx, |
| 197 | mut rx: Rx, |
| 198 | server_dir: ServerDir, |
| 199 | state_ref: StateRef, |
| 200 | senders: &Senders, |
| 201 | end_flag: Arc<Notify>, |
| 202 | ) where |
| 203 | Tx::Error: Debug, |
| 204 | { |
| 205 | while let Some(message_result) = rx.next().await { |
| 206 | match message_result { |
| 207 | Ok(message) => { |
| 208 | let response = match message { |
| 209 | FromClientMessage::Submit(msg, stream_opts) => { |
| 210 | let response = submit::handle_submit(&state_ref, senders, msg); |
| 211 | if !response.is_error() { |
| 212 | senders.events.flush_journal().await; |
| 213 | }; |
| 214 | if let Some(mut stream_opts) = stream_opts |
| 215 | && let ToClientMessage::SubmitResponse(SubmitResponse::Ok { |
| 216 | job, .. |
| 217 | }) = &response |
| 218 | { |
| 219 | if !stream_opts.filter.is_filtering_jobs() { |
| 220 | let mut s = Set::new(); |
| 221 | s.insert(job.info.id); |
| 222 | stream_opts.filter.set_jobs(s); |
| 223 | } |
| 224 | start_streaming( |
| 225 | tx, |
| 226 | rx, |
| 227 | state_ref, |
| 228 | senders, |
| 229 | stream_opts, |
| 230 | Some(response), |
| 231 | ) |
| 232 | .await; |
| 233 | break; |
| 234 | } |
| 235 | response |
| 236 | } |
| 237 | FromClientMessage::JobInfo(msg, stream_opts) => { |
| 238 | let response = |
| 239 | compute_job_info(&state_ref, &msg.selector, msg.include_running_tasks); |
| 240 | if let Some(mut stream_opts) = stream_opts |
| 241 | && let ToClientMessage::JobInfoResponse(JobInfoResponse { jobs }) = |
| 242 | &response |
| 243 | { |
| 244 | if !stream_opts.filter.is_filtering_jobs() { |
| 245 | stream_opts |
| 246 | .filter |
| 247 | .set_jobs(jobs.iter().map(|j| j.id).collect()); |
| 248 | } |
| 249 | start_streaming( |
no test coverage detected