(
mut self,
job_id: JobId,
state: &mut State,
server_ref: &ServerRef,
)
| 52 | |
| 53 | impl RestorerJob { |
| 54 | pub fn restore_job( |
| 55 | mut self, |
| 56 | job_id: JobId, |
| 57 | state: &mut State, |
| 58 | server_ref: &ServerRef, |
| 59 | ) -> crate::Result<Vec<TaskSubmit>> { |
| 60 | log::debug!("Restoring job {job_id}"); |
| 61 | let job = Job::new(job_id, self.job_desc, self.is_open); |
| 62 | state.add_job(job); |
| 63 | let mut result: Vec<TaskSubmit> = Vec::new(); |
| 64 | for submit in self.submit_descs { |
| 65 | if let Some(e) = validate_submit(state.get_job(job_id), &submit.description().task_desc) |
| 66 | { |
| 67 | return Err(HqError::GenericError(format!( |
| 68 | "Job validation failed {e:?}" |
| 69 | ))); |
| 70 | } |
| 71 | let mut new_tasks = submit_job_desc( |
| 72 | state, |
| 73 | server_ref, |
| 74 | job_id, |
| 75 | submit.description().clone(), |
| 76 | submit.submitted_at(), |
| 77 | ); |
| 78 | let job = state.get_job_mut(job_id).unwrap(); |
| 79 | |
| 80 | new_tasks.tasks.retain_mut(|t| { |
| 81 | t.task_deps |
| 82 | .retain(|d| !is_task_completed(&self.tasks, d.job_task_id())); |
| 83 | !is_task_completed(&self.tasks, t.id.job_task_id()) |
| 84 | }); |
| 85 | |
| 86 | for (task_id, job_task) in job.tasks.iter_mut() { |
| 87 | if let Some(task) = self.tasks.get_mut(task_id) { |
| 88 | if task.crash_counter > 0 || task.instance_id.is_some() { |
| 89 | new_tasks.adjust_instance_id_and_crash_counters.insert( |
| 90 | TaskId::new(job_id, *task_id), |
| 91 | ( |
| 92 | task.instance_id.map(|x| x.as_num() + 1).unwrap_or(0).into(), |
| 93 | task.crash_counter, |
| 94 | ), |
| 95 | ); |
| 96 | } |
| 97 | match &task.state { |
| 98 | JobTaskState::Waiting | JobTaskState::Running { .. } => continue, |
| 99 | JobTaskState::Finished { .. } => job.counters.n_finished_tasks += 1, |
| 100 | JobTaskState::Failed { .. } => job.counters.n_failed_tasks += 1, |
| 101 | JobTaskState::Canceled { .. } => job.counters.n_canceled_tasks += 1, |
| 102 | JobTaskState::Aborted { .. } => job.counters.n_aborted_tasks += 1, |
| 103 | } |
| 104 | job_task.state = task.state.clone(); |
| 105 | } |
| 106 | } |
| 107 | if !new_tasks.tasks.is_empty() { |
| 108 | result.push(new_tasks); |
| 109 | } |
| 110 | } |
| 111 | Ok(result) |
no test coverage detected