(
state_ref: &StateRef,
senders: &Senders,
selector: IdSelector,
)
| 577 | } |
| 578 | |
| 579 | fn handle_worker_stop( |
| 580 | state_ref: &StateRef, |
| 581 | senders: &Senders, |
| 582 | selector: IdSelector, |
| 583 | ) -> ToClientMessage { |
| 584 | log::debug!("Client asked for worker termination {selector:?}"); |
| 585 | let mut responses: Vec<(WorkerId, StopWorkerResponse)> = Vec::new(); |
| 586 | |
| 587 | let worker_ids: Vec<WorkerId> = match selector { |
| 588 | IdSelector::Specific(array) => array.iter().map(|id| id.into()).collect(), |
| 589 | IdSelector::All => state_ref |
| 590 | .get() |
| 591 | .get_workers() |
| 592 | .iter() |
| 593 | .filter(|(_, worker)| worker.make_info(None).ended.is_none()) |
| 594 | .map(|(_, worker)| worker.worker_id()) |
| 595 | .collect(), |
| 596 | IdSelector::LastN(n) => { |
| 597 | let mut ids: Vec<_> = state_ref.get().get_workers().keys().copied().collect(); |
| 598 | ids.sort_by_key(|&k| std::cmp::Reverse(k)); |
| 599 | ids.truncate(n as usize); |
| 600 | ids |
| 601 | } |
| 602 | }; |
| 603 | |
| 604 | for worker_id in worker_ids { |
| 605 | if let Some(worker) = state_ref.get().get_worker(worker_id) { |
| 606 | if worker.make_info(None).ended.is_some() { |
| 607 | responses.push((worker_id, StopWorkerResponse::AlreadyStopped)); |
| 608 | continue; |
| 609 | } |
| 610 | } else { |
| 611 | responses.push((worker_id, StopWorkerResponse::InvalidWorker)); |
| 612 | continue; |
| 613 | } |
| 614 | let response = senders.server_control.stop_worker(worker_id); |
| 615 | |
| 616 | match response { |
| 617 | Ok(()) => responses.push((worker_id, StopWorkerResponse::Stopped)), |
| 618 | Err(err) => { |
| 619 | responses.push((worker_id, StopWorkerResponse::Failed(err.to_string()))); |
| 620 | log::error!("Unable to stop worker: {worker_id} error: {err:?}"); |
| 621 | } |
| 622 | } |
| 623 | } |
| 624 | ToClientMessage::StopWorkerResponse(responses) |
| 625 | } |
| 626 | |
| 627 | fn compute_job_detail( |
| 628 | state_ref: &StateRef, |
no test coverage detected