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 refreshable = blocking(move || {
315 let controller = Controller::load()?;
316 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
317 let mut sessions = controller
318 .state
319 .sessions
320 .values()
321 .filter(|session| session.state == SessionState::Stopped)
322 .filter(|session| {
323 chrono::DateTime::parse_from_rfc3339(&session.updated_at)
324 .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
325 })
326 .filter(|session| {
327 session.managed_worktree.as_ref().is_some_and(|owned| {
328 owned.kind == mj_core::state::ManagedCheckoutKind::Clone
329 })
330 })
331 .filter(|session| {
332 session.publication.as_ref().is_none_or(|evidence| {
333 evidence.state != mj_core::state::PublicationState::Published
334 && !evidence.dirty
335 && !evidence.stashed
336 })
337 })
338 .cloned()
339 .collect::<Vec<_>>();
340 sessions.sort_by(|a, b| {
341 a.publication
342 .as_ref()
343 .map(|e| &e.checked_at)
344 .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
345 });
346 sessions.truncate(4);
347 Ok(sessions)
348 })
349 .await?;
350 let mut checks = tokio::task::JoinSet::new();
351 for session in refreshable {
352 checks.spawn_blocking(move || {
353 let assessment =
354 crate::controller::publication::refresh_stopped_clone_publication(&session);
355 (session.id, assessment)
356 });
357 }
358 let mut refreshed = false;
359 while let Some(done) = checks.join_next().await {
360 let (id, assessment) = done.context("publication refresh task failed")?;
361 if let Some(assessment) = assessment {
362 refreshed |= blocking(move || {
363 crate::database::set_publication_assessment_if_current(&id, &assessment)
364 })
365 .await?;
366 }
367 }
368 if refreshed {
369 self.reload_controller().await?;
370 self.publish_revision();
371 }
372 let candidates = blocking(move || {
373 let controller = Controller::load()?;
374 Ok(crate::sessionwiki::sessions_ready_to_archive(
375 &controller.state.sessions,
376 &controller.state.subagents,
377 chrono::Utc::now(),
378 older_than_days,
379 ))
380 })
381 .await
382 .context("select the stopped sessions old enough to archive")?;
383 if candidates.is_empty() {
384 return Ok(0);
385 }
386 let indexed = blocking({
387 let candidates = candidates.clone();
388 move || crate::sessionwiki::indexed_with_messages(&candidates)
389 })
390 .await
391 .context("check the SessionWiki index before archiving")?;
392 let mut archived = 0;
393 for session_id in candidates {
394 if !indexed.contains(&session_id) {
395 tracing::warn!(
396 %session_id,
397 "SessionWiki holds no conversation for this stopped session; keeping it"
398 );
399 continue;
400 }
401 match self.archive_stopped_session(session_id.clone()).await {
402 Ok(()) => {
403 archived += 1;
404 tracing::info!(
405 %session_id,
406 older_than_days,
407 "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
408 );
409 }
410 Err(error) => tracing::warn!(
411 %session_id,
412 error = %format!("{error:#}"),
413 "could not archive a stopped session"
414 ),
415 }
416 }
417 if archived > 0 {
418 self.wiki().request_sync(false);
420 }
421 Ok(archived)
422 }
423
424 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
430 self.tear_down_stopped_session(
431 session_id,
432 LifecycleKind::ArchiveStopped,
433 BranchDisposition::DeleteIfMerged,
434 CheckoutDisposition::Remove,
435 )
436 .await
437 .map(|_| ())
438 }
439
440 async fn tear_down_stopped_session(
442 self: &Arc<Self>,
443 session_id: String,
444 kind: LifecycleKind,
445 branch: BranchDisposition,
446 checkout: CheckoutDisposition,
447 ) -> Result<Option<PathBuf>> {
448 self.index_before_destroy(&session_id).await;
451 let children = blocking({
452 let session_id = session_id.clone();
453 move || {
454 Ok(crate::database::list_subagents(&session_id)?
455 .into_iter()
456 .map(|child| child.child_session_id)
457 .collect::<Vec<_>>())
458 }
459 })
460 .await?;
461 for child_id in children {
462 Box::pin(self.force_destroy_indexed_session(
465 child_id.clone(),
466 BranchDisposition::Keep,
467 LifecycleKind::ForceDestroy,
468 ))
469 .await
470 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
471 }
472 self.wait_for_deferred_cleanup(&session_id).await?;
473 let exists = blocking({
474 let session_id = session_id.clone();
475 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
476 })
477 .await?;
478 if !exists {
479 return Ok(None);
480 }
481 let retained = Arc::new(Mutex::new(None));
484 self.run_lifecycle(session_id, kind, {
485 let retained = retained.clone();
486 move |state, session_id, cancelled| async move {
487 blocking(move || {
488 let mut controller = Controller::load()?;
489 let executor = DaemonStageReportingExecutor::new(
490 CancellableProcessExecutor::new(cancelled),
491 state,
492 session_id.clone(),
493 );
494 let kept = controller.destroy_session_controlled_with_checkout(
495 &session_id,
496 &executor,
497 branch,
498 checkout,
499 )?;
500 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
501 Ok(DaemonLifecycleResult::Done)
502 })
503 .await
504 }
505 })
506 .await?;
507 let kept = retained
508 .lock()
509 .unwrap_or_else(PoisonError::into_inner)
510 .clone();
511 Ok(kept)
512 }
513
514 pub(super) async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
523 use crate::sessionwiki::IndexedBeforeDestroy;
524 let outcome = self
525 .wiki()
526 .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
527 .await;
528 match outcome {
529 IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
530 %session_id,
531 reason,
532 "destroying a session without indexing it in SessionWiki"
533 ),
534 IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
535 %session_id,
536 %reason,
537 "could not index a session in SessionWiki before destroying it"
538 ),
539 outcome => tracing::debug!(
540 %session_id,
541 ?outcome,
542 "indexed a session in SessionWiki before destroying it"
543 ),
544 }
545 }
546
547 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
558 let (mut result, tearing_down) = {
559 let mut lifecycle_owner = self.owner();
560 let lifecycle = &mut lifecycle_owner.lifecycle;
561 let Some(active) = lifecycle.get_mut(session_id) else {
562 return Ok(());
563 };
564 if !active.is_running() {
565 return Ok(());
566 }
567 let tearing_down = active.kind.is_teardown();
568 if !tearing_down {
574 active.request_cancel();
575 }
576 (active.result.clone(), tearing_down)
577 };
578 let mut finished = wait_for_lifecycle_result(&mut result).await;
579 if finished.is_err() && tearing_down {
580 {
581 let mut lifecycle_owner = self.owner();
582 if let Some(active) = lifecycle_owner.lifecycle.get_mut(session_id) {
583 active.request_cancel();
584 }
585 }
586 finished = wait_for_lifecycle_result(&mut result).await;
587 }
588 match finished {
589 Ok(Ok(())) => Ok(()),
592 Ok(Err(())) => bail!(
593 "daemon lifecycle operation stopped without a result for session {session_id}"
594 ),
595 Err(_) => bail!(
596 "session {session_id} still has an operation that did not stop after cancellation; try again"
597 ),
598 }
599 }
600
601 pub async fn force_destroy_session(
606 self: &Arc<Self>,
607 session_id: String,
608 branch: BranchDisposition,
609 ) -> Result<()> {
610 let _upgrade_work = crate::upgrade::destroy_activity(&session_id)?;
613 self.index_before_destroy(&session_id).await;
614 self.force_destroy_indexed_session(session_id, branch, LifecycleKind::ForceDestroy)
615 .await
616 }
617
618 pub(super) async fn force_destroy_indexed_session(
623 self: &Arc<Self>,
624 session_id: String,
625 branch: BranchDisposition,
626 kind: LifecycleKind,
627 ) -> Result<()> {
628 let children = blocking({
629 let session_id = session_id.clone();
630 move || {
631 Ok(crate::database::list_subagents(&session_id)?
632 .into_iter()
633 .map(|child| child.child_session_id)
634 .collect::<Vec<_>>())
635 }
636 })
637 .await?;
638 for child_id in children {
639 Box::pin(self.force_destroy_indexed_session(
641 child_id.clone(),
642 BranchDisposition::Keep,
643 kind,
644 ))
645 .await
646 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
647 }
648 self.preempt_active_lifecycle(&session_id).await?;
649 let exists = blocking({
650 let session_id = session_id.clone();
651 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
652 })
653 .await?;
654 if !exists {
655 return Ok(());
656 }
657 self.run_lifecycle(
658 session_id,
659 kind,
660 move |state, session_id, cancelled| async move {
661 let _recovery_reservation = tokio::task::spawn_blocking({
662 let observer = state.recovery_observer.clone();
663 let session_id = session_id.clone();
664 let cancelled = cancelled.clone();
665 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
666 })
667 .await
668 .context("reserve recovery for daemon force-destroy task")??;
669 blocking({
670 let session_id = session_id.clone();
671 move || {
672 let mut controller = Controller::load()?;
673 let executor = DaemonStageReportingExecutor::new(
674 CancellableProcessExecutor::new(cancelled),
675 state,
676 session_id.clone(),
677 );
678 controller.force_destroy_session(&session_id, &executor, branch)?;
679 crate::controller::move_session::release_move_queue_hold(&session_id);
680 Ok(DaemonLifecycleResult::Done)
681 }
682 })
683 .await
684 },
685 )
686 .await?;
687 Ok(())
688 }
689}
690
691async fn wait_for_lifecycle_result(
694 result: &mut LifecycleWatch,
695) -> std::result::Result<std::result::Result<(), ()>, tokio::time::error::Elapsed> {
696 tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
697 loop {
698 if result.borrow().is_some() {
699 return Ok(());
700 }
701 if result.changed().await.is_err() {
702 return Err(());
703 }
704 }
705 })
706 .await
707}