mj_controller/daemon/
session_move.rs1use super::*;
4use mj_core::state::MoveOperation;
5
6pub(super) fn load_controller_for_resume(request: &ResumeSessionRequest) -> Result<Controller> {
7 let mut controller = Controller::load()?;
8 if let Some(operation) = crate::database::load_move_operation(&request.session_id)?
9 && matches!(
10 operation.phase,
11 mj_core::state::MovePhase::Failed | mj_core::state::MovePhase::Cancelled
12 )
13 && !operation.queue_admission_started
14 && operation.source_profile_id == request.profile_id
15 && operation.source_target_template_id == request.target_template_id
16 {
17 let record = controller
20 .state
21 .sessions
22 .get_mut(&request.session_id)
23 .context("Move recovery session is missing")?;
24 record.resource_allocation = operation.source_resource_allocation;
25 record.additional_mounts = operation.source_additional_mounts;
26 if let Some(previous) = operation.recovery_session {
27 record.container_cpus = previous.container_cpus;
28 record.container_memory = previous.container_memory;
29 }
30 crate::database::save_session(record)?;
31 }
32 Ok(controller)
33}
34
35impl RuntimeState {
36 pub async fn prepare_move_session(
37 self: &Arc<Self>,
38 selection: MoveSelection,
39 ) -> Result<MovePreparation> {
40 let (source_harness, source_active) = {
41 let controller = self
42 .controller
43 .lock()
44 .unwrap_or_else(PoisonError::into_inner);
45 let source = controller
46 .state
47 .sessions
48 .get(&selection.session_id)
49 .context("Move session is missing")?;
50 (
51 source.harness_kind,
52 matches!(
53 source.state,
54 SessionState::Running | SessionState::Disconnected
55 ),
56 )
57 };
58 let snapshot = if source_active {
59 crate::controller::move_session::refresh_move_source(
60 &self.session_manager,
61 &selection.session_id,
62 )
63 .await?
64 } else {
65 None
66 };
67 let runtime = tokio::runtime::Handle::current();
68 let mut preparation = blocking(move || {
69 let controller = Controller::load()?;
70 runtime
71 .block_on(controller.prepare_move_session_controlled(selection, &ProcessExecutor))
72 })
73 .await?;
74 preparation.source_unavailable = source_active
75 && snapshot
76 .as_ref()
77 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
78 preparation.active |= preparation.source_unavailable;
79 if let Some(snapshot) = snapshot {
80 let mut operational = snapshot.operational;
81 operational.queued_prompts.clear();
82 operational.checkpoint_barrier = None;
83 preparation.active |= !operational.safe_to_replace(source_harness);
84 }
85 Ok(preparation)
86 }
87
88 pub async fn move_session(
89 self: &Arc<Self>,
90 request: MoveSessionRequest,
91 ) -> Result<MoveOutcome> {
92 let selection = request.preparation.selection.clone();
93 let operation_id = request.preparation.operation_id.clone();
94 let key =
97 serde_json::to_string(&(&selection, request.queue, request.acknowledge_interruption))?;
98 let session_id = selection.session_id.clone();
99 let result = self.start_or_join_lifecycle_with_key(
100 session_id.clone(),
101 LifecycleKind::Move,
102 None,
103 Some(key),
104 move |state, session_id, cancelled| async move {
105 let runtime = tokio::runtime::Handle::current();
106 blocking(move || {
107 let _reservation = reserve_recovery_or_cancel(
108 &state.recovery_observer,
109 &session_id,
110 &cancelled,
111 )?;
112 let mut controller = Controller::load()?;
113 let executor = DaemonStageReportingExecutor::new(
114 CancellableProcessExecutor::new(cancelled),
115 state.clone(),
116 session_id,
117 );
118 let outcome = runtime.block_on(controller.move_session_managed_controlled(
119 request,
120 &executor,
121 &state.session_manager,
122 ))?;
123 Ok(DaemonLifecycleResult::Move(outcome))
124 })
125 .await
126 },
127 )?;
128 self.set_lifecycle_resume_destination(
129 &session_id,
130 selection.profile_id.clone().unwrap_or_default(),
131 selection.target_template_id.clone().unwrap_or_default(),
132 );
133 let channel = result.clone();
134 let result = Self::wait_lifecycle_result(result).await;
135 self.remove_completed_lifecycle(&channel);
136 match result {
137 Ok(DaemonLifecycleResult::Move(outcome)) => Ok(outcome),
138 Ok(_) => bail!("move returned an unrelated lifecycle result"),
139 Err(error) => Ok(MoveOutcome {
140 operation_id, session_id, profile_id: selection.profile_id.unwrap_or_default(),
141 target_template_id: selection.target_template_id.unwrap_or_default(), outcome: "failed".into(),
142 error: Some(format!("{error:#}")), recovery: Some("Inspect session status and prepare Move again; any verified checkpoint is retained.".into()),
143 }),
144 }
145 }
146
147 pub(super) fn recover_moves(
148 self: &Arc<Self>,
149 operations: Vec<MoveOperation>,
150 ) -> Result<BTreeSet<String>> {
151 let mut owned = BTreeSet::new();
152 for operation in &operations {
153 crate::controller::move_session::restore_move_queue_hold(operation);
154 }
155 for operation in operations.into_iter().filter(|op| {
156 op.is_active()
157 || self
158 .controller
159 .lock()
160 .unwrap_or_else(PoisonError::into_inner)
161 .state
162 .sessions
163 .get(&op.selection.session_id)
164 .is_some_and(|session| {
165 matches!(
166 session.state,
167 SessionState::Closing | SessionState::Destroying
168 )
169 })
170 }) {
171 let id = operation.selection.session_id.clone();
172 let key = format!("recovery:{}", operation.operation_id);
173 let result = self.start_or_join_lifecycle_with_key(
174 id.clone(),
175 LifecycleKind::Move,
176 None,
177 Some(key),
178 move |state, session_id, cancelled| async move {
179 let runtime = tokio::runtime::Handle::current();
180 blocking(move || {
181 let _reservation = reserve_recovery_or_cancel(
182 &state.recovery_observer,
183 &session_id,
184 &cancelled,
185 )?;
186 let mut controller = Controller::load()?;
187 let executor = DaemonStageReportingExecutor::new(
188 CancellableProcessExecutor::new(cancelled),
189 state.clone(),
190 session_id,
191 );
192 runtime
193 .block_on(controller.recover_move_managed_controlled(
194 operation,
195 &executor,
196 &state.session_manager,
197 ))
198 .map(DaemonLifecycleResult::Move)
199 })
200 .await
201 },
202 )?;
203 owned.insert(id.clone());
204 let state = self.clone();
205 tokio::spawn(async move {
206 let channel = result.clone();
207 match Self::wait_lifecycle_result(result).await {
208 Ok(DaemonLifecycleResult::Move(outcome)) if outcome.outcome != "completed" => {
209 state.push_notice(
210 &id,
211 outcome
212 .error
213 .unwrap_or_else(|| "Move recovery needs attention".into()),
214 );
215 }
216 Err(error) => {
217 state.push_notice(&id, format!("Move recovery failed: {error:#}"))
218 }
219 _ => {}
220 }
221 state.remove_completed_lifecycle(&channel);
222 });
223 }
224 Ok(owned)
225 }
226}