1use std::path::PathBuf;
4use std::time::Duration;
5
6use anyhow::{Context, Result, bail, ensure};
7
8use crate::session_manager::{SessionManagerControl, new_command_id};
9use mj_core::state::{CheckpointMetadata, SessionRecord, SessionState};
10
11use crate::targets::{self, CommandExecutor, ProcessExecutor, ProvisionStage, ProvisionStageGuard};
12use mj_core::relay::{RelayCommand, RelayExecutionState};
13
14use super::backend::backend_locator;
15use super::checkpoint::{
16 CheckpointExportPolicy, LatchExclusivity, prune_replaced_checkpoint,
17 release_projection_behind_checkpoint, verify_installed_checkpoint_gate, wait_for_relay_closed,
18};
19use super::worker_restart::WorkerRestartLeftNoWorker;
20use super::worktree::{
21 cleanup_managed_worktree, managed_worktree_checkout_is_dirty, retire_managed_worktree,
22};
23use super::{Controller, now, persist_session_record_transition_or_restore};
24
25#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub enum BranchDisposition {
28 Delete,
31 Keep,
34 DeleteIfMerged,
40}
41
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum CheckoutDisposition {
45 Remove,
48 KeepWhenDirty,
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub(super) enum SourceTargetDisposition {
57 Destroy,
59 RetainForInPlaceSwap,
62}
63
64impl Controller {
65 pub async fn suspend_session(&mut self, session_id: &str) -> Result<()> {
70 self.suspend_session_controlled(session_id, &ProcessExecutor)
71 .await
72 }
73
74 pub async fn suspend_session_controlled(
75 &mut self,
76 session_id: &str,
77 executor: &(impl CommandExecutor + Sync),
78 ) -> Result<()> {
79 if self
80 .suspend_session_controlled_with_manager(
81 session_id,
82 executor,
83 None,
84 None,
85 SourceTargetDisposition::Destroy,
86 )
87 .await?
88 {
89 self.cleanup_stopped_target(session_id, executor)?;
90 }
91 Ok(())
92 }
93
94 pub async fn suspend_session_managed_controlled(
95 &mut self,
96 session_id: &str,
97 executor: &(impl CommandExecutor + Sync),
98 manager: &SessionManagerControl,
99 ) -> Result<bool> {
100 self.suspend_session_controlled_with_manager(
101 session_id,
102 executor,
103 Some(manager),
104 None,
105 SourceTargetDisposition::Destroy,
106 )
107 .await
108 }
109
110 pub(super) async fn suspend_session_for_move(
111 &mut self,
112 session_id: &str,
113 executor: &(impl CommandExecutor + Sync),
114 manager: &SessionManagerControl,
115 operation: &mut mj_core::state::MoveOperation,
116 preparation: Option<&mj_core::state::MovePreparation>,
117 disposition: SourceTargetDisposition,
118 ) -> Result<bool> {
119 self.prepare_move_source_checkpoint(session_id, executor, manager, operation)
120 .await?;
121 self.suspend_session_controlled_with_manager(
122 session_id,
123 executor,
124 Some(manager),
125 Some((operation, preparation)),
126 disposition,
127 )
128 .await
129 }
130
131 async fn suspend_session_controlled_with_manager(
132 &mut self,
133 session_id: &str,
134 executor: &(impl CommandExecutor + Sync),
135 manager: Option<&SessionManagerControl>,
136 move_intent: Option<(
137 &mut mj_core::state::MoveOperation,
138 Option<&mj_core::state::MovePreparation>,
139 )>,
140 disposition: SourceTargetDisposition,
141 ) -> Result<bool> {
142 let previous = self
143 .state
144 .sessions
145 .get(session_id)
146 .with_context(|| format!("unknown session {session_id}"))?
147 .clone();
148 let record = self.state.sessions.get_mut(session_id).unwrap();
149 apply_close_checkpoint_started(record, now());
153 self.persist_session_transition_or_restore(
154 session_id,
155 &previous,
156 "persist closing state before checkpointing the session",
157 )?;
158
159 let mut latched = match self
162 .checkpoint_session_latched(
163 session_id,
164 executor,
165 manager,
166 LatchExclusivity::HoldThroughClose,
167 CheckpointExportPolicy::ReuseUnchangedArchive,
168 )
169 .await
170 {
171 Ok(latched) => latched,
172 Err(error) => {
173 let record = self.state.sessions.get_mut(session_id).unwrap();
174 apply_close_checkpoint_failure(record, &previous, &error, now());
178 return Err(
179 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
180 );
181 }
182 };
183
184 let artifact = latched.artifact.clone();
185 let record = self.state.sessions.get_mut(session_id).unwrap();
186 record.state = SessionState::Closing;
187 record.native_session_id = Some(artifact.native_session_id.clone());
188 record.checkpoint = Some(artifact.metadata.clone());
189 record.updated_at = now();
190 record.last_error = None;
191 record.last_checkpoint_error = None;
192 self.persist_checkpoint_transition_or_restore(
193 session_id,
194 &previous,
195 "persist verified checkpoint and closing state before sealing the relay",
196 )?;
197 if let Some((operation, preparation)) = move_intent {
198 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
201 let record = self.state.sessions.get_mut(session_id).unwrap();
202 record.state = previous.state;
203 record.last_error = Some(format!("{error:#}"));
204 self.persist_session_transition_or_restore(
205 session_id,
206 &previous,
207 "restore source after move preflight failure",
208 )?;
209 return Err(error);
210 }
211 operation.checkpoint = Some(artifact.metadata.clone());
212 operation.updated_at = now();
213 crate::database::save_move_operation(operation)?;
214 }
215 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
216 release_projection_behind_checkpoint(session_id, &artifact.metadata);
219
220 let close_command_id = new_command_id("close")?;
221 let barrier_command_id = latched.barrier_command_id.clone();
222 let close_result = {
223 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
224 latched
225 .relay
226 .connection_mut()
227 .submit(
228 close_command_id,
229 RelayCommand::Close {
230 barrier_command_id: barrier_command_id.clone(),
231 expected: latched.cursor.clone(),
232 },
233 )
234 .await
235 };
236 if let Err(error) = close_result {
237 self.record_interrupted_close(session_id, &error)?;
238 return Err(error.context("seal verified checkpoint for close"));
239 }
240 let close_result = {
241 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
242 latched
243 .relay
244 .connection_mut()
245 .submit(
246 new_command_id("checkpoint-complete")?,
247 RelayCommand::CompleteCheckpoint { barrier_command_id },
248 )
249 .await
250 };
251 if let Err(error) = close_result {
252 self.record_interrupted_close(session_id, &error)?;
253 return Err(error.context("release verified close checkpoint"));
254 }
255 let close_result = {
256 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
257 wait_for_relay_closed(latched.relay.connection_mut()).await
258 };
259 if let Err(error) = close_result {
260 self.record_interrupted_close(session_id, &error)?;
261 return Err(error);
262 }
263 latched.relay.release();
264
265 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
266 return Ok(false);
272 }
273 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
274 Ok(deferred) => Ok(deferred),
275 Err(error) => {
276 self.record_interrupted_close(session_id, &error)?;
277 Err(error)
278 }
279 }
280 }
281
282 pub async fn recover_interrupted_close_managed(
288 &mut self,
289 session_id: &str,
290 executor: &(impl CommandExecutor + Sync),
291 manager: &SessionManagerControl,
292 ) -> Result<bool> {
293 let (state, verified) = {
294 let session = self
295 .state
296 .sessions
297 .get(session_id)
298 .with_context(|| format!("unknown session {session_id}"))?;
299 ensure!(
300 matches!(
301 session.state,
302 SessionState::Closing | SessionState::Destroying
303 ),
304 "session {session_id} has no interrupted close to recover"
305 );
306 (session.state, session.checkpoint.clone())
307 };
308 if state == SessionState::Destroying {
309 let verified = verified.context("destroying session has no verified checkpoint")?;
310 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
311 }
312 ensure!(
313 state == SessionState::Closing,
314 "session {session_id} has no relay close to recover"
315 );
316 let handle = manager
317 .wait_for_session(session_id, Duration::from_secs(5))
318 .await?;
319 let mut lease = handle.lease_connection().await?;
320 let execution = lease.connection_mut().sync().await?.operational.execution;
321 match execution {
322 RelayExecutionState::Closed => {}
323 RelayExecutionState::Closing => {
324 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
325 wait_for_relay_closed(lease.connection_mut()).await?;
326 }
327 RelayExecutionState::Idle | RelayExecutionState::Running => {
328 lease.release();
329 return self
330 .suspend_session_controlled_with_manager(
331 session_id,
332 executor,
333 Some(manager),
334 None,
335 SourceTargetDisposition::Destroy,
336 )
337 .await;
338 }
339 }
340 lease.release();
341 let verified = verified.context("closed relay has no verified checkpoint")?;
342 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
343 }
344
345 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
352 self.fail_interrupted_lifecycle_with(
353 session_id,
354 cause,
355 crate::database::save_lifecycle_session,
356 )
357 }
358
359 fn fail_interrupted_lifecycle_with(
360 &mut self,
361 session_id: &str,
362 cause: &str,
363 persist: impl Fn(&SessionRecord) -> Result<()>,
364 ) -> Result<bool> {
365 let Some(record) = self.state.sessions.get_mut(session_id) else {
366 return Ok(false);
367 };
368 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
369 return Ok(false);
370 }
371 let previous = record.clone();
372 record.state = SessionState::Error;
373 record.updated_at = now();
374 record.last_error = Some(cause.to_owned());
375 persist_session_record_transition_or_restore(
376 &mut self.state,
377 session_id,
378 &previous,
379 "persist the failure of an interrupted lifecycle state",
380 &persist,
381 )?;
382 Ok(true)
383 }
384
385 pub fn fail_unready_session(
394 &mut self,
395 session_id: &str,
396 cause: &str,
397 observed_updated_at: &str,
398 ) -> Result<bool> {
399 self.fail_unready_session_with(
400 session_id,
401 cause,
402 observed_updated_at,
403 crate::database::save_lifecycle_session,
404 )
405 }
406
407 fn fail_unready_session_with(
408 &mut self,
409 session_id: &str,
410 cause: &str,
411 observed_updated_at: &str,
412 persist: impl Fn(&SessionRecord) -> Result<()>,
413 ) -> Result<bool> {
414 let Some(record) = self.state.sessions.get_mut(session_id) else {
415 return Ok(false);
416 };
417 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
418 return Ok(false);
419 }
420 let previous = record.clone();
421 record.state = SessionState::Error;
422 record.updated_at = now();
423 record.last_error = Some(cause.to_owned());
424 persist_session_record_transition_or_restore(
425 &mut self.state,
426 session_id,
427 &previous,
428 "persist a session whose harness never became usable",
429 &persist,
430 )?;
431 Ok(true)
432 }
433
434 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
438 let Some(record) = self.state.sessions.get(session_id) else {
439 return Ok(false);
440 };
441 let previous = record.clone();
444 let record = self.state.sessions.get_mut(session_id).unwrap();
445 record.last_error = Some(cause.to_owned());
446 record.updated_at = now();
447 persist_session_record_transition_or_restore(
448 &mut self.state,
449 session_id,
450 &previous,
451 "persist the reason a close did not finish",
452 &crate::database::save_lifecycle_session,
453 )?;
454 Ok(true)
455 }
456
457 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
462 let Some(record) = self.state.sessions.get(session_id) else {
463 return Ok(false);
464 };
465 if record.public_error().is_none() {
466 return Ok(false);
467 }
468 let previous = record.clone();
469 let record = self.state.sessions.get_mut(session_id).unwrap();
470 record.last_error = None;
471 record.updated_at = now();
472 persist_session_record_transition_or_restore(
473 &mut self.state,
474 session_id,
475 &previous,
476 "clear the reason a close did not finish",
477 &crate::database::save_lifecycle_session,
478 )?;
479 Ok(true)
480 }
481
482 pub fn suspend_session_without_checkpoint(
495 &mut self,
496 session_id: &str,
497 executor: &impl CommandExecutor,
498 ) -> Result<bool> {
499 self.suspend_session_without_checkpoint_with(
500 session_id,
501 executor,
502 crate::database::save_lifecycle_session,
503 )
504 }
505
506 fn suspend_session_without_checkpoint_with(
507 &mut self,
508 session_id: &str,
509 executor: &impl CommandExecutor,
510 persist: impl Fn(&SessionRecord) -> Result<()>,
511 ) -> Result<bool> {
512 let session = self
513 .state
514 .sessions
515 .get(session_id)
516 .with_context(|| format!("unknown session {session_id}"))?
517 .clone();
518 ensure!(
519 has_nothing_to_checkpoint(&session),
520 "session {session_id} has a workspace to checkpoint; close it gracefully instead"
521 );
522 self.stop_target_and_settle(session_id, &session, executor, &persist)
523 }
524
525 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
526 let record = self.state.sessions.get_mut(session_id).unwrap();
527 apply_interrupted_close_error(record, error, &now());
528 self.persist_session_state(session_id)
529 }
530
531 fn destroy_after_verified_checkpoint(
534 &mut self,
535 session_id: &str,
536 verified: &CheckpointMetadata,
537 executor: &impl CommandExecutor,
538 ) -> Result<bool> {
539 self.destroy_after_verified_checkpoint_with(
540 session_id,
541 verified,
542 executor,
543 crate::database::save_lifecycle_session,
544 )
545 }
546
547 fn destroy_after_verified_checkpoint_with(
548 &mut self,
549 session_id: &str,
550 verified: &CheckpointMetadata,
551 executor: &impl CommandExecutor,
552 persist: impl Fn(&SessionRecord) -> Result<()>,
553 ) -> Result<bool> {
554 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
555 let _target_guard = target_mutex.lock().map_err(|_| {
556 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
557 })?;
558 let session = self
559 .state
560 .sessions
561 .get(session_id)
562 .with_context(|| format!("unknown session {session_id}"))?
563 .clone();
564 ensure!(
565 matches!(
566 session.state,
567 SessionState::Closing | SessionState::Destroying
568 ),
569 "refusing to destroy session {session_id}: it is not closing or destroying"
570 );
571 ensure!(
572 session.checkpoint.as_ref() == Some(verified),
573 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
574 );
575 if session.state == SessionState::Closing {
576 let record = self.state.sessions.get_mut(session_id).unwrap();
577 record.state = SessionState::Destroying;
578 record.updated_at = now();
579 record.last_error = None;
580 persist_session_record_transition_or_restore(
581 &mut self.state,
582 session_id,
583 &session,
584 "persist destroying state before target cleanup",
585 &persist,
586 )?;
587 }
588
589 let destroying = self
590 .state
591 .sessions
592 .get(session_id)
593 .expect("destroying session disappeared")
594 .clone();
595 {
596 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
597 verify_installed_checkpoint_gate(session_id, verified)?;
598 }
599 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
604 tracing::warn!(
605 session_id,
606 error = format!("{error:#}"),
607 "could not record that the second-opinion conversation ends with this target"
608 );
609 }
610 let locator = destroying
611 .target
612 .as_ref()
613 .context("session has no target")?;
614 let backend = backend_locator(locator, &destroying, &self.config)?;
615 let deferred = if self.state.subagents.contains_key(session_id) {
616 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
617 false
618 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
619 plan.execute(executor)?;
620 true
621 } else {
622 execute_target_cleanup(&backend, session_id, executor)?;
623 false
624 };
625 if let Some(worktree) = &destroying.managed_worktree {
626 retire_managed_worktree(executor, worktree)
627 .context("retire managed raw-session worktree after verified close")?;
628 }
629 let record = self.state.sessions.get_mut(session_id).unwrap();
630 record.state = SessionState::Stopped;
631 if !deferred {
632 record.target = None;
633 }
634 record.updated_at = now();
635 record.last_error = None;
636 persist_session_record_transition_or_restore(
637 &mut self.state,
638 session_id,
639 &destroying,
640 "persist stopped state after target cleanup",
641 &persist,
642 )?;
643 Ok(deferred)
644 }
645
646 pub fn cleanup_stopped_target(
650 &mut self,
651 session_id: &str,
652 executor: &impl CommandExecutor,
653 ) -> Result<()> {
654 self.cleanup_stopped_target_with(
655 session_id,
656 executor,
657 crate::database::save_lifecycle_session,
658 )
659 }
660
661 fn cleanup_stopped_target_with(
662 &mut self,
663 session_id: &str,
664 executor: &impl CommandExecutor,
665 persist: impl Fn(&SessionRecord) -> Result<()>,
666 ) -> Result<()> {
667 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
668 let _target_guard = target_mutex.lock().map_err(|_| {
669 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
670 })?;
671 let previous = self
672 .state
673 .sessions
674 .get(session_id)
675 .with_context(|| format!("unknown session {session_id}"))?
676 .clone();
677 ensure!(
678 previous.state == SessionState::Stopped,
679 "refusing deferred cleanup for active session {session_id}"
680 );
681 let Some(locator) = previous.target.as_ref() else {
682 return Ok(());
683 };
684 let backend = backend_locator(locator, &previous, &self.config)?;
685 ensure!(
686 targets::quiesce_plan(&backend, session_id)?.is_some(),
687 "session {session_id} retained a non-Podman target after stopping"
688 );
689 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
690 let record = self.state.sessions.get_mut(session_id).unwrap();
691 record.updated_at = now();
692 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
693 let persisted = persist_session_record_transition_or_restore(
694 &mut self.state,
695 session_id,
696 &previous,
697 "persist deferred target cleanup failure",
698 &persist,
699 );
700 return match persisted {
701 Ok(()) => Err(error),
702 Err(persist_error) => Err(error.context(format!(
703 "also failed to persist deferred target cleanup failure: {persist_error:#}"
704 ))),
705 };
706 }
707 let record = self.state.sessions.get_mut(session_id).unwrap();
708 record.target = None;
709 record.updated_at = now();
710 record.last_error = None;
711 persist_session_record_transition_or_restore(
712 &mut self.state,
713 session_id,
714 &previous,
715 "persist completion of deferred Podman target cleanup",
716 &persist,
717 )
718 }
719
720 pub fn force_stop(
723 &mut self,
724 session_id: &str,
725 executor: &impl CommandExecutor,
726 ) -> Result<bool> {
727 self.force_stop_with(
728 session_id,
729 executor,
730 crate::database::save_lifecycle_session,
731 )
732 }
733
734 fn force_stop_with(
735 &mut self,
736 session_id: &str,
737 executor: &impl CommandExecutor,
738 persist: impl Fn(&SessionRecord) -> Result<()>,
739 ) -> Result<bool> {
740 let session = self
741 .state
742 .sessions
743 .get(session_id)
744 .with_context(|| format!("unknown session {session_id}"))?
745 .clone();
746 ensure!(
747 session.state.is_active(),
748 "session {session_id} is already inactive"
749 );
750 let checkpoint = session
751 .checkpoint
752 .as_ref()
753 .context("force stop requires an existing recovery archive")?;
754 {
757 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
758 verify_installed_checkpoint_gate(session_id, checkpoint)
759 .context("verify the recovery archive before force stopping")?;
760 }
761 self.stop_target_and_settle(session_id, &session, executor, &persist)
762 }
763
764 fn stop_target_and_settle(
769 &mut self,
770 session_id: &str,
771 session: &SessionRecord,
772 executor: &impl CommandExecutor,
773 persist: &impl Fn(&SessionRecord) -> Result<()>,
774 ) -> Result<bool> {
775 let mut deferred = false;
776 if let Some(locator) = &session.target {
777 let backend = backend_locator(locator, session, &self.config)?;
778 if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
779 plan.execute(executor)?;
780 deferred = true;
781 } else {
782 execute_target_cleanup(&backend, session_id, executor)?;
783 }
784 }
785 if let Some(worktree) = &session.managed_worktree {
786 retire_managed_worktree(executor, worktree)
787 .context("retire managed raw-session worktree after stopping the target")?;
788 }
789 let record = self.state.sessions.get_mut(session_id).unwrap();
790 record.state = SessionState::Stopped;
791 if !deferred {
792 record.target = None;
793 }
794 record.updated_at = now();
795 record.last_error = None;
796 record.last_checkpoint_error = None;
797 persist_session_record_transition_or_restore(
798 &mut self.state,
799 session_id,
800 session,
801 "persist stopped state after tearing down the current target",
802 persist,
803 )?;
804 Ok(deferred)
805 }
806
807 pub fn destroy_session_controlled(
811 &mut self,
812 session_id: &str,
813 executor: &impl CommandExecutor,
814 ) -> Result<()> {
815 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
816 }
817
818 pub fn destroy_session_controlled_with(
826 &mut self,
827 session_id: &str,
828 executor: &impl CommandExecutor,
829 branch: BranchDisposition,
830 ) -> Result<()> {
831 self.destroy_session_controlled_with_checkout(
832 session_id,
833 executor,
834 branch,
835 CheckoutDisposition::Remove,
836 )
837 .map(|_| ())
838 }
839
840 pub(crate) fn destroy_session_controlled_with_checkout(
846 &mut self,
847 session_id: &str,
848 executor: &impl CommandExecutor,
849 branch: BranchDisposition,
850 checkout: CheckoutDisposition,
851 ) -> Result<Option<PathBuf>> {
852 let session = self
853 .state
854 .sessions
855 .get(session_id)
856 .with_context(|| format!("unknown session {session_id}"))?
857 .clone();
858 if session.state.is_active() {
859 bail!("refusing to destroy active session {session_id}");
860 }
861 let mut retained_checkout = None;
862 if let Some(worktree) = &session.managed_worktree {
863 let keep = checkout == CheckoutDisposition::KeepWhenDirty
864 && managed_worktree_checkout_is_dirty(executor, worktree)
865 .context("check the managed raw-session worktree for uncommitted changes")?;
866 if keep {
867 retained_checkout = Some(worktree.worktree_root.clone());
871 } else {
872 cleanup_managed_worktree(executor, worktree, branch)
873 .context("remove managed raw-session worktree")?;
874 }
875 }
876 if let Some(checkpoint) = &session.checkpoint
877 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
878 && error.kind() != std::io::ErrorKind::NotFound
879 {
880 return Err(error).with_context(|| {
881 format!(
882 "remove session recovery archive {}",
883 checkpoint.archive_path.display()
884 )
885 });
886 }
887 mj_core::attachment::AttachmentStore::controller(session_id)?
888 .remove_session_data()
889 .context("remove session image attachments")?;
890 crate::database::delete_session(session_id)
891 .context("destroy stopped session in database")?;
892 self.state.subagents.remove(session_id);
893 self.state.destroy_stopped_session(session_id)?;
894 Ok(retained_checkout)
895 }
896
897 pub fn force_destroy_session(
909 &mut self,
910 session_id: &str,
911 executor: &impl CommandExecutor,
912 branch: BranchDisposition,
913 ) -> Result<()> {
914 self.force_destroy_session_with(
915 session_id,
916 executor,
917 branch,
918 crate::database::delete_session,
919 )
920 }
921
922 fn force_destroy_session_with(
923 &mut self,
924 session_id: &str,
925 executor: &impl CommandExecutor,
926 branch: BranchDisposition,
927 delete: impl Fn(&str) -> Result<()>,
928 ) -> Result<()> {
929 let session = self
930 .state
931 .sessions
932 .get(session_id)
933 .with_context(|| format!("unknown session {session_id}"))?
934 .clone();
935 if let Some(locator) = &session.target {
939 let backend = backend_locator(locator, &session, &self.config)?;
940 if self.state.subagents.contains_key(session_id) {
941 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
942 } else {
943 execute_target_cleanup(&backend, session_id, executor)?;
944 }
945 }
946 if let Some(worktree) = &session.managed_worktree {
947 cleanup_managed_worktree(executor, worktree, branch)
948 .context("remove managed raw-session worktree")?;
949 }
950 if let Some(checkpoint) = &session.checkpoint
951 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
952 && error.kind() != std::io::ErrorKind::NotFound
953 {
954 return Err(error).with_context(|| {
955 format!(
956 "remove session recovery archive {}",
957 checkpoint.archive_path.display()
958 )
959 });
960 }
961 mj_core::attachment::AttachmentStore::controller(session_id)?
962 .remove_session_data()
963 .context("remove session image attachments")?;
964 delete(session_id).context("force destroy session in database")?;
965 self.state.subagents.remove(session_id);
966 self.state.destroy_session_force(session_id)?;
967 Ok(())
968 }
969}
970
971pub fn has_nothing_to_checkpoint(session: &SessionRecord) -> bool {
980 match session.state {
981 SessionState::Provisioning => true,
982 SessionState::Closing
983 | SessionState::Destroying
984 | SessionState::Error
985 | SessionState::Lost => session.target.is_none(),
986 SessionState::Running
987 | SessionState::Disconnected
988 | SessionState::Checkpointing
989 | SessionState::Stopped
990 | SessionState::DestroyedWithDataLoss => false,
991 }
992}
993
994fn execute_target_cleanup(
995 backend: &targets::TargetLocator,
996 session_id: &str,
997 executor: &impl CommandExecutor,
998) -> Result<()> {
999 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1000 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1001 Ok(true) => {
1002 tracing::warn!(
1003 session_id,
1004 error = format!("{cleanup_error:#}"),
1005 "target cleanup command failed, but the target was confirmed absent"
1006 );
1007 }
1008 Ok(false) => {
1009 tracing::error!(
1010 session_id,
1011 error = format!("{cleanup_error:#}"),
1012 "target cleanup failed and the target is still present"
1013 );
1014 return Err(cleanup_error);
1015 }
1016 Err(probe_error) => {
1017 tracing::error!(
1018 session_id,
1019 cleanup_error = format!("{cleanup_error:#}"),
1020 probe_error = format!("{probe_error:#}"),
1021 "target cleanup failed and exact absence could not be confirmed"
1022 );
1023 return Err(cleanup_error.context(format!(
1024 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1025 )));
1026 }
1027 }
1028 }
1029 Ok(())
1030}
1031
1032fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1033 record.state = SessionState::Closing;
1034 record.updated_at = updated_at;
1035 record.last_checkpoint_error = None;
1036}
1037
1038fn apply_close_checkpoint_failure(
1045 record: &mut SessionRecord,
1046 previous: &SessionRecord,
1047 error: &anyhow::Error,
1048 updated_at: String,
1049) {
1050 if WorkerRestartLeftNoWorker::marks(error) {
1051 record.state = SessionState::Error;
1052 record.last_error = Some(format!(
1053 "close failed and left the session without a live worker; retry the close, \
1054 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1055 ));
1056 } else {
1057 record.state = if previous.state == SessionState::Closing {
1058 SessionState::Running
1059 } else {
1060 previous.state
1061 };
1062 }
1063 record.last_checkpoint_error = Some(format!("{error:#}"));
1064 record.updated_at = updated_at;
1065}
1066
1067fn apply_interrupted_close_error(
1068 record: &mut SessionRecord,
1069 error: &anyhow::Error,
1070 updated_at: &str,
1071) {
1072 let destroying = record.state == SessionState::Destroying;
1073 if !destroying {
1074 record.state = SessionState::Closing;
1075 }
1076 record.updated_at = updated_at.to_owned();
1077 record.last_error = Some(if destroying {
1078 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1079 } else {
1080 format!("close is safely resumable from its verified checkpoint: {error:#}")
1081 });
1082}
1083
1084#[cfg(test)]
1085mod tests;