1use super::*;
2
3impl RuntimeState {
4 pub async fn resume_session(self: &Arc<Self>, request: ResumeSessionRequest) -> Result<()> {
12 let _admission = self
13 .workspace_resume_gate(&request.workspace_id)
14 .read_owned()
15 .await;
16 let session_id = request.session_id.clone();
17 self.wait_for_deferred_cleanup(&session_id).await?;
18 let already_running = blocking({
21 let session_id = session_id.clone();
22 move || {
23 let controller = Controller::load()?;
24 Ok(controller
25 .state
26 .sessions
27 .get(&session_id)
28 .is_some_and(|session| session.state == SessionState::Running))
29 }
30 })
31 .await?;
32 if already_running {
33 return Ok(());
34 }
35 let profile_id = request.profile_id.clone();
36 let target_template_id = request.target_template_id.clone();
37 let workspace_id = request.workspace_id.clone();
38 let rebind_workspace_id = workspace_id.clone();
39 let operation_session_id = session_id.clone();
40 let result = self.start_or_join_lifecycle_for_workspace(
41 session_id,
42 LifecycleKind::Resume,
43 Some(workspace_id),
44 move |state, session_id, cancelled| async move {
45 let _recovery_reservation = tokio::task::spawn_blocking({
46 let observer = state.recovery_observer.clone();
47 let session_id = session_id.clone();
48 let cancelled = cancelled.clone();
49 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
50 })
51 .await
52 .context("reserve recovery for daemon resume task")??;
53 blocking({
54 let session_id = session_id.clone();
55 move || {
56 crate::database::reassign_resumable_session_workspace(
57 &session_id,
58 &rebind_workspace_id,
59 )
60 }
61 })
62 .await?;
63 let restore_request = request.clone();
64 let mut controller = tokio::task::spawn_blocking(move || {
65 session_move::load_controller_for_resume(&restore_request)
66 })
67 .await
68 .context("load controller for daemon resume task")??;
69 let executor = DaemonStageReportingExecutor::new(
70 CancellableProcessExecutor::new(cancelled),
71 state.clone(),
72 session_id.clone(),
73 );
74 let materialized = controller
75 .resume_session_controlled_with_repository_preflight(
76 &session_id,
77 &request.profile_id,
78 &request.target_template_id,
79 SessionResumeOptions {
80 additional_mounts: request.additional_mounts,
81 resource_allocation: request.resource_allocation,
82 discard_queue: request.discard_queue,
83 },
84 request.repository_preflight,
85 &executor,
86 )
87 .await?;
88 let _ = materialized;
92 Ok(DaemonLifecycleResult::Done)
93 },
94 )?;
95 self.set_lifecycle_resume_destination(
96 &operation_session_id,
97 profile_id,
98 target_template_id,
99 );
100 let channel = result.clone();
101 let result = Self::wait_lifecycle_result(result).await;
102 self.remove_completed_lifecycle(&channel);
103 match result? {
104 DaemonLifecycleResult::Done => {}
105 DaemonLifecycleResult::Move(_) => unreachable!("resume cannot return a move outcome"),
106 DaemonLifecycleResult::DeferredCleanup => {
107 unreachable!("session resume cannot schedule target cleanup")
108 }
109 }
110 blocking(move || {
111 if let Some(mut operation) =
112 crate::database::load_move_operation(&operation_session_id)?
113 && !operation.queue_admission_started
114 {
115 operation.phase = mj_core::state::MovePhase::Cancelled;
116 operation.queue_admission_finished = true;
117 operation.updated_at = chrono::Utc::now().to_rfc3339();
118 operation.error = Some("Recovered through an explicit Resume operation".into());
119 crate::database::save_move_operation(&operation)?;
120 }
121 Ok(())
122 })
123 .await?;
124 Ok(())
125 }
126
127 pub(super) async fn discard_since_checkpoint(
128 self: &Arc<Self>,
129 session_id: String,
130 checkpoint: mj_core::state::CheckpointMetadata,
131 ) -> Result<()> {
132 let children = blocking({
133 let session_id = session_id.clone();
134 move || {
135 let controller = Controller::load()?;
136 Ok(active_child_session_ids(&controller.state, &session_id))
137 }
138 })
139 .await?;
140 for child_id in children {
141 Box::pin(self.suspend_session(child_id.clone()))
142 .await
143 .with_context(|| {
144 format!("suspend sub-agent {child_id} before discarding parent changes")
145 })?;
146 }
147 let operation_session_id = session_id.clone();
148 let result = self
149 .run_lifecycle(
150 operation_session_id,
151 LifecycleKind::ForceStop,
152 move |state, session_id, cancelled| async move {
153 blocking(move || {
154 let mut controller = Controller::load()?;
155 let executor = DaemonStageReportingExecutor::new(
156 CancellableProcessExecutor::new(cancelled),
157 state,
158 session_id.clone(),
159 );
160 ensure!(
161 controller
162 .state
163 .sessions
164 .get(&session_id)
165 .and_then(|s| s.checkpoint.as_ref())
166 == Some(&checkpoint),
167 "the recovery copy changed; review it before discarding changes"
168 );
169 let deferred = controller.force_stop(&session_id, &executor)?;
170 Ok(if deferred {
171 DaemonLifecycleResult::DeferredCleanup
172 } else {
173 DaemonLifecycleResult::Done
174 })
175 })
176 .await
177 },
178 )
179 .await?;
180 let _ = result; Ok(())
182 }
183
184 pub(super) async fn destroy_stopped_session(
185 self: &Arc<Self>,
186 session_id: String,
187 branch: BranchDisposition,
188 ) -> Result<()> {
189 self.tear_down_stopped_session(
190 session_id,
191 LifecycleKind::DestroyStopped,
192 branch,
193 CheckoutDisposition::Remove,
194 )
195 .await
196 .map(|_| ())
197 }
198
199 pub(crate) async fn discard_lost_session(self: &Arc<Self>, session_id: String) {
209 let short = mj_core::state::short_id(&session_id).to_owned();
210 match self
211 .tear_down_stopped_session(
212 session_id.clone(),
213 LifecycleKind::DestroyStopped,
214 BranchDisposition::Keep,
215 CheckoutDisposition::KeepWhenDirty,
216 )
217 .await
218 {
219 Ok(retained) => {
220 let mut text = format!(
221 "Session {short} was lost because its managed target no longer exists; its record was removed."
222 );
223 if let Some(path) = retained {
224 text.push_str(&format!(
225 " Its checkout has uncommitted changes, so it was kept at {}.",
226 path.display()
227 ));
228 }
229 self.push_notice(&session_id, text);
230 }
231 Err(error) => {
232 tracing::warn!(
233 %session_id,
234 error = format!("{error:#}"),
235 "could not discard the record of a lost session"
236 );
237 self.push_notice(
238 &session_id,
239 format!(
240 "Session {short} was lost, but its record could not be removed: {error:#}"
241 ),
242 );
243 }
244 }
245 }
246
247 pub(crate) async fn archive_aged_sessions(
253 self: &Arc<Self>,
254 older_than_days: u32,
255 ) -> Result<usize> {
256 self.wiki()
257 .sync_now(true)
258 .await
259 .context("sync SessionWiki before archiving stopped sessions")?;
260 let candidates = blocking(move || {
261 let controller = Controller::load()?;
262 Ok(crate::sessionwiki::sessions_ready_to_archive(
263 &controller.state.sessions,
264 &controller.state.subagents,
265 chrono::Utc::now(),
266 older_than_days,
267 ))
268 })
269 .await
270 .context("select the stopped sessions old enough to archive")?;
271 if candidates.is_empty() {
272 return Ok(0);
273 }
274 let indexed = blocking({
275 let candidates = candidates.clone();
276 move || crate::sessionwiki::indexed_with_messages(&candidates)
277 })
278 .await
279 .context("check the SessionWiki index before archiving")?;
280 let mut archived = 0;
281 for session_id in candidates {
282 if !indexed.contains(&session_id) {
283 tracing::warn!(
284 %session_id,
285 "SessionWiki holds no conversation for this stopped session; keeping it"
286 );
287 continue;
288 }
289 match self.archive_stopped_session(session_id.clone()).await {
290 Ok(()) => {
291 archived += 1;
292 tracing::info!(
293 %session_id,
294 older_than_days,
295 "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
296 );
297 }
298 Err(error) => tracing::warn!(
299 %session_id,
300 error = %format!("{error:#}"),
301 "could not archive a stopped session"
302 ),
303 }
304 }
305 if archived > 0 {
306 self.wiki().request_sync(false);
308 }
309 Ok(archived)
310 }
311
312 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
318 self.tear_down_stopped_session(
319 session_id,
320 LifecycleKind::ArchiveStopped,
321 BranchDisposition::DeleteIfMerged,
322 CheckoutDisposition::Remove,
323 )
324 .await
325 .map(|_| ())
326 }
327
328 async fn tear_down_stopped_session(
330 self: &Arc<Self>,
331 session_id: String,
332 kind: LifecycleKind,
333 branch: BranchDisposition,
334 checkout: CheckoutDisposition,
335 ) -> Result<Option<PathBuf>> {
336 let children = blocking({
337 let session_id = session_id.clone();
338 move || {
339 Ok(crate::database::list_subagents(&session_id)?
340 .into_iter()
341 .map(|child| child.child_session_id)
342 .collect::<Vec<_>>())
343 }
344 })
345 .await?;
346 for child_id in children {
347 Box::pin(self.force_destroy_session(child_id.clone(), BranchDisposition::Keep))
350 .await
351 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
352 }
353 self.wait_for_deferred_cleanup(&session_id).await?;
354 let exists = blocking({
355 let session_id = session_id.clone();
356 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
357 })
358 .await?;
359 if !exists {
360 return Ok(None);
361 }
362 let retained = Arc::new(Mutex::new(None));
365 self.run_lifecycle(session_id, kind, {
366 let retained = retained.clone();
367 move |state, session_id, cancelled| async move {
368 blocking(move || {
369 let mut controller = Controller::load()?;
370 let executor = DaemonStageReportingExecutor::new(
371 CancellableProcessExecutor::new(cancelled),
372 state,
373 session_id.clone(),
374 );
375 let kept = controller.destroy_session_controlled_with_checkout(
376 &session_id,
377 &executor,
378 branch,
379 checkout,
380 )?;
381 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
382 Ok(DaemonLifecycleResult::Done)
383 })
384 .await
385 }
386 })
387 .await?;
388 let kept = retained
389 .lock()
390 .unwrap_or_else(PoisonError::into_inner)
391 .clone();
392 Ok(kept)
393 }
394
395 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
406 let mut result = {
407 let lifecycle = self
408 .lifecycle
409 .lock()
410 .unwrap_or_else(PoisonError::into_inner);
411 let Some(active) = lifecycle.get(session_id) else {
412 return Ok(());
413 };
414 if !active.result.borrow().is_none() {
415 return Ok(());
416 }
417 active.request_cancel();
418 active.result.clone()
419 };
420 let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
421 loop {
422 if result.borrow().is_some() {
423 return Ok(());
424 }
425 if result.changed().await.is_err() {
426 return Err(());
427 }
428 }
429 })
430 .await;
431 match finished {
432 Ok(Ok(())) => Ok(()),
435 Ok(Err(())) => bail!(
436 "daemon lifecycle operation stopped without a result for session {session_id}"
437 ),
438 Err(_) => bail!(
439 "session {session_id} still has an operation that did not stop after cancellation; try again"
440 ),
441 }
442 }
443
444 pub async fn force_destroy_session(
449 self: &Arc<Self>,
450 session_id: String,
451 branch: BranchDisposition,
452 ) -> Result<()> {
453 let children = blocking({
454 let session_id = session_id.clone();
455 move || {
456 Ok(crate::database::list_subagents(&session_id)?
457 .into_iter()
458 .map(|child| child.child_session_id)
459 .collect::<Vec<_>>())
460 }
461 })
462 .await?;
463 for child_id in children {
464 Box::pin(self.force_destroy_session(child_id.clone(), BranchDisposition::Keep))
466 .await
467 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
468 }
469 self.preempt_active_lifecycle(&session_id).await?;
470 let exists = blocking({
471 let session_id = session_id.clone();
472 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
473 })
474 .await?;
475 if !exists {
476 return Ok(());
477 }
478 self.run_lifecycle(
479 session_id,
480 LifecycleKind::ForceDestroy,
481 move |state, session_id, cancelled| async move {
482 let _recovery_reservation = tokio::task::spawn_blocking({
483 let observer = state.recovery_observer.clone();
484 let session_id = session_id.clone();
485 let cancelled = cancelled.clone();
486 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
487 })
488 .await
489 .context("reserve recovery for daemon force-destroy task")??;
490 blocking({
491 let session_id = session_id.clone();
492 move || {
493 let mut controller = Controller::load()?;
494 let executor = DaemonStageReportingExecutor::new(
495 CancellableProcessExecutor::new(cancelled),
496 state,
497 session_id.clone(),
498 );
499 controller.force_destroy_session(&session_id, &executor, branch)?;
500 crate::controller::move_session::release_move_queue_hold(&session_id);
501 Ok(DaemonLifecycleResult::Done)
502 }
503 })
504 .await
505 },
506 )
507 .await?;
508 Ok(())
509 }
510
511 pub async fn force_delete_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
520 ensure!(
521 !self.workspace_has_active_resume(&workspace_id),
522 "workspace has a session resume in progress"
523 );
524 let sessions = blocking({
525 let workspace_id = workspace_id.clone();
526 move || {
527 let controller = Controller::load()?;
528 Ok(active_sessions_for_force_destruction(
529 &controller,
530 &workspace_id,
531 ))
532 }
533 })
534 .await?;
535 for (index, session_id) in sessions.iter().enumerate() {
536 if let Err(error) = self
539 .force_destroy_session(session_id.clone(), BranchDisposition::Keep)
540 .await
541 {
542 let remaining = sessions.len() - index - 1;
543 bail!(
544 "force-destroying session {session_id} failed: {error:#}; \
545 {remaining} session(s) in the workspace remain"
546 );
547 }
548 }
549 blocking({
550 let workspace_id = workspace_id.clone();
551 move || crate::database::force_delete_workspace(&workspace_id)
552 })
553 .await?;
554 refresh_runtime_workspaces(self).await?;
555 Ok(())
556 }
557}