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, ManagedCheckoutKind, 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
64pub type BeforeClose = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>;
71
72impl Controller {
73 pub async fn suspend_session(&mut self, session_id: &str) -> Result<()> {
78 self.suspend_session_controlled(session_id, &ProcessExecutor)
79 .await
80 }
81
82 pub async fn suspend_session_controlled(
83 &mut self,
84 session_id: &str,
85 executor: &(impl CommandExecutor + Sync),
86 ) -> Result<()> {
87 if self
88 .suspend_session_controlled_with_manager(
89 session_id,
90 executor,
91 None,
92 None,
93 SourceTargetDisposition::Destroy,
94 true,
95 None,
96 )
97 .await?
98 {
99 self.cleanup_stopped_target(session_id, executor)?;
100 }
101 Ok(())
102 }
103
104 pub async fn suspend_session_managed_controlled(
105 &mut self,
106 session_id: &str,
107 executor: &(impl CommandExecutor + Sync),
108 manager: &SessionManagerControl,
109 acknowledge_unpublished_work: bool,
110 before_close: Option<BeforeClose>,
111 ) -> Result<bool> {
112 self.suspend_session_controlled_with_manager(
113 session_id,
114 executor,
115 Some(manager),
116 None,
117 SourceTargetDisposition::Destroy,
118 acknowledge_unpublished_work,
119 before_close,
120 )
121 .await
122 }
123
124 pub(super) async fn suspend_session_for_move(
125 &mut self,
126 session_id: &str,
127 executor: &(impl CommandExecutor + Sync),
128 manager: &SessionManagerControl,
129 operation: &mut mj_core::state::MoveOperation,
130 preparation: Option<&mj_core::state::MovePreparation>,
131 disposition: SourceTargetDisposition,
132 ) -> Result<bool> {
133 self.prepare_move_source_checkpoint(session_id, executor, manager, operation)
134 .await?;
135 self.suspend_session_controlled_with_manager(
136 session_id,
137 executor,
138 Some(manager),
139 Some((operation, preparation)),
140 disposition,
141 true,
142 None,
143 )
144 .await
145 }
146
147 #[allow(clippy::too_many_arguments)]
148 async fn suspend_session_controlled_with_manager(
149 &mut self,
150 session_id: &str,
151 executor: &(impl CommandExecutor + Sync),
152 manager: Option<&SessionManagerControl>,
153 move_intent: Option<(
154 &mut mj_core::state::MoveOperation,
155 Option<&mj_core::state::MovePreparation>,
156 )>,
157 disposition: SourceTargetDisposition,
158 acknowledge_unpublished_work: bool,
159 before_close: Option<BeforeClose>,
160 ) -> Result<bool> {
161 let previous = self
162 .state
163 .sessions
164 .get(session_id)
165 .with_context(|| format!("unknown session {session_id}"))?
166 .clone();
167 let record = self.state.sessions.get_mut(session_id).unwrap();
168 apply_close_checkpoint_started(record, now());
172 self.persist_session_transition_or_restore(
173 session_id,
174 &previous,
175 "persist closing state before checkpointing the session",
176 )?;
177
178 let mut latched = match self
181 .checkpoint_session_latched(
182 session_id,
183 executor,
184 manager,
185 LatchExclusivity::HoldThroughClose,
186 CheckpointExportPolicy::ReuseUnchangedArchive,
187 )
188 .await
189 {
190 Ok(latched) => latched,
191 Err(error) => {
192 let record = self.state.sessions.get_mut(session_id).unwrap();
193 apply_close_checkpoint_failure(record, &previous, &error, now());
197 return Err(
198 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
199 );
200 }
201 };
202
203 let artifact = latched.artifact.clone();
204 let publication = match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
205 Some(ManagedCheckoutKind::Clone) => Some(super::publication::assess_clone_checkpoint(
206 &previous,
207 &artifact.metadata,
208 )),
209 Some(ManagedCheckoutKind::Worktree) => None,
210 None if previous.project_directory.is_none() => {
211 Some(super::publication::assess_network_checkpoint(
212 &previous,
213 &artifact.metadata,
214 &self.config,
215 ))
216 }
217 None => None,
218 };
219 if !acknowledge_unpublished_work
220 && publication
221 .as_ref()
222 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
223 {
224 latched.relay.cancel_abandoned_barrier().await?;
225 let mut restored = previous.clone();
229 restored.state = state_after_unsealed_close(&previous);
230 restored.updated_at = now();
231 self.state.sessions.insert(session_id.to_owned(), restored);
232 self.persist_session_transition_or_restore(
233 session_id,
234 &previous,
235 "restore session after unpublished-work preflight",
236 )?;
237 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
238 return Err(mj_core::refusal::Refusal::precondition(
239 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
240 ).into());
241 }
242 let record = self.state.sessions.get_mut(session_id).unwrap();
243 record.state = SessionState::Closing;
244 record.native_session_id = Some(artifact.native_session_id.clone());
245 record.checkpoint = Some(artifact.metadata.clone());
246 record.publication = publication;
247 record.updated_at = now();
248 record.last_error = None;
249 record.last_checkpoint_error = None;
250 self.persist_checkpoint_transition_or_restore(
251 session_id,
252 &previous,
253 "persist verified checkpoint and closing state before sealing the relay",
254 )?;
255 if let Some((operation, preparation)) = move_intent {
256 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
259 let record = self.state.sessions.get_mut(session_id).unwrap();
260 record.state = previous.state;
261 record.last_error = Some(format!("{error:#}"));
262 self.persist_session_transition_or_restore(
263 session_id,
264 &previous,
265 "restore source after move preflight failure",
266 )?;
267 return Err(error);
268 }
269 operation.checkpoint = Some(artifact.metadata.clone());
270 operation.updated_at = now();
271 crate::database::save_move_operation(operation)?;
272 }
273 if let Some(before_close) = before_close
274 && let Err(error) = before_close.await
275 {
276 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
279 tracing::warn!(
280 session_id,
281 error = format!("{cancel:#}"),
282 "could not release the checkpoint barrier of a close that did not proceed"
283 );
284 }
285 let record = self.state.sessions.get_mut(session_id).unwrap();
286 record.state = state_after_unsealed_close(&previous);
287 record.last_error = Some(format!("{error:#}"));
288 record.updated_at = now();
289 self.persist_session_transition_or_restore(
290 session_id,
291 &previous,
292 "restore a session whose close did not proceed past its checkpoint",
293 )?;
294 return Err(error);
295 }
296 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
297 release_projection_behind_checkpoint(session_id, &artifact.metadata);
300
301 let close_command_id = new_command_id("close")?;
302 let barrier_command_id = latched.barrier_command_id.clone();
303 let close_result = {
304 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
305 latched
306 .relay
307 .connection_mut()
308 .submit(
309 close_command_id,
310 RelayCommand::Close {
311 barrier_command_id: barrier_command_id.clone(),
312 expected: latched.cursor.clone(),
313 },
314 )
315 .await
316 };
317 if let Err(error) = close_result {
318 self.record_interrupted_close(session_id, &error)?;
319 return Err(error.context("seal verified checkpoint for close"));
320 }
321 let close_result = {
322 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
323 latched
324 .relay
325 .connection_mut()
326 .submit(
327 new_command_id("checkpoint-complete")?,
328 RelayCommand::CompleteCheckpoint { barrier_command_id },
329 )
330 .await
331 };
332 if let Err(error) = close_result {
333 self.record_interrupted_close(session_id, &error)?;
334 return Err(error.context("release verified close checkpoint"));
335 }
336 let close_result = {
337 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
338 wait_for_relay_closed(latched.relay.connection_mut()).await
339 };
340 if let Err(error) = close_result {
341 self.record_interrupted_close(session_id, &error)?;
342 return Err(error);
343 }
344 latched.relay.release();
345
346 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
347 return Ok(false);
353 }
354 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
355 Ok(deferred) => Ok(deferred),
356 Err(error) => {
357 self.record_interrupted_close(session_id, &error)?;
358 Err(error)
359 }
360 }
361 }
362
363 pub async fn recover_interrupted_close_managed(
375 &mut self,
376 session_id: &str,
377 executor: &(impl CommandExecutor + Sync),
378 manager: &SessionManagerControl,
379 acknowledge_unpublished_work: bool,
380 before_close: Option<BeforeClose>,
381 ) -> Result<bool> {
382 let (state, verified) = {
383 let session = self
384 .state
385 .sessions
386 .get(session_id)
387 .with_context(|| format!("unknown session {session_id}"))?;
388 ensure!(
389 matches!(
390 session.state,
391 SessionState::Closing | SessionState::Destroying
392 ),
393 "session {session_id} has no interrupted close to recover"
394 );
395 (session.state, session.checkpoint.clone())
396 };
397 if state == SessionState::Destroying {
398 let verified = verified.context("destroying session has no verified checkpoint")?;
399 if let Some(before_close) = before_close {
400 before_close.await?;
401 }
402 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
403 }
404 ensure!(
405 state == SessionState::Closing,
406 "session {session_id} has no relay close to recover"
407 );
408 let handle = manager
409 .wait_for_session(session_id, Duration::from_secs(5))
410 .await?;
411 let mut lease = handle.lease_connection().await?;
412 let execution = lease.connection_mut().sync().await?.operational.execution;
413 match execution {
414 RelayExecutionState::Closed => {}
415 RelayExecutionState::Closing => {
416 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
417 wait_for_relay_closed(lease.connection_mut()).await?;
418 }
419 RelayExecutionState::Idle | RelayExecutionState::Running => {
420 lease.release();
421 return self
422 .suspend_session_controlled_with_manager(
423 session_id,
424 executor,
425 Some(manager),
426 None,
427 SourceTargetDisposition::Destroy,
428 acknowledge_unpublished_work,
429 before_close,
430 )
431 .await;
432 }
433 }
434 lease.release();
435 let verified = verified.context("closed relay has no verified checkpoint")?;
436 if let Some(before_close) = before_close {
438 before_close.await?;
439 }
440 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
441 }
442
443 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
450 self.fail_interrupted_lifecycle_with(
451 session_id,
452 cause,
453 crate::database::save_lifecycle_session,
454 )
455 }
456
457 fn fail_interrupted_lifecycle_with(
458 &mut self,
459 session_id: &str,
460 cause: &str,
461 persist: impl Fn(&SessionRecord) -> Result<()>,
462 ) -> Result<bool> {
463 let Some(record) = self.state.sessions.get_mut(session_id) else {
464 return Ok(false);
465 };
466 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
467 return Ok(false);
468 }
469 let previous = record.clone();
470 record.state = SessionState::Error;
471 record.updated_at = now();
472 record.last_error = Some(cause.to_owned());
473 persist_session_record_transition_or_restore(
474 &mut self.state,
475 session_id,
476 &previous,
477 "persist the failure of an interrupted lifecycle state",
478 &persist,
479 )?;
480 Ok(true)
481 }
482
483 pub fn fail_unready_session(
492 &mut self,
493 session_id: &str,
494 cause: &str,
495 observed_updated_at: &str,
496 ) -> Result<bool> {
497 self.fail_unready_session_with(
498 session_id,
499 cause,
500 observed_updated_at,
501 crate::database::save_lifecycle_session,
502 )
503 }
504
505 fn fail_unready_session_with(
506 &mut self,
507 session_id: &str,
508 cause: &str,
509 observed_updated_at: &str,
510 persist: impl Fn(&SessionRecord) -> Result<()>,
511 ) -> Result<bool> {
512 let Some(record) = self.state.sessions.get_mut(session_id) else {
513 return Ok(false);
514 };
515 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
516 return Ok(false);
517 }
518 let previous = record.clone();
519 record.state = SessionState::Error;
520 record.updated_at = now();
521 record.last_error = Some(cause.to_owned());
522 persist_session_record_transition_or_restore(
523 &mut self.state,
524 session_id,
525 &previous,
526 "persist a session whose harness never became usable",
527 &persist,
528 )?;
529 Ok(true)
530 }
531
532 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
536 let Some(record) = self.state.sessions.get(session_id) else {
537 return Ok(false);
538 };
539 let previous = record.clone();
542 let record = self.state.sessions.get_mut(session_id).unwrap();
543 record.last_error = Some(cause.to_owned());
544 record.updated_at = now();
545 persist_session_record_transition_or_restore(
546 &mut self.state,
547 session_id,
548 &previous,
549 "persist the reason a close did not finish",
550 &crate::database::save_lifecycle_session,
551 )?;
552 Ok(true)
553 }
554
555 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
560 let Some(record) = self.state.sessions.get(session_id) else {
561 return Ok(false);
562 };
563 if record.public_error().is_none() {
564 return Ok(false);
565 }
566 let previous = record.clone();
567 let record = self.state.sessions.get_mut(session_id).unwrap();
568 record.last_error = None;
569 record.updated_at = now();
570 persist_session_record_transition_or_restore(
571 &mut self.state,
572 session_id,
573 &previous,
574 "clear the reason a close did not finish",
575 &crate::database::save_lifecycle_session,
576 )?;
577 Ok(true)
578 }
579
580 pub fn suspend_session_without_checkpoint(
593 &mut self,
594 session_id: &str,
595 executor: &impl CommandExecutor,
596 ) -> Result<bool> {
597 self.suspend_session_without_checkpoint_with(
598 session_id,
599 executor,
600 crate::database::save_lifecycle_session,
601 )
602 }
603
604 fn suspend_session_without_checkpoint_with(
605 &mut self,
606 session_id: &str,
607 executor: &impl CommandExecutor,
608 persist: impl Fn(&SessionRecord) -> Result<()>,
609 ) -> Result<bool> {
610 let session = self
611 .state
612 .sessions
613 .get(session_id)
614 .with_context(|| format!("unknown session {session_id}"))?
615 .clone();
616 ensure!(
617 has_nothing_to_checkpoint(&session),
618 "session {session_id} has a workspace to checkpoint; suspend it instead"
619 );
620 self.stop_target_and_settle(session_id, &session, executor, &persist)
621 }
622
623 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
624 let record = self.state.sessions.get_mut(session_id).unwrap();
625 apply_interrupted_close_error(record, error, &now());
626 self.persist_session_state(session_id)
627 }
628
629 fn destroy_after_verified_checkpoint(
632 &mut self,
633 session_id: &str,
634 verified: &CheckpointMetadata,
635 executor: &impl CommandExecutor,
636 ) -> Result<bool> {
637 self.destroy_after_verified_checkpoint_with(
638 session_id,
639 verified,
640 executor,
641 crate::database::save_lifecycle_session,
642 )
643 }
644
645 fn destroy_after_verified_checkpoint_with(
646 &mut self,
647 session_id: &str,
648 verified: &CheckpointMetadata,
649 executor: &impl CommandExecutor,
650 persist: impl Fn(&SessionRecord) -> Result<()>,
651 ) -> Result<bool> {
652 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
653 let _target_guard = target_mutex.lock().map_err(|_| {
654 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
655 })?;
656 let session = self
657 .state
658 .sessions
659 .get(session_id)
660 .with_context(|| format!("unknown session {session_id}"))?
661 .clone();
662 ensure!(
663 matches!(
664 session.state,
665 SessionState::Closing | SessionState::Destroying
666 ),
667 "refusing to destroy session {session_id}: it is not closing or destroying"
668 );
669 ensure!(
670 session.checkpoint.as_ref() == Some(verified),
671 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
672 );
673 if session.state == SessionState::Closing {
674 let record = self.state.sessions.get_mut(session_id).unwrap();
675 record.state = SessionState::Destroying;
676 record.updated_at = now();
677 record.last_error = None;
678 persist_session_record_transition_or_restore(
679 &mut self.state,
680 session_id,
681 &session,
682 "persist destroying state before target cleanup",
683 &persist,
684 )?;
685 }
686
687 let destroying = self
688 .state
689 .sessions
690 .get(session_id)
691 .expect("destroying session disappeared")
692 .clone();
693 {
694 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
695 verify_installed_checkpoint_gate(session_id, verified)?;
696 }
697 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
702 tracing::warn!(
703 session_id,
704 error = format!("{error:#}"),
705 "could not record that the second-opinion conversation ends with this target"
706 );
707 }
708 let locator = destroying
709 .target
710 .as_ref()
711 .context("session has no target")?;
712 let backend = backend_locator(locator, &destroying, &self.config)?;
713 let deferred = if self.state.subagents.contains_key(session_id) {
714 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
715 false
716 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
717 plan.execute(executor)?;
718 true
719 } else {
720 execute_target_cleanup(&backend, session_id, executor)?;
721 false
722 };
723 if let Some(worktree) = &destroying.managed_worktree {
724 retire_managed_worktree(executor, worktree)
725 .context("retire managed raw-session worktree after verified close")?;
726 }
727 let record = self.state.sessions.get_mut(session_id).unwrap();
728 record.state = SessionState::Stopped;
729 if !deferred {
730 record.target = None;
731 }
732 record.updated_at = now();
733 record.last_error = None;
734 persist_session_record_transition_or_restore(
735 &mut self.state,
736 session_id,
737 &destroying,
738 "persist stopped state after target cleanup",
739 &persist,
740 )?;
741 Ok(deferred)
742 }
743
744 pub fn cleanup_stopped_target(
748 &mut self,
749 session_id: &str,
750 executor: &impl CommandExecutor,
751 ) -> Result<()> {
752 self.cleanup_stopped_target_with(
753 session_id,
754 executor,
755 crate::database::save_lifecycle_session,
756 )
757 }
758
759 fn cleanup_stopped_target_with(
760 &mut self,
761 session_id: &str,
762 executor: &impl CommandExecutor,
763 persist: impl Fn(&SessionRecord) -> Result<()>,
764 ) -> Result<()> {
765 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
766 let _target_guard = target_mutex.lock().map_err(|_| {
767 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
768 })?;
769 let previous = self
770 .state
771 .sessions
772 .get(session_id)
773 .with_context(|| format!("unknown session {session_id}"))?
774 .clone();
775 ensure!(
776 previous.state == SessionState::Stopped,
777 "refusing deferred cleanup for active session {session_id}"
778 );
779 let Some(locator) = previous.target.as_ref() else {
780 return Ok(());
781 };
782 let backend = backend_locator(locator, &previous, &self.config)?;
783 ensure!(
784 targets::quiesce_plan(&backend, session_id)?.is_some(),
785 "session {session_id} retained a non-Podman target after stopping"
786 );
787 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
788 let record = self.state.sessions.get_mut(session_id).unwrap();
789 record.updated_at = now();
790 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
791 let persisted = persist_session_record_transition_or_restore(
792 &mut self.state,
793 session_id,
794 &previous,
795 "persist deferred target cleanup failure",
796 &persist,
797 );
798 return match persisted {
799 Ok(()) => Err(error),
800 Err(persist_error) => Err(error.context(format!(
801 "also failed to persist deferred target cleanup failure: {persist_error:#}"
802 ))),
803 };
804 }
805 let record = self.state.sessions.get_mut(session_id).unwrap();
806 record.target = None;
807 record.updated_at = now();
808 record.last_error = None;
809 persist_session_record_transition_or_restore(
810 &mut self.state,
811 session_id,
812 &previous,
813 "persist completion of deferred Podman target cleanup",
814 &persist,
815 )
816 }
817
818 pub fn force_stop(
821 &mut self,
822 session_id: &str,
823 executor: &impl CommandExecutor,
824 ) -> Result<bool> {
825 self.force_stop_with(
826 session_id,
827 executor,
828 crate::database::save_lifecycle_session,
829 )
830 }
831
832 fn force_stop_with(
833 &mut self,
834 session_id: &str,
835 executor: &impl CommandExecutor,
836 persist: impl Fn(&SessionRecord) -> Result<()>,
837 ) -> Result<bool> {
838 let session = self
839 .state
840 .sessions
841 .get(session_id)
842 .with_context(|| format!("unknown session {session_id}"))?
843 .clone();
844 ensure!(
845 session.state.is_active(),
846 "session {session_id} is already inactive"
847 );
848 let checkpoint = session
849 .checkpoint
850 .as_ref()
851 .context("force stop requires an existing recovery archive")?;
852 {
855 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
856 verify_installed_checkpoint_gate(session_id, checkpoint)
857 .context("verify the recovery archive before force stopping")?;
858 }
859 self.stop_target_and_settle(session_id, &session, executor, &persist)
860 }
861
862 fn stop_target_and_settle(
867 &mut self,
868 session_id: &str,
869 session: &SessionRecord,
870 executor: &impl CommandExecutor,
871 persist: &impl Fn(&SessionRecord) -> Result<()>,
872 ) -> Result<bool> {
873 let mut deferred = false;
874 if let Some(locator) = &session.target {
875 let backend = backend_locator(locator, session, &self.config)?;
876 if self.state.subagents.contains_key(session_id) {
880 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
881 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
882 plan.execute(executor)?;
883 deferred = true;
884 } else {
885 execute_target_cleanup(&backend, session_id, executor)?;
886 }
887 }
888 if let Some(worktree) = &session.managed_worktree {
889 retire_managed_worktree(executor, worktree)
890 .context("retire managed raw-session worktree after stopping the target")?;
891 }
892 let record = self.state.sessions.get_mut(session_id).unwrap();
893 record.state = SessionState::Stopped;
894 if !deferred {
895 record.target = None;
896 }
897 record.updated_at = now();
898 record.last_error = None;
899 record.last_checkpoint_error = None;
900 persist_session_record_transition_or_restore(
901 &mut self.state,
902 session_id,
903 session,
904 "persist stopped state after tearing down the current target",
905 persist,
906 )?;
907 Ok(deferred)
908 }
909
910 pub fn destroy_session_controlled(
914 &mut self,
915 session_id: &str,
916 executor: &impl CommandExecutor,
917 ) -> Result<()> {
918 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
919 }
920
921 pub fn destroy_session_controlled_with(
929 &mut self,
930 session_id: &str,
931 executor: &impl CommandExecutor,
932 branch: BranchDisposition,
933 ) -> Result<()> {
934 self.destroy_session_controlled_with_checkout(
935 session_id,
936 executor,
937 branch,
938 CheckoutDisposition::Remove,
939 )
940 .map(|_| ())
941 }
942
943 pub(crate) fn destroy_session_controlled_with_checkout(
949 &mut self,
950 session_id: &str,
951 executor: &impl CommandExecutor,
952 branch: BranchDisposition,
953 checkout: CheckoutDisposition,
954 ) -> Result<Option<PathBuf>> {
955 let session = self
956 .state
957 .sessions
958 .get(session_id)
959 .with_context(|| format!("unknown session {session_id}"))?
960 .clone();
961 if session.state.is_active() {
962 bail!("refusing to destroy active session {session_id}");
963 }
964 let mut retained_checkout = None;
965 if let Some(worktree) = &session.managed_worktree {
966 let keep = checkout == CheckoutDisposition::KeepWhenDirty
967 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
968 || managed_worktree_checkout_is_dirty(executor, worktree).context(
969 "check the managed raw-session worktree for uncommitted changes",
970 )?);
971 if keep {
972 retained_checkout = Some(worktree.worktree_root.clone());
976 } else {
977 cleanup_managed_worktree(executor, worktree, branch)
978 .context("remove managed raw-session worktree")?;
979 }
980 }
981 if let Some(checkpoint) = &session.checkpoint
982 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
983 && error.kind() != std::io::ErrorKind::NotFound
984 {
985 return Err(error).with_context(|| {
986 format!(
987 "remove session recovery archive {}",
988 checkpoint.archive_path.display()
989 )
990 });
991 }
992 mj_core::attachment::AttachmentStore::controller(session_id)?
993 .remove_session_data()
994 .context("remove session image attachments")?;
995 crate::database::delete_session(session_id)
996 .context("destroy stopped session in database")?;
997 self.state.subagents.remove(session_id);
998 self.state.destroy_stopped_session(session_id)?;
999 Ok(retained_checkout)
1000 }
1001
1002 pub fn force_destroy_session(
1014 &mut self,
1015 session_id: &str,
1016 executor: &impl CommandExecutor,
1017 branch: BranchDisposition,
1018 ) -> Result<()> {
1019 self.force_destroy_session_with(
1020 session_id,
1021 executor,
1022 branch,
1023 crate::database::delete_session,
1024 )
1025 }
1026
1027 fn force_destroy_session_with(
1028 &mut self,
1029 session_id: &str,
1030 executor: &impl CommandExecutor,
1031 branch: BranchDisposition,
1032 delete: impl Fn(&str) -> Result<()>,
1033 ) -> Result<()> {
1034 let session = self
1035 .state
1036 .sessions
1037 .get(session_id)
1038 .with_context(|| format!("unknown session {session_id}"))?
1039 .clone();
1040 if let Some(locator) = &session.target {
1044 let backend = backend_locator(locator, &session, &self.config)?;
1045 if self.state.subagents.contains_key(session_id) {
1046 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
1047 } else {
1048 execute_target_cleanup(&backend, session_id, executor)?;
1049 }
1050 }
1051 if let Some(worktree) = &session.managed_worktree {
1052 cleanup_managed_worktree(executor, worktree, branch)
1053 .context("remove managed raw-session worktree")?;
1054 }
1055 if let Some(checkpoint) = &session.checkpoint
1056 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1057 && error.kind() != std::io::ErrorKind::NotFound
1058 {
1059 return Err(error).with_context(|| {
1060 format!(
1061 "remove session recovery archive {}",
1062 checkpoint.archive_path.display()
1063 )
1064 });
1065 }
1066 mj_core::attachment::AttachmentStore::controller(session_id)?
1067 .remove_session_data()
1068 .context("remove session image attachments")?;
1069 delete(session_id).context("force destroy session in database")?;
1070 self.state.subagents.remove(session_id);
1071 self.state.destroy_session_force(session_id)?;
1072 Ok(())
1073 }
1074}
1075
1076pub fn has_nothing_to_checkpoint(session: &SessionRecord) -> bool {
1085 match session.state {
1086 SessionState::Provisioning => true,
1087 SessionState::Parked => true,
1091 SessionState::Closing
1092 | SessionState::Destroying
1093 | SessionState::Error
1094 | SessionState::Lost => session.target.is_none(),
1095 SessionState::Running
1096 | SessionState::Disconnected
1097 | SessionState::Checkpointing
1098 | SessionState::Stopped
1099 | SessionState::DestroyedWithDataLoss => false,
1100 }
1101}
1102
1103fn execute_target_cleanup(
1104 backend: &targets::TargetLocator,
1105 session_id: &str,
1106 executor: &impl CommandExecutor,
1107) -> Result<()> {
1108 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1109 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1110 Ok(true) => {
1111 tracing::warn!(
1112 session_id,
1113 error = format!("{cleanup_error:#}"),
1114 "target cleanup command failed, but the target was confirmed absent"
1115 );
1116 }
1117 Ok(false) => {
1118 tracing::error!(
1119 session_id,
1120 error = format!("{cleanup_error:#}"),
1121 "target cleanup failed and the target is still present"
1122 );
1123 return Err(cleanup_error);
1124 }
1125 Err(probe_error) => {
1126 tracing::error!(
1127 session_id,
1128 cleanup_error = format!("{cleanup_error:#}"),
1129 probe_error = format!("{probe_error:#}"),
1130 "target cleanup failed and exact absence could not be confirmed"
1131 );
1132 return Err(cleanup_error.context(format!(
1133 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1134 )));
1135 }
1136 }
1137 }
1138 Ok(())
1139}
1140
1141fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1142 record.state = SessionState::Closing;
1143 record.updated_at = updated_at;
1144 record.last_checkpoint_error = None;
1145}
1146
1147fn apply_close_checkpoint_failure(
1154 record: &mut SessionRecord,
1155 previous: &SessionRecord,
1156 error: &anyhow::Error,
1157 updated_at: String,
1158) {
1159 if WorkerRestartLeftNoWorker::marks(error) {
1160 record.state = SessionState::Error;
1161 record.last_error = Some(format!(
1162 "close failed and left the session without a live worker; retry the close, \
1163 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1164 ));
1165 } else {
1166 record.state = state_after_unsealed_close(previous);
1167 }
1168 record.last_checkpoint_error = Some(format!("{error:#}"));
1169 record.updated_at = updated_at;
1170}
1171
1172fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1177 if previous.state == SessionState::Closing {
1178 SessionState::Running
1179 } else {
1180 previous.state
1181 }
1182}
1183
1184fn apply_interrupted_close_error(
1185 record: &mut SessionRecord,
1186 error: &anyhow::Error,
1187 updated_at: &str,
1188) {
1189 let destroying = record.state == SessionState::Destroying;
1190 if !destroying {
1191 record.state = SessionState::Closing;
1192 }
1193 record.updated_at = updated_at.to_owned();
1194 record.last_error = Some(if destroying {
1195 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1196 } else {
1197 format!("close is safely resumable from its verified checkpoint: {error:#}")
1198 });
1199}
1200
1201#[cfg(test)]
1202mod tests;