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 mut preparation = blocking(move || {
68 let controller = Controller::load()?;
69 mj_core::runtime::block_on(
70 controller.prepare_move_session_controlled(selection, &ProcessExecutor),
71 )?
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 blocking(move || {
106 let _reservation = reserve_recovery_or_cancel(
107 &state.recovery_observer,
108 &session_id,
109 &cancelled,
110 )?;
111 let mut controller = Controller::load()?;
112 let executor = DaemonStageReportingExecutor::new(
113 CancellableProcessExecutor::new(cancelled),
114 state.clone(),
115 session_id,
116 );
117 let outcome =
118 mj_core::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 blocking(move || {
180 let _reservation = reserve_recovery_or_cancel(
181 &state.recovery_observer,
182 &session_id,
183 &cancelled,
184 )?;
185 let mut controller = Controller::load()?;
186 let executor = DaemonStageReportingExecutor::new(
187 CancellableProcessExecutor::new(cancelled),
188 state.clone(),
189 session_id,
190 );
191 mj_core::runtime::block_on(controller.recover_move_managed_controlled(
192 operation,
193 &executor,
194 &state.session_manager,
195 ))?
196 .map(DaemonLifecycleResult::Move)
197 })
198 .await
199 },
200 )?;
201 owned.insert(id.clone());
202 let state = self.clone();
203 tokio::spawn(async move {
204 let channel = result.clone();
205 match Self::wait_lifecycle_result(result).await {
206 Ok(DaemonLifecycleResult::Move(outcome)) if outcome.outcome != "completed" => {
207 state.push_notice(
208 &id,
209 outcome
210 .error
211 .unwrap_or_else(|| "Move recovery needs attention".into()),
212 );
213 }
214 Err(error) => {
215 state.push_notice(&id, format!("Move recovery failed: {error:#}"))
216 }
217 _ => {}
218 }
219 state.remove_completed_lifecycle(&channel);
220 });
221 }
222 Ok(owned)
223 }
224}