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 let owned_clone = blocking({
211 let session_id = session_id.clone();
212 move || {
213 let controller = Controller::load()?;
214 let Some(checkout) = controller
215 .state
216 .sessions
217 .get(&session_id)
218 .and_then(|session| session.managed_worktree.as_ref())
219 .filter(|checkout| checkout.kind == mj_core::state::ManagedCheckoutKind::Clone)
220 else {
221 return Ok(None);
222 };
223 if crate::controller::path_exists_on_managed_target(
224 &crate::targets::ProcessExecutor,
225 &checkout.target,
226 &checkout.worktree_root,
227 )? {
228 Ok(Some(checkout.worktree_root.clone()))
229 } else {
230 Ok(None)
231 }
232 }
233 })
234 .await;
235 match owned_clone {
236 Ok(Some(path)) => {
237 self.push_notice(&session_id, format!(
238 "Session {short} lost its target, but its owned clone remains at {}; keeping the record for inspection",
239 path.display(),
240 ));
241 return;
242 }
243 Err(error) => {
244 tracing::warn!(%session_id, error = %format!("{error:#}"), "could not inspect owned clone before lost-session cleanup");
245 return;
246 }
247 Ok(None) => {}
248 }
249 match self
250 .tear_down_stopped_session(
251 session_id.clone(),
252 LifecycleKind::DestroyStopped,
253 BranchDisposition::Keep,
254 CheckoutDisposition::KeepWhenDirty,
255 )
256 .await
257 {
258 Ok(retained) => {
259 let mut text = format!(
260 "Session {short} was lost because its managed target no longer exists; its record was removed."
261 );
262 if let Some(path) = retained {
263 text.push_str(&format!(
264 " Its checkout has uncommitted changes, so it was kept at {}.",
265 path.display()
266 ));
267 }
268 self.push_notice(&session_id, text);
269 }
270 Err(error) => {
271 tracing::warn!(
272 %session_id,
273 error = format!("{error:#}"),
274 "could not discard the record of a lost session"
275 );
276 self.push_notice(
277 &session_id,
278 format!(
279 "Session {short} was lost, but its record could not be removed: {error:#}"
280 ),
281 );
282 }
283 }
284 }
285
286 pub(crate) async fn archive_aged_sessions(
292 self: &Arc<Self>,
293 older_than_days: u32,
294 ) -> Result<usize> {
295 self.wiki()
296 .sync_now(true)
297 .await
298 .context("sync SessionWiki before archiving stopped sessions")?;
299 let refreshable = blocking(move || {
302 let controller = Controller::load()?;
303 let cutoff = chrono::Utc::now() - chrono::Duration::days(i64::from(older_than_days));
304 let mut sessions = controller
305 .state
306 .sessions
307 .values()
308 .filter(|session| session.state == SessionState::Stopped)
309 .filter(|session| {
310 chrono::DateTime::parse_from_rfc3339(&session.updated_at)
311 .is_ok_and(|time| time.with_timezone(&chrono::Utc) <= cutoff)
312 })
313 .filter(|session| {
314 session.managed_worktree.as_ref().is_some_and(|owned| {
315 owned.kind == mj_core::state::ManagedCheckoutKind::Clone
316 })
317 })
318 .filter(|session| {
319 session.publication.as_ref().is_none_or(|evidence| {
320 evidence.state != mj_core::state::PublicationState::Published
321 && !evidence.dirty
322 && !evidence.stashed
323 })
324 })
325 .cloned()
326 .collect::<Vec<_>>();
327 sessions.sort_by(|a, b| {
328 a.publication
329 .as_ref()
330 .map(|e| &e.checked_at)
331 .cmp(&b.publication.as_ref().map(|e| &e.checked_at))
332 });
333 sessions.truncate(4);
334 Ok(sessions)
335 })
336 .await?;
337 let mut checks = tokio::task::JoinSet::new();
338 for session in refreshable {
339 checks.spawn_blocking(move || {
340 let assessment =
341 crate::controller::publication::refresh_stopped_clone_publication(&session);
342 (session.id, assessment)
343 });
344 }
345 let mut refreshed = false;
346 while let Some(done) = checks.join_next().await {
347 let (id, assessment) = done.context("publication refresh task failed")?;
348 if let Some(assessment) = assessment {
349 refreshed |= blocking(move || {
350 crate::database::set_publication_assessment_if_current(&id, &assessment)
351 })
352 .await?;
353 }
354 }
355 if refreshed {
356 self.reload_controller().await?;
357 self.publish_revision();
358 }
359 let candidates = blocking(move || {
360 let controller = Controller::load()?;
361 Ok(crate::sessionwiki::sessions_ready_to_archive(
362 &controller.state.sessions,
363 &controller.state.subagents,
364 chrono::Utc::now(),
365 older_than_days,
366 ))
367 })
368 .await
369 .context("select the stopped sessions old enough to archive")?;
370 if candidates.is_empty() {
371 return Ok(0);
372 }
373 let indexed = blocking({
374 let candidates = candidates.clone();
375 move || crate::sessionwiki::indexed_with_messages(&candidates)
376 })
377 .await
378 .context("check the SessionWiki index before archiving")?;
379 let mut archived = 0;
380 for session_id in candidates {
381 if !indexed.contains(&session_id) {
382 tracing::warn!(
383 %session_id,
384 "SessionWiki holds no conversation for this stopped session; keeping it"
385 );
386 continue;
387 }
388 match self.archive_stopped_session(session_id.clone()).await {
389 Ok(()) => {
390 archived += 1;
391 tracing::info!(
392 %session_id,
393 older_than_days,
394 "archived a stopped session: SessionWiki keeps the conversation, and the repository keeps the branch unless another branch already contains it"
395 );
396 }
397 Err(error) => tracing::warn!(
398 %session_id,
399 error = %format!("{error:#}"),
400 "could not archive a stopped session"
401 ),
402 }
403 }
404 if archived > 0 {
405 self.wiki().request_sync(false);
407 }
408 Ok(archived)
409 }
410
411 async fn archive_stopped_session(self: &Arc<Self>, session_id: String) -> Result<()> {
417 self.tear_down_stopped_session(
418 session_id,
419 LifecycleKind::ArchiveStopped,
420 BranchDisposition::DeleteIfMerged,
421 CheckoutDisposition::Remove,
422 )
423 .await
424 .map(|_| ())
425 }
426
427 async fn tear_down_stopped_session(
429 self: &Arc<Self>,
430 session_id: String,
431 kind: LifecycleKind,
432 branch: BranchDisposition,
433 checkout: CheckoutDisposition,
434 ) -> Result<Option<PathBuf>> {
435 self.index_before_destroy(&session_id).await;
438 let children = blocking({
439 let session_id = session_id.clone();
440 move || {
441 Ok(crate::database::list_subagents(&session_id)?
442 .into_iter()
443 .map(|child| child.child_session_id)
444 .collect::<Vec<_>>())
445 }
446 })
447 .await?;
448 for child_id in children {
449 Box::pin(self.force_destroy_indexed_session(child_id.clone(), BranchDisposition::Keep))
452 .await
453 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
454 }
455 self.wait_for_deferred_cleanup(&session_id).await?;
456 let exists = blocking({
457 let session_id = session_id.clone();
458 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
459 })
460 .await?;
461 if !exists {
462 return Ok(None);
463 }
464 let retained = Arc::new(Mutex::new(None));
467 self.run_lifecycle(session_id, kind, {
468 let retained = retained.clone();
469 move |state, session_id, cancelled| async move {
470 blocking(move || {
471 let mut controller = Controller::load()?;
472 let executor = DaemonStageReportingExecutor::new(
473 CancellableProcessExecutor::new(cancelled),
474 state,
475 session_id.clone(),
476 );
477 let kept = controller.destroy_session_controlled_with_checkout(
478 &session_id,
479 &executor,
480 branch,
481 checkout,
482 )?;
483 *retained.lock().unwrap_or_else(PoisonError::into_inner) = kept;
484 Ok(DaemonLifecycleResult::Done)
485 })
486 .await
487 }
488 })
489 .await?;
490 let kept = retained
491 .lock()
492 .unwrap_or_else(PoisonError::into_inner)
493 .clone();
494 Ok(kept)
495 }
496
497 async fn index_before_destroy(self: &Arc<Self>, session_id: &str) {
506 use crate::sessionwiki::IndexedBeforeDestroy;
507 let outcome = self
508 .wiki()
509 .index_before_destroy(session_id, crate::sessionwiki::DESTROY_SYNC_WAIT)
510 .await;
511 match outcome {
512 IndexedBeforeDestroy::Unavailable(reason) => tracing::info!(
513 %session_id,
514 reason,
515 "destroying a session without indexing it in SessionWiki"
516 ),
517 IndexedBeforeDestroy::Failed(reason) => tracing::warn!(
518 %session_id,
519 %reason,
520 "could not index a session in SessionWiki before destroying it"
521 ),
522 outcome => tracing::debug!(
523 %session_id,
524 ?outcome,
525 "indexed a session in SessionWiki before destroying it"
526 ),
527 }
528 }
529
530 pub(super) async fn preempt_active_lifecycle(self: &Arc<Self>, session_id: &str) -> Result<()> {
541 let mut result = {
542 let lifecycle = self
543 .lifecycle
544 .lock()
545 .unwrap_or_else(PoisonError::into_inner);
546 let Some(active) = lifecycle.get(session_id) else {
547 return Ok(());
548 };
549 if !active.result.borrow().is_none() {
550 return Ok(());
551 }
552 active.request_cancel();
553 active.result.clone()
554 };
555 let finished = tokio::time::timeout(FORCE_DESTROY_PREEMPT_TIMEOUT, async {
556 loop {
557 if result.borrow().is_some() {
558 return Ok(());
559 }
560 if result.changed().await.is_err() {
561 return Err(());
562 }
563 }
564 })
565 .await;
566 match finished {
567 Ok(Ok(())) => Ok(()),
570 Ok(Err(())) => bail!(
571 "daemon lifecycle operation stopped without a result for session {session_id}"
572 ),
573 Err(_) => bail!(
574 "session {session_id} still has an operation that did not stop after cancellation; try again"
575 ),
576 }
577 }
578
579 pub async fn force_destroy_session(
584 self: &Arc<Self>,
585 session_id: String,
586 branch: BranchDisposition,
587 ) -> Result<()> {
588 self.index_before_destroy(&session_id).await;
589 self.force_destroy_indexed_session(session_id, branch).await
590 }
591
592 async fn force_destroy_indexed_session(
595 self: &Arc<Self>,
596 session_id: String,
597 branch: BranchDisposition,
598 ) -> Result<()> {
599 let children = blocking({
600 let session_id = session_id.clone();
601 move || {
602 Ok(crate::database::list_subagents(&session_id)?
603 .into_iter()
604 .map(|child| child.child_session_id)
605 .collect::<Vec<_>>())
606 }
607 })
608 .await?;
609 for child_id in children {
610 Box::pin(self.force_destroy_indexed_session(child_id.clone(), BranchDisposition::Keep))
612 .await
613 .with_context(|| format!("destroy sub-agent {child_id} before its parent"))?;
614 }
615 self.preempt_active_lifecycle(&session_id).await?;
616 let exists = blocking({
617 let session_id = session_id.clone();
618 move || Ok(Controller::load()?.state.sessions.contains_key(&session_id))
619 })
620 .await?;
621 if !exists {
622 return Ok(());
623 }
624 self.run_lifecycle(
625 session_id,
626 LifecycleKind::ForceDestroy,
627 move |state, session_id, cancelled| async move {
628 let _recovery_reservation = tokio::task::spawn_blocking({
629 let observer = state.recovery_observer.clone();
630 let session_id = session_id.clone();
631 let cancelled = cancelled.clone();
632 move || reserve_recovery_or_cancel(&observer, &session_id, &cancelled)
633 })
634 .await
635 .context("reserve recovery for daemon force-destroy task")??;
636 blocking({
637 let session_id = session_id.clone();
638 move || {
639 let mut controller = Controller::load()?;
640 let executor = DaemonStageReportingExecutor::new(
641 CancellableProcessExecutor::new(cancelled),
642 state,
643 session_id.clone(),
644 );
645 controller.force_destroy_session(&session_id, &executor, branch)?;
646 crate::controller::move_session::release_move_queue_hold(&session_id);
647 Ok(DaemonLifecycleResult::Done)
648 }
649 })
650 .await
651 },
652 )
653 .await?;
654 Ok(())
655 }
656
657 pub async fn force_delete_workspace(self: &Arc<Self>, workspace_id: String) -> Result<()> {
666 ensure!(
667 !self.workspace_has_active_resume(&workspace_id),
668 "workspace has a session resume in progress"
669 );
670 let sessions = blocking({
671 let workspace_id = workspace_id.clone();
672 move || {
673 let controller = Controller::load()?;
674 Ok(active_sessions_for_force_destruction(
675 &controller,
676 &workspace_id,
677 ))
678 }
679 })
680 .await?;
681 for (index, session_id) in sessions.iter().enumerate() {
682 if let Err(error) = self
685 .force_destroy_session(session_id.clone(), BranchDisposition::Keep)
686 .await
687 {
688 let remaining = sessions.len() - index - 1;
689 let remaining = if remaining == 1 {
691 "1 session in the workspace remains".to_owned()
692 } else {
693 format!("{remaining} sessions in the workspace remain")
694 };
695 bail!("force-destroying session {session_id} failed: {error:#}; {remaining}");
696 }
697 }
698 blocking({
699 let workspace_id = workspace_id.clone();
700 move || crate::database::force_delete_workspace(&workspace_id)
701 })
702 .await?;
703 refresh_runtime_workspaces(self).await?;
704 Ok(())
705 }
706}