MCPcopy Create free account
hub / github.com/It4innovations/hyperqueue / client_rpc_loop

Function client_rpc_loop

crates/hyperqueue/src/server/client/mod.rs:192–368  ·  view source on GitHub ↗
(
    mut tx: Tx,
    mut rx: Rx,
    server_dir: ServerDir,
    state_ref: StateRef,
    senders: &Senders,
    end_flag: Arc<Notify>,
)

Source from the content-addressed store, hash-verified

190}
191
192pub 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(

Callers 1

handle_clientFunction · 0.85

Calls 15

handle_submitFunction · 0.85
start_streamingFunction · 0.85
compute_job_infoFunction · 0.85
handle_get_listFunction · 0.85
handle_worker_infoFunction · 0.85
handle_worker_stopFunction · 0.85
handle_job_cancelFunction · 0.85
handle_job_forgetFunction · 0.85
compute_job_detailFunction · 0.85
handle_autoalloc_messageFunction · 0.85
handle_open_jobFunction · 0.85
handle_job_closeFunction · 0.85

Tested by

no test coverage detected