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

Method restore_job

crates/hyperqueue/src/server/restore.rs:54–112  ·  view source on GitHub ↗
(
        mut self,
        job_id: JobId,
        state: &mut State,
        server_ref: &ServerRef,
    )

Source from the content-addressed store, hash-verified

52
53impl 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)

Callers 1

Calls 15

validate_submitFunction · 0.85
submit_job_descFunction · 0.85
is_task_completedFunction · 0.85
get_jobMethod · 0.80
descriptionMethod · 0.80
submitted_atMethod · 0.80
job_task_idMethod · 0.80
iter_mutMethod · 0.80
intoMethod · 0.80
pushMethod · 0.80
add_jobMethod · 0.45
cloneMethod · 0.45

Tested by

no test coverage detected