1use 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_resumed_session(record, None)?;
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_owner = self.owner();
42 let controller = controller_owner.controller();
43 let source = controller
44 .state
45 .sessions
46 .get(&selection.session_id)
47 .context("Move session is missing")?;
48 (
49 source.harness_kind,
50 matches!(
51 source.state,
52 SessionState::Running | SessionState::Disconnected
53 ),
54 )
55 };
56 let snapshot = if source_active {
57 crate::controller::move_session::refresh_move_source(
58 &self.session_manager,
59 &selection.session_id,
60 )
61 .await?
62 } else {
63 None
64 };
65 let mut preparation = blocking(move || {
66 let controller = Controller::load()?;
67 mj_core::runtime::block_on(
68 controller.prepare_move_session_controlled(selection, &ProcessExecutor),
69 )?
70 })
71 .await?;
72 preparation.source_unavailable = source_active
73 && snapshot
74 .as_ref()
75 .is_none_or(|snapshot| !snapshot.operational.native_session_is_ready());
76 preparation.active |= preparation.source_unavailable;
77 if let Some(snapshot) = snapshot {
78 let mut operational = snapshot.operational;
79 operational.queued_prompts.clear();
80 operational.checkpoint_barrier = None;
81 preparation.active |= !operational.safe_to_replace(source_harness);
82 }
83 Ok(preparation)
84 }
85
86 pub(crate) fn start_move_session(self: &Arc<Self>, request: MoveSessionRequest) -> Result<()> {
88 self.admit_move_session(request).map(|_| ())
89 }
90
91 fn admit_move_session(self: &Arc<Self>, request: MoveSessionRequest) -> Result<LifecycleWatch> {
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.admit_lifecycle(
100 session_id.clone(),
101 LifecycleKind::Move,
102 super::lifecycle::LifecycleStart {
103 resume_workspace_id: None,
104 request_key: Some(key),
105 create_control: None,
106 phase: LifecyclePhase::Executing,
107 move_operation_id: Some(operation_id.clone()),
108 },
109 move |state, session_id, cancelled| async move {
110 let stopping = std::time::Instant::now();
115 state.stop_subagents_for_suspend(&session_id).await?;
116 tracing::info!(
117 %session_id,
118 phase = "stop sub-agents",
119 elapsed_ms = stopping.elapsed().as_millis() as u64,
120 "move phase finished"
121 );
122 let result = blocking({
123 let state = state.clone();
124 let session_id = session_id.clone();
125 move || {
126 let reserving = std::time::Instant::now();
127 let _reservation = reserve_recovery_or_cancel(
128 &state.recovery_observer,
129 &session_id,
130 &cancelled,
131 )?;
132 tracing::info!(
133 %session_id,
134 phase = "recovery reservation",
135 elapsed_ms = reserving.elapsed().as_millis() as u64,
136 "move phase finished"
137 );
138 let loading = std::time::Instant::now();
139 let mut controller = Controller::load()?;
140 tracing::info!(
141 %session_id,
142 phase = "load controller state",
143 elapsed_ms = loading.elapsed().as_millis() as u64,
144 "move phase finished"
145 );
146 let executor = DaemonStageReportingExecutor::new(
147 CancellableProcessExecutor::new(cancelled),
148 state.clone(),
149 session_id,
150 );
151 let outcome = mj_core::runtime::block_on(
152 controller.move_session_managed_controlled(
153 request,
154 &executor,
155 &state.session_manager,
156 ),
157 )??;
158 Ok(DaemonLifecycleResult::Move(outcome))
159 }
160 })
161 .await;
162 if result.is_err() {
163 state
164 .tell_live_parent_about_stopped_subagents(&session_id)
165 .await;
166 }
167 result
168 },
169 )?;
170 self.set_lifecycle_resume_destination(
171 &session_id,
172 selection.profile_id.clone().unwrap_or_default(),
173 selection.target_template_id.clone().unwrap_or_default(),
174 );
175 Ok(result)
176 }
177
178 pub async fn move_session(
179 self: &Arc<Self>,
180 request: MoveSessionRequest,
181 ) -> Result<MoveOutcome> {
182 let selection = request.preparation.selection.clone();
183 let operation_id = request.preparation.operation_id.clone();
184 let session_id = selection.session_id.clone();
185 let result = self.admit_move_session(request)?;
186 let channel = result.clone();
187 let result = Self::wait_lifecycle_result(result).await;
188 self.remove_completed_lifecycle(&channel);
189 match result {
190 Ok(DaemonLifecycleResult::Move(outcome)) => Ok(outcome),
191 Ok(_) => bail!("move returned an unrelated lifecycle result"),
192 Err(error) => Ok(MoveOutcome {
193 operation_id, session_id, profile_id: selection.profile_id.unwrap_or_default(),
194 target_template_id: selection.target_template_id.unwrap_or_default(), outcome: "failed".into(),
195 error: Some(format!("{error:#}")), recovery: Some("Inspect session status and prepare Move again; any verified checkpoint is retained.".into()),
196 }),
197 }
198 }
199
200 pub(super) fn recover_moves(
201 self: &Arc<Self>,
202 operations: Vec<MoveOperation>,
203 ) -> Result<BTreeSet<String>> {
204 let mut owned = BTreeSet::new();
205 for operation in &operations {
206 crate::controller::move_session::restore_move_queue_hold(operation);
207 }
208 for operation in operations.into_iter().filter(|op| {
209 op.is_active()
210 || self
211 .owner()
212 .controller()
213 .state
214 .sessions
215 .get(&op.selection.session_id)
216 .is_some_and(|session| {
217 matches!(
218 session.state,
219 SessionState::Closing | SessionState::Destroying
220 )
221 })
222 }) {
223 let id = operation.selection.session_id.clone();
224 let key = format!("recovery:{}", operation.operation_id);
225 let phase = if operation.recovery_session.is_some() {
226 LifecyclePhase::MovingDestination
227 } else {
228 LifecyclePhase::Executing
229 };
230 let result = self.admit_lifecycle(
231 id.clone(),
232 LifecycleKind::Move,
233 super::lifecycle::LifecycleStart {
234 resume_workspace_id: None,
235 request_key: Some(key),
236 create_control: None,
237 phase,
238 move_operation_id: None,
239 },
240 move |state, session_id, cancelled| async move {
241 blocking(move || {
242 let _reservation = reserve_recovery_or_cancel(
243 &state.recovery_observer,
244 &session_id,
245 &cancelled,
246 )?;
247 let mut controller = Controller::load()?;
248 let executor = DaemonStageReportingExecutor::new(
249 CancellableProcessExecutor::new(cancelled),
250 state.clone(),
251 session_id,
252 );
253 mj_core::runtime::block_on(controller.recover_move_managed_controlled(
254 operation,
255 &executor,
256 &state.session_manager,
257 ))?
258 .map(DaemonLifecycleResult::Move)
259 })
260 .await
261 },
262 )?;
263 owned.insert(id.clone());
264 let state = self.clone();
265 tokio::spawn(async move {
266 let channel = result.clone();
267 match Self::wait_lifecycle_result(result).await {
268 Ok(DaemonLifecycleResult::Move(outcome))
269 if !matches!(outcome.outcome.as_str(), "completed" | "interrupted") =>
270 {
271 state.push_notice(
272 &id,
273 outcome
274 .error
275 .unwrap_or_else(|| "Move recovery needs attention".into()),
276 );
277 }
278 Err(error) => {
279 state.push_notice(&id, format!("Move recovery failed: {error:#}"))
280 }
281 _ => {}
282 }
283 state.remove_completed_lifecycle(&channel);
284 });
285 }
286 Ok(owned)
287 }
288}