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 Some(checkout) = controller
228 .state
229 .sessions
230 .get(&session_id)
231 .and_then(|session| session.managed_worktree.as_ref())
232 .filter(|checkout| checkout.kind == mj_core::state::ManagedCheckoutKind::Clone)
233 else {
234 return Ok(None);
235 };
236 if crate::controller::path_exists_on_managed_target(
237 &crate::targets::ProcessExecutor,
238 &checkout.target,
239 &checkout.worktree_root,
240 )? {
241 Ok(Some(checkout.worktree_root.clone()))
242 } else {
243 Ok(None)
244 }
245 }
246 })
247 .await;
248 match owned_clone {
249 Ok(Some(path)) => {
250 self.push_notice(&session_id, format!(
251 "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
252 path.display(),
253 ));
254 return;
255 }
256 Err(error) => {
257 tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
258 return;
259 }
260 Ok(None) => {}
261 }
262 match self
263 .tear_down_stopped_session(
264 session_id.clone(),
265 LifecycleKind::DestroyStopped,
266 BranchDisposition::Keep,
267 CheckoutDisposition::KeepWhenDirty,
268 )
269 .await
270 {
271 Ok(retained) => {
272 let mut text = format!(
273 "Session {short} was lost because its managed target no longer exists; its record was removed."
274 );
275 if let Some(path) = retained {
276 text.push_str(&format!(
277 " Its checkout has uncommitted changes, so it was kept at {}.",
278 path.display()
279 ));
280 }
281 self.push_notice(&session_id, text);
282 }
283 Err(error) => {
284 tracing::warn!(
285 %session_id,
286 error = format!("{error:#}"),
287 "could not discard the record of a lost session"
288 );
289 self.push_notice(
290 &session_id,
291 format!(
292 "Session {short} was lost, but its record could not be removed: {error:#}"
293 ),
294 );
295 }
296 }
297 }
298
299 pub(crate) async fn archive_aged_sessions(
305 self: &Arc<Self>,
306 older_than_days: u32,
307 ) -> Result<usize> {
308 self.wiki()
309 .sync_now(true)
310 .await
311 .context("sync SessionWiki before archiving stopped sessions")?;
312 let Ok(_work) = crate::upgrade::activity_unless_draining("SessionWiki archive") else {
316 tracing::debug!(
317 "a daemon upgrade is waiting; the SessionWiki archive pass waits for the next daemon"
318 );
319 return Ok(0);
320 };
321 let refreshable = blocking(move || {
324 let controller = Controller::load()?;
325 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
326 let mut sessions = controller
327 .state
328 .sessions
329 .values()
330 .filter(|session| session.state == SessionState::Stopped)
331 .filter(|session| {
332 chrono::DateTime::parse_from_rfc3339(&session.updated_at)
333 .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
334 })
335 .filter(|session| {
336 session.managed_worktree.as_ref().is_some_and(|owned| {
337 owned.kind == mj_core::state::ManagedCheckoutKind::Clone
338 })
339 })
340 .filter(|session| {
341 session.publication.as_ref().is_none_or(|evidence| {
342 evidence.state != mj_core::state::PublicationState::Published
343 && !evidence.dirty
344 && !evidence.stashed
345 })
346 })
347 .cloned()
348 .collect::<Vec<_>>();
349 sessions.sort_by(|a, b| {
350 a.publication
351 .as_ref()
352 .map(|e| &e.checked_at)
353 .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
354 });
355 sessions.truncate(4);
356 Ok(sessions)
357 })
358 .await?;
359 let mut checks = tokio::task::JoinSet::new();
360 for session in refreshable {
361 checks.spawn_blocking(move || {
362 let assessment =
363 crate::controller::publication::refresh_stopped_clone_publication(&session);
364 (session.id, assessment)
365 });
366 }
367 let mut refreshed = false;
368 while let Some(done) = checks.join_next().await {
369 let (id, assessment) = done.context("publication refresh task failed")?;
370 if let Some(assessment) = assessment {
371 refreshed |= blocking(move || {
372 crate::database::set_publication_assessment_if_current(&id, &assessment)
373 })
374 .await?;
375 }
376 }
377 if refreshed {
378 self.reload_controller().await?;
379 self.publish_revision();
380 }
381 let (candidates, children) = blocking(move || {
382 let controller = Controller::load()?;
383 let candidates = crate::sessionwiki::sessions_ready_to_archive(
384 &controller.state.sessions,
385 &controller.state.subagents,
386 chrono::Utc::now(),
387 older_than_days,
388 );
389 let children: std::collections::BTreeSet<String> = candidates
390 .iter()
391 .filter(|id| controller.state.is_subagent_session(id))
392 .cloned()
393 .collect();
394 Ok((candidates, children))
395 })
396 .await
397 .context("select the stopped sessions old enough to archive")?;
398 if candidates.is_empty() {
399 return Ok(0);
400 }
401 let indexed = blocking({
402 let candidates = candidates
403 .iter()
404 .filter(|id| !children.contains(*id))
405 .cloned()
406 .collect::<Vec<_>>();
407 move || crate::sessionwiki::indexed_with_messages(&candidates)
408 })
409 .await
410 .context("check the SessionWiki index before archiving")?;
411 let mut archived = 0;
412 for session_id in candidates {
413 if !children.contains(&session_id) && !indexed.contains(&session_id) {
414 tracing::warn!(
415 %session_id,
416 "SessionWiki holds no conversation for this stopped session; keeping it"
417 );
418 continue;
419 }
420 match self.archive_stopped_session(session_id.clone()).await {
421 Ok(()) => {
422 archived += 1;
423 tracing::info!(
424 %session_id,
425 older_than_days,
426 child = children.contains(&session_id),
427 "archived a stopped session; top-level conversations are retained in SessionWiki"
428 );
429 }
430 Err(error) => tracing::warn!(
431 %session_id,
432 error = %format!("{error:#}"),
433 "could not archive a stopped session"
434 ),
435 }
436 }
437 if archived > 0 {
438 self.wiki().request_sync(false);
440 }
441 Ok(archived)
442 }
443
444 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
450 self.tear_down_stopped_session(
451 session_id,
452 LifecycleKind::ArchiveStopped,
453 BranchDisposition::DeleteIfMerged,
454 CheckoutDisposition::Remove,
455 )
456 .await
457 .map(|_| ())
458 }
459
460 async fn tear_down_stopped_session(
462 self: &Arc<Self>,
463 session_id: String,
464 kind: LifecycleKind,
465 branch: BranchDisposition,
466 checkout: CheckoutDisposition,
467 ) -> Result<Option<PathBuf>> {
468 self.index_before_destroy(&session_id).await;
471 let children = blocking({
472 let session_id = session_id.clone();
473 move || {
474 Ok(crate::database::list_subagents(&session_id)?
475 .into_iter()
476 .map(|child| child.child_session_id)
477 .collect::<Vec<_>>())
478 }
479 })
480 .await?;
481 for child_id in children {
482 Box::pin(self.force_destroy_indexed_session(
485 child_id.clone(),
486 BranchDisposition::Keep,
487 LifecycleKind::ForceDestroy,
488 ))
489 .await
490 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
491 }
492 self.wait_for_deferred_cleanup(&session_id).await?;
493 let exists = blocking({
494 let session_id = session_id.clone();
495 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
496 })
497 .await?;
498 if !exists {
499 return Ok(None);
500 }
501 let retained = Arc::new(Mutex::new(None));
504 self.run_lifecycle(session_id, kind, {
505 let retained = retained.clone();
506 move |state, session_id, cancelled| async move {
507 blocking(move || {
508 let mut controller = Controller::load()?;
509 let executor = DaemonStageReportingExecutor::new(
510 CancellableProcessExecutor::new(cancelled),
511 state,
512 session_id.clone(),
513 );
514 let kept = controller.destroy_session_controlled_with_checkout(
515 &session_id,
516 &executor,
517 branch,
518 checkout,
519 )?;
520 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
521 Ok(DaemonLifecycleResult::Done)
522 })
523 .await
524 }
525 })
526 .await?;
527 let kept = retained
528 .lock()
529 .unwrap_or_else(PoisonError::into_inner)
530 .clone();
531 Ok(kept)
532 }
533
534 pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
542 use crate::sessionwiki::IndexedBeforeDestroy;
543 let outcome = self
544 .wiki()
545 .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
546 .await;
547 match outcome {
548 IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
549 %session_id,
550 reason,
551 "destroying a session without indexing it in SessionWiki"
552 ),
553 IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
554 %session_id,
555 %reason,
556 "could not index a session in SessionWiki before destroying it"
557 ),
558 outcome => tracing::debug!(
559 %session_id,
560 ?outcome,
561 "indexed a session in SessionWiki before destroying it"
562 ),
563 }
564 }
565
566 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
577 let (mut result, tearing_down) = {
578 let mut lifecycle_owner = self.owner();
579 let lifecycle = &mut lifecycle_owner.lifecycle;
580 let Some(active) = lifecycle.get_mut(session_id) else {
581 return Ok(());
582 };
583 if !active.is_running() {
584 return Ok(());
585 }
586 let tearing_down = active.kind.is_teardown();
587 if !tearing_down {
593 active.request_cancel();
594 }
595 (active.result.clone(), tearing_down)
596 };
597 let mut finished = wait_for_lifecycle_result(&mut result).await;
598 if finished.is_err() && tearing_down {
599 {
600 let mut lifecycle_owner = self.owner();
601 if let Some(active) = lifecycle_owner.lifecycle.get_mut(session_id) {
602 active.request_cancel();
603 }
604 }
605 finished = wait_for_lifecycle_result(&mut result).await;
606 }
607 match finished {
608 Ok(Ok(())) => Ok(()),
611 Ok(Err(())) => bail!(
612 "daemon lifecycle operation stopped without a result for session {session_id}"
613 ),
614 Err(_) => bail!(
615 "session {session_id} still has an operation that did not stop after cancellation; try again"
616 ),
617 }
618 }
619
620 pub async fn force_destroy_session(
625 self: &Arc<Self>,
626 session_id: String,
627 branch: BranchDisposition,
628 ) -> Result<()> {
629 let _upgrade_work = crate::upgrade::destroy_activity(&session_id)?;
632 self.index_before_destroy(&session_id).await;
633 self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
634 .await
635 }
636
637 pub(super) async fn force_destroy_indexed_session(
642 self: &Arc<Self>,
643 session_id: String,
644 branch: BranchDisposition,
645 kind: LifecycleKind,
646 ) -> Result<()> {
647 let children = blocking({
648 let session_id = session_id.clone();
649 move || {
650 Ok(crate::database::list_subagents(&session_id)?
651 .into_iter()
652 .map(|child| child.child_session_id)
653 .collect::<Vec<_>>())
654 }
655 })
656 .await?;
657 for child_id in children {
658 Box::pin(self.force_destroy_indexed_session(
660 child_id.clone(),
661 BranchDisposition::Keep,
662 kind,
663 ))
664 .await
665 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
666 }
667 self.preempt_active_lifecycle(&session_id).await?;
668 let exists = blocking({
669 let session_id = session_id.clone();
670 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
671 })
672 .await?;
673 if !exists {
674 return Ok(());
675 }
676 self.run_lifecycle(
677 session_id,
678 kind,
679 move |state, session_id, cancelled| async move {
680 let _recovery_reservation = tokio::task::spawn_blocking({
681 let observer = state.recovery_observer.clone();
682 let session_id = session_id.clone();
683 let cancelled = cancelled.clone();
684 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
685 })
686 .await
687 .context("reserve recovery for daemon force-destroy task")??;
688 blocking({
689 let session_id = session_id.clone();
690 move || {
691 let mut controller = Controller::load()?;
692 let executor = DaemonStageReportingExecutor::new(
693 CancellableProcessExecutor::new(cancelled),
694 state,
695 session_id.clone(),
696 );
697 controller.force_destroy_session(&session_id, &executor, branch)?;
698 crate::controller::move_session::release_move_queue_hold(&session_id);
699 Ok(DaemonLifecycleResult::Done)
700 }
701 })
702 .await
703 },
704 )
705 .await?;
706 Ok(())
707 }
708}
709
710async fn wait_for_lifecycle_result(
713 result: &mut LifecycleWatch,
714) -> std::result::Result<std::result::Result<(), ()>, tokio::time::error::Elapsed> {
715 tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
716 loop {
717 if result.borrow().is_some() {
718 return Ok(());
719 }
720 if result.changed().await.is_err() {
721 return Err(());
722 }
723 }
724 })
725 .await
726}