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 operation_session_id = session_id.clone();
133 let result = self
134 .run_lifecycle(
135 operation_session_id,
136 LifecycleKind::ForceStop,
137 move |state, session_id, cancelled| async move {
138 blocking({
140 let session_id = session_id.clone();
141 let checkpoint = checkpoint.clone();
142 move || {
143 ensure!(
144 Controller::load()?
145 .state
146 .sessions
147 .get(&session_id)
148 .and_then(|s| s.checkpoint.as_ref())
149 == Some(&checkpoint),
150 "the recovery copy changed; review it before discarding changes"
151 );
152 Ok(())
153 }
154 })
155 .await?;
156 state.stop_subagents_for_suspend(&session_id).await?;
160 blocking(move || {
161 let mut controller = Controller::load()?;
162 let executor = DaemonStageReportingExecutor::new(
163 CancellableProcessExecutor::new(cancelled),
164 state,
165 session_id.clone(),
166 );
167 ensure!(
168 controller
169 .state
170 .sessions
171 .get(&session_id)
172 .and_then(|s| s.checkpoint.as_ref())
173 == Some(&checkpoint),
174 "the recovery copy changed; review it before discarding changes"
175 );
176 let deferred = controller.force_stop(&session_id, &executor)?;
177 Ok(if deferred {
178 DaemonLifecycleResult::DeferredCleanup
179 } else {
180 DaemonLifecycleResult::Done
181 })
182 })
183 .await
184 },
185 )
186 .await;
187 if result.is_err() {
188 self.tell_live_parent_about_stopped_subagents(&session_id)
189 .await;
190 }
191 let _ = result?; Ok(())
193 }
194
195 pub(super) async fn destroy_stopped_session(
196 self: &Arc<Self>,
197 session_id: String,
198 branch: BranchDisposition,
199 ) -> Result<()> {
200 self.tear_down_stopped_session(
201 session_id,
202 LifecycleKind::DestroyStopped,
203 branch,
204 CheckoutDisposition::Remove,
205 )
206 .await
207 .map(|_| ())
208 }
209
210 pub(crate) async fn discard_lost_session(self: &Arc<Self>, session_id: String) {
220 let short = mj_core::state::short_id(&session_id).to_owned();
221 let owned_clone = blocking({
222 let session_id = session_id.clone();
223 move || {
224 let controller = Controller::load()?;
225 let Some(checkout) = controller
226 .state
227 .sessions
228 .get(&session_id)
229 .and_then(|session| session.managed_worktree.as_ref())
230 .filter(|checkout| checkout.kind == mj_core::state::ManagedCheckoutKind::Clone)
231 else {
232 return Ok(None);
233 };
234 if crate::controller::path_exists_on_managed_target(
235 &crate::targets::ProcessExecutor,
236 &checkout.target,
237 &checkout.worktree_root,
238 )? {
239 Ok(Some(checkout.worktree_root.clone()))
240 } else {
241 Ok(None)
242 }
243 }
244 })
245 .await;
246 match owned_clone {
247 Ok(Some(path)) => {
248 self.push_notice(&session_id, format!(
249 "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
250 path.display(),
251 ));
252 return;
253 }
254 Err(error) => {
255 tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
256 return;
257 }
258 Ok(None) => {}
259 }
260 match self
261 .tear_down_stopped_session(
262 session_id.clone(),
263 LifecycleKind::DestroyStopped,
264 BranchDisposition::Keep,
265 CheckoutDisposition::KeepWhenDirty,
266 )
267 .await
268 {
269 Ok(retained) => {
270 let mut text = format!(
271 "Session {short} was lost because its managed target no longer exists; its record was removed."
272 );
273 if let Some(path) = retained {
274 text.push_str(&format!(
275 " Its checkout has uncommitted changes, so it was kept at {}.",
276 path.display()
277 ));
278 }
279 self.push_notice(&session_id, text);
280 }
281 Err(error) => {
282 tracing::warn!(
283 %session_id,
284 error = format!("{error:#}"),
285 "could not discard the record of a lost session"
286 );
287 self.push_notice(
288 &session_id,
289 format!(
290 "Session {short} was lost, but its record could not be removed: {error:#}"
291 ),
292 );
293 }
294 }
295 }
296
297 pub(crate) async fn archive_aged_sessions(
303 self: &Arc<Self>,
304 older_than_days: u32,
305 ) -> Result<usize> {
306 self.wiki()
307 .sync_now(true)
308 .await
309 .context("sync SessionWiki before archiving stopped sessions")?;
310 let refreshable = blocking(move || {
313 let controller = Controller::load()?;
314 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
315 let mut sessions = controller
316 .state
317 .sessions
318 .values()
319 .filter(|session| session.state == SessionState::Stopped)
320 .filter(|session| {
321 chrono::DateTime::parse_from_rfc3339(&session.updated_at)
322 .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
323 })
324 .filter(|session| {
325 session.managed_worktree.as_ref().is_some_and(|owned| {
326 owned.kind == mj_core::state::ManagedCheckoutKind::Clone
327 })
328 })
329 .filter(|session| {
330 session.publication.as_ref().is_none_or(|evidence| {
331 evidence.state != mj_core::state::PublicationState::Published
332 && !evidence.dirty
333 && !evidence.stashed
334 })
335 })
336 .cloned()
337 .collect::<Vec<_>>();
338 sessions.sort_by(|a, b| {
339 a.publication
340 .as_ref()
341 .map(|e| &e.checked_at)
342 .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
343 });
344 sessions.truncate(4);
345 Ok(sessions)
346 })
347 .await?;
348 let mut checks = tokio::task::JoinSet::new();
349 for session in refreshable {
350 checks.spawn_blocking(move || {
351 let assessment =
352 crate::controller::publication::refresh_stopped_clone_publication(&session);
353 (session.id, assessment)
354 });
355 }
356 let mut refreshed = false;
357 while let Some(done) = checks.join_next().await {
358 let (id, assessment) = done.context("publication refresh task failed")?;
359 if let Some(assessment) = assessment {
360 refreshed |= blocking(move || {
361 crate::database::set_publication_assessment_if_current(&id, &assessment)
362 })
363 .await?;
364 }
365 }
366 if refreshed {
367 self.reload_controller().await?;
368 self.publish_revision();
369 }
370 let candidates = blocking(move || {
371 let controller = Controller::load()?;
372 Ok(crate::sessionwiki::sessions_ready_to_archive(
373 &controller.state.sessions,
374 &controller.state.subagents,
375 chrono::Utc::now(),
376 older_than_days,
377 ))
378 })
379 .await
380 .context("select the stopped sessions old enough to archive")?;
381 if candidates.is_empty() {
382 return Ok(0);
383 }
384 let indexed = blocking({
385 let candidates = candidates.clone();
386 move || crate::sessionwiki::indexed_with_messages(&candidates)
387 })
388 .await
389 .context("check the SessionWiki index before archiving")?;
390 let mut archived = 0;
391 for session_id in candidates {
392 if !indexed.contains(&session_id) {
393 tracing::warn!(
394 %session_id,
395 "SessionWiki holds no conversation for this stopped session; keeping it"
396 );
397 continue;
398 }
399 match self.archive_stopped_session(session_id.clone()).await {
400 Ok(()) => {
401 archived += 1;
402 tracing::info!(
403 %session_id,
404 older_than_days,
405 "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
406 );
407 }
408 Err(error) => tracing::warn!(
409 %session_id,
410 error = %format!("{error:#}"),
411 "could not archive a stopped session"
412 ),
413 }
414 }
415 if archived > 0 {
416 self.wiki().request_sync(false);
418 }
419 Ok(archived)
420 }
421
422 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
428 self.tear_down_stopped_session(
429 session_id,
430 LifecycleKind::ArchiveStopped,
431 BranchDisposition::DeleteIfMerged,
432 CheckoutDisposition::Remove,
433 )
434 .await
435 .map(|_| ())
436 }
437
438 async fn tear_down_stopped_session(
440 self: &Arc<Self>,
441 session_id: String,
442 kind: LifecycleKind,
443 branch: BranchDisposition,
444 checkout: CheckoutDisposition,
445 ) -> Result<Option<PathBuf>> {
446 self.index_before_destroy(&session_id).await;
449 let children = blocking({
450 let session_id = session_id.clone();
451 move || {
452 Ok(crate::database::list_subagents(&session_id)?
453 .into_iter()
454 .map(|child| child.child_session_id)
455 .collect::<Vec<_>>())
456 }
457 })
458 .await?;
459 for child_id in children {
460 Box::pin(self.force_destroy_indexed_session(
463 child_id.clone(),
464 BranchDisposition::Keep,
465 LifecycleKind::ForceDestroy,
466 ))
467 .await
468 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
469 }
470 self.wait_for_deferred_cleanup(&session_id).await?;
471 let exists = blocking({
472 let session_id = session_id.clone();
473 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
474 })
475 .await?;
476 if !exists {
477 return Ok(None);
478 }
479 let retained = Arc::new(Mutex::new(None));
482 self.run_lifecycle(session_id, kind, {
483 let retained = retained.clone();
484 move |state, session_id, cancelled| async move {
485 blocking(move || {
486 let mut controller = Controller::load()?;
487 let executor = DaemonStageReportingExecutor::new(
488 CancellableProcessExecutor::new(cancelled),
489 state,
490 session_id.clone(),
491 );
492 let kept = controller.destroy_session_controlled_with_checkout(
493 &session_id,
494 &executor,
495 branch,
496 checkout,
497 )?;
498 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
499 Ok(DaemonLifecycleResult::Done)
500 })
501 .await
502 }
503 })
504 .await?;
505 let kept = retained
506 .lock()
507 .unwrap_or_else(PoisonError::into_inner)
508 .clone();
509 Ok(kept)
510 }
511
512 pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
521 use crate::sessionwiki::IndexedBeforeDestroy;
522 let outcome = self
523 .wiki()
524 .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
525 .await;
526 match outcome {
527 IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
528 %session_id,
529 reason,
530 "destroying a session without indexing it in SessionWiki"
531 ),
532 IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
533 %session_id,
534 %reason,
535 "could not index a session in SessionWiki before destroying it"
536 ),
537 outcome => tracing::debug!(
538 %session_id,
539 ?outcome,
540 "indexed a session in SessionWiki before destroying it"
541 ),
542 }
543 }
544
545 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
556 let mut result = {
557 let lifecycle = self
558 .lifecycle
559 .lock()
560 .unwrap_or_else(PoisonError::into_inner);
561 let Some(active) = lifecycle.get(session_id) else {
562 return Ok(());
563 };
564 if !active.result.borrow().is_none() {
565 return Ok(());
566 }
567 active.request_cancel();
568 active.result.clone()
569 };
570 let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
571 loop {
572 if result.borrow().is_some() {
573 return Ok(());
574 }
575 if result.changed().await.is_err() {
576 return Err(());
577 }
578 }
579 })
580 .await;
581 match finished {
582 Ok(Ok(())) => Ok(()),
585 Ok(Err(())) => bail!(
586 "daemon lifecycle operation stopped without a result for session {session_id}"
587 ),
588 Err(_) => bail!(
589 "session {session_id} still has an operation that did not stop after cancellation; try again"
590 ),
591 }
592 }
593
594 pub async fn force_destroy_session(
599 self: &Arc<Self>,
600 session_id: String,
601 branch: BranchDisposition,
602 ) -> Result<()> {
603 self.index_before_destroy(&session_id).await;
604 self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
605 .await
606 }
607
608 pub(super) async fn force_destroy_indexed_session(
613 self: &Arc<Self>,
614 session_id: String,
615 branch: BranchDisposition,
616 kind: LifecycleKind,
617 ) -> Result<()> {
618 let children = blocking({
619 let session_id = session_id.clone();
620 move || {
621 Ok(crate::database::list_subagents(&session_id)?
622 .into_iter()
623 .map(|child| child.child_session_id)
624 .collect::<Vec<_>>())
625 }
626 })
627 .await?;
628 for child_id in children {
629 Box::pin(self.force_destroy_indexed_session(
631 child_id.clone(),
632 BranchDisposition::Keep,
633 kind,
634 ))
635 .await
636 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
637 }
638 self.preempt_active_lifecycle(&session_id).await?;
639 let exists = blocking({
640 let session_id = session_id.clone();
641 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
642 })
643 .await?;
644 if !exists {
645 return Ok(());
646 }
647 self.run_lifecycle(
648 session_id,
649 kind,
650 move |state, session_id, cancelled| async move {
651 let _recovery_reservation = tokio::task::spawn_blocking({
652 let observer = state.recovery_observer.clone();
653 let session_id = session_id.clone();
654 let cancelled = cancelled.clone();
655 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
656 })
657 .await
658 .context("reserve recovery for daemon force-destroy task")??;
659 blocking({
660 let session_id = session_id.clone();
661 move || {
662 let mut controller = Controller::load()?;
663 let executor = DaemonStageReportingExecutor::new(
664 CancellableProcessExecutor::new(cancelled),
665 state,
666 session_id.clone(),
667 );
668 controller.force_destroy_session(&session_id, &executor, branch)?;
669 crate::controller::move_session::release_move_queue_hold(&session_id);
670 Ok(DaemonLifecycleResult::Done)
671 }
672 })
673 .await
674 },
675 )
676 .await?;
677 Ok(())
678 }
679}