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(_) | DaemonLifecycleResult::Park(_) => {
106 unreachable!("resume cannot return a move outcome")
107 }
108 DaemonLifecycleResult::DeferredCleanup => {
109 unreachable!("session resume cannot schedule target cleanup")
110 }
111 }
112 blocking(move || {
113 if let Some(mut operation) =
114 crate::database::load_move_operation(&operation_session_id)?
115 && !operation.queue_admission_started
116 {
117 operation.phase = mj_core::state::MovePhase::Cancelled;
118 operation.queue_admission_finished = true;
119 operation.updated_at = chrono::Utc::now().to_rfc3339();
120 operation.error = Some("Recovered through an explicit Resume operation".into());
121 crate::database::save_move_operation(&operation)?;
122 }
123 Ok(())
124 })
125 .await?;
126 Ok(())
127 }
128
129 pub(super) async fn discard_since_checkpoint(
130 self: &Arc<Self>,
131 session_id: String,
132 checkpoint: mj_core::state::CheckpointMetadata,
133 ) -> Result<()> {
134 let operation_session_id = session_id.clone();
135 let result = self
136 .run_lifecycle(
137 operation_session_id,
138 LifecycleKind::ForceStop,
139 move |state, session_id, cancelled| async move {
140 blocking({
142 let session_id = session_id.clone();
143 let checkpoint = checkpoint.clone();
144 move || {
145 ensure!(
146 Controller::load()?
147 .state
148 .sessions
149 .get(&session_id)
150 .and_then(|s| s.checkpoint.as_ref())
151 == Some(&checkpoint),
152 "the recovery copy changed; review it before discarding changes"
153 );
154 Ok(())
155 }
156 })
157 .await?;
158 state.stop_subagents_for_suspend(&session_id).await?;
162 blocking(move || {
163 let mut controller = Controller::load()?;
164 let executor = DaemonStageReportingExecutor::new(
165 CancellableProcessExecutor::new(cancelled),
166 state,
167 session_id.clone(),
168 );
169 ensure!(
170 controller
171 .state
172 .sessions
173 .get(&session_id)
174 .and_then(|s| s.checkpoint.as_ref())
175 == Some(&checkpoint),
176 "the recovery copy changed; review it before discarding changes"
177 );
178 let deferred = controller.force_stop(&session_id, &executor)?;
179 Ok(if deferred {
180 DaemonLifecycleResult::DeferredCleanup
181 } else {
182 DaemonLifecycleResult::Done
183 })
184 })
185 .await
186 },
187 )
188 .await;
189 if result.is_err() {
190 self.tell_live_parent_about_stopped_subagents(&session_id)
191 .await;
192 }
193 let _ = result?; Ok(())
195 }
196
197 pub(super) async fn destroy_stopped_session(
198 self: &Arc<Self>,
199 session_id: String,
200 branch: BranchDisposition,
201 ) -> Result<()> {
202 self.tear_down_stopped_session(
203 session_id,
204 LifecycleKind::DestroyStopped,
205 branch,
206 CheckoutDisposition::Remove,
207 )
208 .await
209 .map(|_| ())
210 }
211
212 pub(crate) async fn discard_lost_session(self: &Arc<Self>, session_id: String) {
222 let short = mj_core::state::short_id(&session_id).to_owned();
223 let owned_clone = blocking({
224 let session_id = session_id.clone();
225 move || {
226 let controller = Controller::load()?;
227 let checkout = match controller.state.checkout(&session_id)? {
228 mj_core::state::Checkout::ManagedWorktree { worktree, .. }
229 if worktree.kind == mj_core::state::ManagedCheckoutKind::Clone =>
230 {
231 worktree
232 }
233 _ => return Ok(None),
234 };
235 if !controller.state.sessions.contains_key(&session_id) {
236 return Ok(None);
237 }
238 if crate::controller::path_exists_on_managed_target(
239 &crate::targets::ProcessExecutor,
240 &checkout.target,
241 &checkout.worktree_root,
242 )? {
243 Ok(Some(checkout.worktree_root.clone()))
244 } else {
245 Ok(None)
246 }
247 }
248 })
249 .await;
250 match owned_clone {
251 Ok(Some(path)) => {
252 self.push_notice(&session_id, format!(
253 "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
254 path.display(),
255 ));
256 return;
257 }
258 Err(error) => {
259 tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
260 return;
261 }
262 Ok(None) => {}
263 }
264 match self
265 .tear_down_stopped_session(
266 session_id.clone(),
267 LifecycleKind::DestroyStopped,
268 BranchDisposition::Keep,
269 CheckoutDisposition::KeepWhenDirty,
270 )
271 .await
272 {
273 Ok(retained) => {
274 let mut text = format!(
275 "Session {short} was lost because its managed target no longer exists; its record was removed."
276 );
277 if let Some(path) = retained {
278 text.push_str(&format!(
279 " Its checkout has uncommitted changes, so it was kept at {}.",
280 path.display()
281 ));
282 }
283 self.push_notice(&session_id, text);
284 }
285 Err(error) => {
286 tracing::warn!(
287 %session_id,
288 error = format!("{error:#}"),
289 "could not discard the record of a lost session"
290 );
291 self.push_notice(
292 &session_id,
293 format!(
294 "Session {short} was lost, but its record could not be removed: {error:#}"
295 ),
296 );
297 }
298 }
299 }
300
301 pub(crate) async fn archive_aged_sessions(
307 self: &Arc<Self>,
308 older_than_days: u32,
309 ) -> Result<usize> {
310 self.wiki()
311 .sync_now(true)
312 .await
313 .context("sync SessionWiki before archiving stopped sessions")?;
314 let Ok(_work) = crate::upgrade::activity_unless_draining("SessionWiki archive") else {
318 tracing::debug!(
319 "a daemon upgrade is waiting; the SessionWiki archive pass waits for the next daemon"
320 );
321 return Ok(0);
322 };
323 let refreshable = blocking(move || {
326 let controller = Controller::load()?;
327 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
328 let mut sessions = controller
329 .state
330 .sessions
331 .values()
332 .filter(|session| session.state == SessionState::Stopped)
333 .filter(|session| {
334 chrono::DateTime::parse_from_rfc3339(&session.updated_at)
335 .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
336 })
337 .filter(|session| {
338 session.publication.as_ref().is_none_or(|evidence| {
339 evidence.state != mj_core::state::PublicationState::Published
340 && !evidence.dirty
341 && !evidence.stashed
342 })
343 })
344 .filter_map(|session| {
345 let mj_core::state::Checkout::ManagedWorktree { worktree, .. } =
346 controller.state.checkout(&session.id).ok()?
347 else {
348 return None;
349 };
350 (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone)
351 .then(|| (session.clone(), worktree.clone()))
352 })
353 .collect::<Vec<_>>();
354 sessions.sort_by(|(a, _), (b, _)| {
355 a.publication
356 .as_ref()
357 .map(|e| &e.checked_at)
358 .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
359 });
360 sessions.truncate(4);
361 Ok(sessions)
362 })
363 .await?;
364 let mut checks = tokio::task::JoinSet::new();
365 for (session, checkout) in refreshable {
366 checks.spawn_blocking(move || {
367 let assessment =
368 crate::controller::publication::refresh_stopped_clone_publication_for_worktree(
369 &session, &checkout,
370 );
371 (session.id, assessment)
372 });
373 }
374 let mut refreshed = false;
375 while let Some(done) = checks.join_next().await {
376 let (id, assessment) = done.context("publication refresh task failed")?;
377 if let Some(assessment) = assessment {
378 refreshed |= blocking(move || {
379 crate::database::set_publication_assessment_if_current(&id, &assessment)
380 })
381 .await?;
382 }
383 }
384 if refreshed {
385 self.reload_controller().await?;
386 self.publish_revision();
387 }
388 let (candidates, children) = blocking(move || {
389 let controller = Controller::load()?;
390 let candidates = crate::sessionwiki::sessions_ready_to_archive_from_state(
391 &controller.state,
392 chrono::Utc::now(),
393 older_than_days,
394 );
395 let children: std::collections::BTreeSet<String> = candidates
396 .iter()
397 .filter(|id| controller.state.is_subagent_session(id))
398 .cloned()
399 .collect();
400 Ok((candidates, children))
401 })
402 .await
403 .context("select the stopped sessions old enough to archive")?;
404 if candidates.is_empty() {
405 return Ok(0);
406 }
407 let indexed = blocking({
408 let candidates = candidates
409 .iter()
410 .filter(|id| !children.contains(*id))
411 .cloned()
412 .collect::<Vec<_>>();
413 move || crate::sessionwiki::indexed_with_messages(&candidates)
414 })
415 .await
416 .context("check the SessionWiki index before archiving")?;
417 let mut archived = 0;
418 for session_id in candidates {
419 if !children.contains(&session_id) && !indexed.contains(&session_id) {
420 tracing::warn!(
421 %session_id,
422 "SessionWiki holds no conversation for this stopped session; keeping it"
423 );
424 continue;
425 }
426 match self.archive_stopped_session(session_id.clone()).await {
427 Ok(()) => {
428 archived += 1;
429 tracing::info!(
430 %session_id,
431 older_than_days,
432 child = children.contains(&session_id),
433 "archived a stopped session; top-level conversations are retained in SessionWiki"
434 );
435 }
436 Err(error) => tracing::warn!(
437 %session_id,
438 error = %format!("{error:#}"),
439 "could not archive a stopped session"
440 ),
441 }
442 }
443 if archived > 0 {
444 self.wiki().request_sync(false);
446 }
447 Ok(archived)
448 }
449
450 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
456 self.tear_down_stopped_session(
457 session_id,
458 LifecycleKind::ArchiveStopped,
459 BranchDisposition::DeleteIfMerged,
460 CheckoutDisposition::Remove,
461 )
462 .await
463 .map(|_| ())
464 }
465
466 async fn tear_down_stopped_session(
468 self: &Arc<Self>,
469 session_id: String,
470 kind: LifecycleKind,
471 branch: BranchDisposition,
472 checkout: CheckoutDisposition,
473 ) -> Result<Option<PathBuf>> {
474 self.index_before_destroy(&session_id).await;
477 let children = blocking({
478 let session_id = session_id.clone();
479 move || {
480 Ok(crate::database::list_subagents(&session_id)?
481 .into_iter()
482 .map(|child| child.child_session_id)
483 .collect::<Vec<_>>())
484 }
485 })
486 .await?;
487 for child_id in children {
488 Box::pin(self.force_destroy_indexed_session(
491 child_id.clone(),
492 BranchDisposition::Keep,
493 LifecycleKind::ForceDestroy,
494 ))
495 .await
496 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
497 }
498 self.wait_for_deferred_cleanup(&session_id).await?;
499 let exists = blocking({
500 let session_id = session_id.clone();
501 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
502 })
503 .await?;
504 if !exists {
505 return Ok(None);
506 }
507 let retained = Arc::new(Mutex::new(None));
510 self.run_lifecycle(session_id, kind, {
511 let retained = retained.clone();
512 move |state, session_id, cancelled| async move {
513 blocking(move || {
514 let mut controller = Controller::load()?;
515 let executor = DaemonStageReportingExecutor::new(
516 CancellableProcessExecutor::new(cancelled),
517 state,
518 session_id.clone(),
519 );
520 let kept = controller.destroy_session_controlled_with_checkout(
521 &session_id,
522 &executor,
523 branch,
524 checkout,
525 )?;
526 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
527 Ok(DaemonLifecycleResult::Done)
528 })
529 .await
530 }
531 })
532 .await?;
533 let kept = retained
534 .lock()
535 .unwrap_or_else(PoisonError::into_inner)
536 .clone();
537 Ok(kept)
538 }
539
540 pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
548 use crate::sessionwiki::IndexedBeforeDestroy;
549 let outcome = self
550 .wiki()
551 .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
552 .await;
553 match outcome {
554 IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
555 %session_id,
556 reason,
557 "destroying a session without indexing it in SessionWiki"
558 ),
559 IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
560 %session_id,
561 %reason,
562 "could not index a session in SessionWiki before destroying it"
563 ),
564 outcome => tracing::debug!(
565 %session_id,
566 ?outcome,
567 "indexed a session in SessionWiki before destroying it"
568 ),
569 }
570 }
571
572 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
583 let (mut result, tearing_down, cancelled_restart) = {
584 let mut lifecycle_owner = self.owner();
585 let lifecycle = &mut lifecycle_owner.lifecycle;
586 let Some(active) = lifecycle.get_mut(session_id) else {
587 return Ok(());
588 };
589 if !active.is_running() {
590 return Ok(());
591 }
592 let tearing_down = active.kind.is_teardown();
593 if !tearing_down {
599 active.request_cancel();
600 }
601 (
602 active.result.clone(),
603 tearing_down,
604 active.kind == LifecycleKind::Restart && !tearing_down,
605 )
606 };
607 if cancelled_restart {
608 blocking({
609 let session_id = session_id.to_owned();
610 move || crate::database::cancel_session_restart(&session_id)
611 })
612 .await?;
613 }
614 let mut finished = wait_for_lifecycle_result(&mut result).await;
615 if finished.is_err() && tearing_down {
616 {
617 let mut lifecycle_owner = self.owner();
618 if let Some(active) = lifecycle_owner.lifecycle.get_mut(session_id) {
619 active.request_cancel();
620 }
621 }
622 finished = wait_for_lifecycle_result(&mut result).await;
623 }
624 match finished {
625 Ok(Ok(())) => Ok(()),
628 Ok(Err(())) => bail!(
629 "daemon lifecycle operation stopped without a result for session {session_id}"
630 ),
631 Err(_) => bail!(
632 "session {session_id} still has an operation that did not stop after cancellation; try again"
633 ),
634 }
635 }
636
637 pub async fn force_destroy_session(
642 self: &Arc<Self>,
643 session_id: String,
644 branch: BranchDisposition,
645 ) -> Result<()> {
646 let _upgrade_work = crate::upgrade::destroy_activity(&session_id)?;
649 self.index_before_destroy(&session_id).await;
650 self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
651 .await
652 }
653
654 pub(super) async fn force_destroy_indexed_session(
659 self: &Arc<Self>,
660 session_id: String,
661 branch: BranchDisposition,
662 kind: LifecycleKind,
663 ) -> Result<()> {
664 let children = blocking({
665 let session_id = session_id.clone();
666 move || {
667 Ok(crate::database::list_subagents(&session_id)?
668 .into_iter()
669 .map(|child| child.child_session_id)
670 .collect::<Vec<_>>())
671 }
672 })
673 .await?;
674 for child_id in children {
675 Box::pin(self.force_destroy_indexed_session(
677 child_id.clone(),
678 BranchDisposition::Keep,
679 kind,
680 ))
681 .await
682 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
683 }
684 self.preempt_active_lifecycle(&session_id).await?;
685 let exists = blocking({
686 let session_id = session_id.clone();
687 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
688 })
689 .await?;
690 if !exists {
691 return Ok(());
692 }
693 self.run_lifecycle(
694 session_id,
695 kind,
696 move |state, session_id, cancelled| async move {
697 let _recovery_reservation = tokio::task::spawn_blocking({
698 let observer = state.recovery_observer.clone();
699 let session_id = session_id.clone();
700 let cancelled = cancelled.clone();
701 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
702 })
703 .await
704 .context("reserve recovery for daemon force-destroy task")??;
705 blocking({
706 let session_id = session_id.clone();
707 move || {
708 let mut controller = Controller::load()?;
709 let executor = DaemonStageReportingExecutor::new(
710 CancellableProcessExecutor::new(cancelled),
711 state,
712 session_id.clone(),
713 );
714 controller.force_destroy_session(&session_id, &executor, branch)?;
715 crate::controller::move_session::release_move_queue_hold(&session_id);
716 Ok(DaemonLifecycleResult::Done)
717 }
718 })
719 .await
720 },
721 )
722 .await?;
723 Ok(())
724 }
725}
726
727async fn wait_for_lifecycle_result(
730 result: &mut LifecycleWatch,
731) -> std::result::Result<std::result::Result<(), ()>, tokio::time::error::Elapsed> {
732 tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
733 loop {
734 if result.borrow().is_some() {
735 return Ok(());
736 }
737 if result.changed().await.is_err() {
738 return Err(());
739 }
740 }
741 })
742 .await
743}