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 None,
97 )
98 .await?
99 {
100 self.cleanup_stopped_target(session_id, executor)?;
101 }
102 Ok(())
103 }
104
105 pub async fn suspend_session_managed_controlled(
106 &mut self,
107 session_id: &str,
108 executor: &(impl CommandExecutor + Sync),
109 manager: &SessionManagerControl,
110 acknowledge_unpublished_work: bool,
111 before_close: Option<BeforeClose>,
112 ) -> Result<bool> {
113 self.suspend_session_controlled_with_manager(
114 session_id,
115 executor,
116 Some(manager),
117 None,
118 SourceTargetDisposition::Destroy,
119 acknowledge_unpublished_work,
120 before_close,
121 None,
122 )
123 .await
124 }
125
126 #[allow(clippy::too_many_arguments)]
127 pub(super) async fn suspend_session_for_move(
128 &mut self,
129 session_id: &str,
130 executor: &(impl CommandExecutor + Sync),
131 manager: &SessionManagerControl,
132 operation: &mut mj_core::state::MoveOperation,
133 preparation: Option<&mj_core::state::MovePreparation>,
134 disposition: SourceTargetDisposition,
135 mut source_relay: super::move_session::MoveSourceRelay,
136 ) -> Result<bool> {
137 self.prepare_move_source_checkpoint(
138 session_id,
139 executor,
140 manager,
141 operation,
142 &mut source_relay,
143 )
144 .await?;
145 if disposition == SourceTargetDisposition::RetainForInPlaceSwap
146 || operation.workspace_transfer.is_some()
147 {
148 self.seal_move_handoff(
149 session_id,
150 executor,
151 manager,
152 operation,
153 preparation,
154 source_relay.take(),
155 )
156 .await?;
157 return Ok(false);
158 }
159 self.suspend_session_controlled_with_manager(
160 session_id,
161 executor,
162 Some(manager),
163 Some((operation, preparation)),
164 disposition,
165 true,
166 None,
167 source_relay.take(),
168 )
169 .await
170 }
171
172 #[allow(clippy::too_many_arguments)]
173 async fn suspend_session_controlled_with_manager(
174 &mut self,
175 session_id: &str,
176 executor: &(impl CommandExecutor + Sync),
177 manager: Option<&SessionManagerControl>,
178 move_intent: Option<(
179 &mut mj_core::state::MoveOperation,
180 Option<&mj_core::state::MovePreparation>,
181 )>,
182 disposition: SourceTargetDisposition,
183 acknowledge_unpublished_work: bool,
184 before_close: Option<BeforeClose>,
185 held_relay: Option<super::checkpoint::ControllerRelayLease>,
186 ) -> Result<bool> {
187 let previous = self
188 .state
189 .sessions
190 .get(session_id)
191 .with_context(|| format!("unknown session {session_id}"))?
192 .clone();
193 let record = self.state.sessions.get_mut(session_id).unwrap();
194 apply_close_checkpoint_started(record, now());
198 self.persist_session_transition_or_restore(
199 session_id,
200 &previous,
201 "persist closing state before checkpointing the session",
202 )?;
203
204 let mut latched = match self
207 .checkpoint_session_latched_for_operation(
208 session_id,
209 executor,
210 manager,
211 LatchExclusivity::HoldThroughClose,
212 CheckpointExportPolicy::ReuseUnchangedArchive,
213 None,
214 held_relay,
215 )
216 .await
217 {
218 Ok(latched) => latched,
219 Err(error) => {
220 let record = self.state.sessions.get_mut(session_id).unwrap();
221 apply_close_checkpoint_failure(record, &previous, &error, now());
225 return Err(
226 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
227 );
228 }
229 };
230
231 let artifact = latched.artifact.clone();
232 let publication = if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
235 None
236 } else {
237 match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
238 Some(ManagedCheckoutKind::Clone) => Some(
239 super::publication::assess_clone_checkpoint(&previous, &artifact.metadata),
240 ),
241 Some(ManagedCheckoutKind::Worktree) => None,
242 None if previous.project_directory.is_none() => {
243 Some(super::publication::assess_network_checkpoint(
244 &previous,
245 &artifact.metadata,
246 &self.config,
247 ))
248 }
249 None => None,
250 }
251 };
252 if !acknowledge_unpublished_work
253 && publication
254 .as_ref()
255 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
256 {
257 latched.relay.cancel_abandoned_barrier().await?;
258 let mut restored = previous.clone();
262 restored.state = state_after_unsealed_close(&previous);
263 restored.updated_at = now();
264 self.state.sessions.insert(session_id.to_owned(), restored);
265 self.persist_session_transition_or_restore(
266 session_id,
267 &previous,
268 "restore session after unpublished-work preflight",
269 )?;
270 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
271 return Err(mj_core::refusal::Refusal::precondition(
272 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
273 ).into());
274 }
275 let record = self.state.sessions.get_mut(session_id).unwrap();
276 record.state = SessionState::Closing;
277 record.native_session_id = Some(artifact.native_session_id.clone());
278 record.checkpoint = Some(artifact.metadata.clone());
279 record.publication = publication;
280 record.updated_at = now();
281 record.last_error = None;
282 record.last_checkpoint_error = None;
283 self.persist_checkpoint_transition_or_restore(
284 session_id,
285 &previous,
286 "persist verified checkpoint and closing state before sealing the relay",
287 )?;
288 if let Some((operation, preparation)) = move_intent {
289 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
292 let record = self.state.sessions.get_mut(session_id).unwrap();
293 record.state = previous.state;
294 record.last_error = Some(format!("{error:#}"));
295 self.persist_session_transition_or_restore(
296 session_id,
297 &previous,
298 "restore source after move preflight failure",
299 )?;
300 return Err(error);
301 }
302 operation.checkpoint = Some(artifact.metadata.clone());
303 operation.updated_at = now();
304 crate::database::save_move_operation(operation)?;
305 }
306 if let Some(before_close) = before_close
307 && let Err(error) = before_close.await
308 {
309 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
312 tracing::warn!(
313 session_id,
314 error = format!("{cancel:#}"),
315 "could not release the checkpoint barrier of a close that did not proceed"
316 );
317 }
318 let record = self.state.sessions.get_mut(session_id).unwrap();
319 record.state = state_after_unsealed_close(&previous);
320 record.last_error = Some(format!("{error:#}"));
321 record.updated_at = now();
322 self.persist_session_transition_or_restore(
323 session_id,
324 &previous,
325 "restore a session whose close did not proceed past its checkpoint",
326 )?;
327 return Err(error);
328 }
329 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
330 release_projection_behind_checkpoint(session_id, &artifact.metadata);
333
334 let close_command_id = new_command_id("close")?;
335 let barrier_command_id = latched.barrier_command_id.clone();
336 let close_result = {
337 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
338 latched
339 .relay
340 .connection_mut()
341 .submit(
342 close_command_id,
343 RelayCommand::Close {
344 barrier_command_id: barrier_command_id.clone(),
345 expected: latched.cursor.clone(),
346 },
347 )
348 .await
349 };
350 if let Err(error) = close_result {
351 let error = error.context("seal verified checkpoint for close");
352 if !close_was_refused(&error) {
353 self.record_interrupted_close(session_id, &error)?;
356 return Err(error);
357 }
358 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
363 tracing::warn!(
364 session_id,
365 error = format!("{cancel:#}"),
366 "could not release the checkpoint barrier of a refused close"
367 );
368 }
369 let record = self.state.sessions.get_mut(session_id).unwrap();
370 record.state = state_after_unsealed_close(&previous);
371 record.last_error = Some(format!("{error:#}"));
372 record.updated_at = now();
373 self.persist_session_transition_or_restore(
374 session_id,
375 &previous,
376 "restore a session whose worker refused its close",
377 )?;
378 return Err(error);
379 }
380 let close_result = {
381 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
382 latched
383 .relay
384 .connection_mut()
385 .submit(
386 new_command_id("checkpoint-complete")?,
387 RelayCommand::CompleteCheckpoint { barrier_command_id },
388 )
389 .await
390 };
391 if let Err(error) = close_result {
392 self.record_interrupted_close(session_id, &error)?;
393 return Err(error.context("release verified close checkpoint"));
394 }
395 let close_result = {
396 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
397 wait_for_relay_closed(latched.relay.connection_mut()).await
398 };
399 if let Err(error) = close_result {
400 self.record_interrupted_close(session_id, &error)?;
401 return Err(error);
402 }
403 latched.relay.release();
404
405 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
406 return Ok(false);
412 }
413 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
414 Ok(deferred) => Ok(deferred),
415 Err(error) => {
416 self.record_interrupted_close(session_id, &error)?;
417 Err(error)
418 }
419 }
420 }
421
422 pub async fn recover_interrupted_close_managed(
434 &mut self,
435 session_id: &str,
436 executor: &(impl CommandExecutor + Sync),
437 manager: &SessionManagerControl,
438 acknowledge_unpublished_work: bool,
439 before_close: Option<BeforeClose>,
440 ) -> Result<bool> {
441 let (state, verified) = {
442 let session = self
443 .state
444 .sessions
445 .get(session_id)
446 .with_context(|| format!("unknown session {session_id}"))?;
447 ensure!(
448 matches!(
449 session.state,
450 SessionState::Closing | SessionState::Destroying
451 ),
452 "session {session_id} has no interrupted close to recover"
453 );
454 (session.state, session.checkpoint.clone())
455 };
456 if state == SessionState::Destroying {
457 let verified = verified.context("destroying session has no verified checkpoint")?;
458 if let Some(before_close) = before_close {
459 before_close.await?;
460 }
461 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
462 }
463 ensure!(
464 state == SessionState::Closing,
465 "session {session_id} has no relay close to recover"
466 );
467 let handle = manager
468 .wait_for_session(session_id, Duration::from_secs(5))
469 .await?;
470 let mut lease = handle.lease_connection().await?;
471 let execution = lease.connection_mut().sync().await?.operational.execution;
472 match execution {
473 RelayExecutionState::Closed => {}
474 RelayExecutionState::Closing => {
475 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
476 wait_for_relay_closed(lease.connection_mut()).await?;
477 }
478 RelayExecutionState::Idle | RelayExecutionState::Running => {
479 lease.release();
480 return self
481 .suspend_session_controlled_with_manager(
482 session_id,
483 executor,
484 Some(manager),
485 None,
486 SourceTargetDisposition::Destroy,
487 acknowledge_unpublished_work,
488 before_close,
489 None,
490 )
491 .await;
492 }
493 }
494 lease.release();
495 let verified = verified.context("closed relay has no verified checkpoint")?;
496 if let Some(before_close) = before_close {
498 before_close.await?;
499 }
500 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
501 }
502
503 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
510 self.fail_interrupted_lifecycle_with(
511 session_id,
512 cause,
513 crate::database::save_lifecycle_session,
514 )
515 }
516
517 fn fail_interrupted_lifecycle_with(
518 &mut self,
519 session_id: &str,
520 cause: &str,
521 persist: impl Fn(&SessionRecord) -> Result<()>,
522 ) -> Result<bool> {
523 let Some(record) = self.state.sessions.get_mut(session_id) else {
524 return Ok(false);
525 };
526 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
527 return Ok(false);
528 }
529 let previous = record.clone();
530 record.state = SessionState::Error;
531 record.updated_at = now();
532 record.last_error = Some(cause.to_owned());
533 persist_session_record_transition_or_restore(
534 &mut self.state,
535 session_id,
536 &previous,
537 "persist the failure of an interrupted lifecycle state",
538 &persist,
539 )?;
540 Ok(true)
541 }
542
543 pub fn fail_unready_session(
552 &mut self,
553 session_id: &str,
554 cause: &str,
555 observed_updated_at: &str,
556 ) -> Result<bool> {
557 self.fail_unready_session_with(
558 session_id,
559 cause,
560 observed_updated_at,
561 crate::database::save_lifecycle_session,
562 )
563 }
564
565 fn fail_unready_session_with(
566 &mut self,
567 session_id: &str,
568 cause: &str,
569 observed_updated_at: &str,
570 persist: impl Fn(&SessionRecord) -> Result<()>,
571 ) -> Result<bool> {
572 let Some(record) = self.state.sessions.get_mut(session_id) else {
573 return Ok(false);
574 };
575 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
576 return Ok(false);
577 }
578 let previous = record.clone();
579 record.state = SessionState::Error;
580 record.updated_at = now();
581 record.last_error = Some(cause.to_owned());
582 persist_session_record_transition_or_restore(
583 &mut self.state,
584 session_id,
585 &previous,
586 "persist a session whose harness never became usable",
587 &persist,
588 )?;
589 Ok(true)
590 }
591
592 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
596 let Some(record) = self.state.sessions.get(session_id) else {
597 return Ok(false);
598 };
599 let previous = record.clone();
602 let record = self.state.sessions.get_mut(session_id).unwrap();
603 record.last_error = Some(cause.to_owned());
604 record.updated_at = now();
605 persist_session_record_transition_or_restore(
606 &mut self.state,
607 session_id,
608 &previous,
609 "persist the reason a close did not finish",
610 &crate::database::save_lifecycle_session,
611 )?;
612 Ok(true)
613 }
614
615 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
620 let Some(record) = self.state.sessions.get(session_id) else {
621 return Ok(false);
622 };
623 if record.public_error().is_none() {
624 return Ok(false);
625 }
626 let previous = record.clone();
627 let record = self.state.sessions.get_mut(session_id).unwrap();
628 record.last_error = None;
629 record.updated_at = now();
630 persist_session_record_transition_or_restore(
631 &mut self.state,
632 session_id,
633 &previous,
634 "clear the reason a close did not finish",
635 &crate::database::save_lifecycle_session,
636 )?;
637 Ok(true)
638 }
639
640 pub fn suspend_session_without_checkpoint(
653 &mut self,
654 session_id: &str,
655 executor: &impl CommandExecutor,
656 ) -> Result<bool> {
657 self.suspend_session_without_checkpoint_with(
658 session_id,
659 executor,
660 crate::database::save_lifecycle_session,
661 )
662 }
663
664 fn suspend_session_without_checkpoint_with(
665 &mut self,
666 session_id: &str,
667 executor: &impl CommandExecutor,
668 persist: impl Fn(&SessionRecord) -> Result<()>,
669 ) -> Result<bool> {
670 let session = self
671 .state
672 .sessions
673 .get(session_id)
674 .with_context(|| format!("unknown session {session_id}"))?
675 .clone();
676 ensure!(
677 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
678 "session {session_id} has a workspace to checkpoint; suspend it instead"
679 );
680 self.stop_target_and_settle(session_id, &session, executor, &persist)
681 }
682
683 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
684 let record = self.state.sessions.get_mut(session_id).unwrap();
685 apply_interrupted_close_error(record, error, &now());
686 self.persist_session_state(session_id)
687 }
688
689 fn destroy_after_verified_checkpoint(
692 &mut self,
693 session_id: &str,
694 verified: &CheckpointMetadata,
695 executor: &impl CommandExecutor,
696 ) -> Result<bool> {
697 self.destroy_after_verified_checkpoint_with(
698 session_id,
699 verified,
700 executor,
701 crate::database::save_lifecycle_session,
702 )
703 }
704
705 fn destroy_after_verified_checkpoint_with(
706 &mut self,
707 session_id: &str,
708 verified: &CheckpointMetadata,
709 executor: &impl CommandExecutor,
710 persist: impl Fn(&SessionRecord) -> Result<()>,
711 ) -> Result<bool> {
712 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
713 let _target_guard = target_mutex.lock().map_err(|_| {
714 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
715 })?;
716 let session = self
717 .state
718 .sessions
719 .get(session_id)
720 .with_context(|| format!("unknown session {session_id}"))?
721 .clone();
722 ensure!(
723 matches!(
724 session.state,
725 SessionState::Closing | SessionState::Destroying
726 ),
727 "refusing to destroy session {session_id}: it is not closing or destroying"
728 );
729 ensure!(
730 session.checkpoint.as_ref() == Some(verified),
731 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
732 );
733 if session.state == SessionState::Closing {
734 let record = self.state.sessions.get_mut(session_id).unwrap();
735 record.state = SessionState::Destroying;
736 record.updated_at = now();
737 record.last_error = None;
738 persist_session_record_transition_or_restore(
739 &mut self.state,
740 session_id,
741 &session,
742 "persist destroying state before target cleanup",
743 &persist,
744 )?;
745 }
746
747 let destroying = self
748 .state
749 .sessions
750 .get(session_id)
751 .expect("destroying session disappeared")
752 .clone();
753 {
754 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
755 verify_installed_checkpoint_gate(session_id, verified)?;
756 }
757 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
762 tracing::warn!(
763 session_id,
764 error = format!("{error:#}"),
765 "could not record that the second-opinion conversation ends with this target"
766 );
767 }
768 let locator = destroying
769 .target
770 .as_ref()
771 .context("session has no target")?;
772 let backend = backend_locator(locator, &destroying, &self.config)?;
773 let deferred = if self.state.subagents.contains_key(session_id) {
774 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
775 false
776 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
777 plan.execute(executor)?;
778 true
779 } else {
780 execute_target_cleanup(&backend, session_id, executor)?;
781 false
782 };
783 if let Some(worktree) = &destroying.managed_worktree {
784 retire_managed_worktree(executor, worktree)
785 .context("retire managed raw-session worktree after verified close")?;
786 }
787 let record = self.state.sessions.get_mut(session_id).unwrap();
788 record.state = SessionState::Stopped;
789 if !deferred {
790 record.target = None;
791 }
792 record.updated_at = now();
793 record.last_error = None;
794 persist_session_record_transition_or_restore(
795 &mut self.state,
796 session_id,
797 &destroying,
798 "persist stopped state after target cleanup",
799 &persist,
800 )?;
801 Ok(deferred)
802 }
803
804 pub fn cleanup_stopped_target(
808 &mut self,
809 session_id: &str,
810 executor: &impl CommandExecutor,
811 ) -> Result<()> {
812 self.cleanup_stopped_target_with(
813 session_id,
814 executor,
815 crate::database::save_lifecycle_session,
816 )
817 }
818
819 fn cleanup_stopped_target_with(
820 &mut self,
821 session_id: &str,
822 executor: &impl CommandExecutor,
823 persist: impl Fn(&SessionRecord) -> Result<()>,
824 ) -> Result<()> {
825 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
826 let _target_guard = target_mutex.lock().map_err(|_| {
827 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
828 })?;
829 let previous = self
830 .state
831 .sessions
832 .get(session_id)
833 .with_context(|| format!("unknown session {session_id}"))?
834 .clone();
835 ensure!(
836 previous.state == SessionState::Stopped,
837 "refusing deferred cleanup for active session {session_id}"
838 );
839 let Some(locator) = previous.target.as_ref() else {
840 return Ok(());
841 };
842 let backend = backend_locator(locator, &previous, &self.config)?;
843 ensure!(
844 targets::quiesce_plan(&backend, session_id)?.is_some(),
845 "session {session_id} retained a non-Podman target after stopping"
846 );
847 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
848 let record = self.state.sessions.get_mut(session_id).unwrap();
849 record.updated_at = now();
850 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
851 let persisted = persist_session_record_transition_or_restore(
852 &mut self.state,
853 session_id,
854 &previous,
855 "persist deferred target cleanup failure",
856 &persist,
857 );
858 return match persisted {
859 Ok(()) => Err(error),
860 Err(persist_error) => Err(error.context(format!(
861 "also failed to persist deferred target cleanup failure: {persist_error:#}"
862 ))),
863 };
864 }
865 let record = self.state.sessions.get_mut(session_id).unwrap();
866 record.target = None;
867 record.updated_at = now();
868 record.last_error = None;
869 persist_session_record_transition_or_restore(
870 &mut self.state,
871 session_id,
872 &previous,
873 "persist completion of deferred Podman target cleanup",
874 &persist,
875 )
876 }
877
878 pub fn force_stop(
881 &mut self,
882 session_id: &str,
883 executor: &impl CommandExecutor,
884 ) -> Result<bool> {
885 self.force_stop_with(
886 session_id,
887 executor,
888 crate::database::save_lifecycle_session,
889 )
890 }
891
892 fn force_stop_with(
893 &mut self,
894 session_id: &str,
895 executor: &impl CommandExecutor,
896 persist: impl Fn(&SessionRecord) -> Result<()>,
897 ) -> Result<bool> {
898 let session = self
899 .state
900 .sessions
901 .get(session_id)
902 .with_context(|| format!("unknown session {session_id}"))?
903 .clone();
904 ensure!(
905 session.state.is_active(),
906 "session {session_id} is already inactive"
907 );
908 let checkpoint = session
909 .checkpoint
910 .as_ref()
911 .context("force stop requires an existing recovery archive")?;
912 {
915 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
916 verify_installed_checkpoint_gate(session_id, checkpoint)
917 .context("verify the recovery archive before force stopping")?;
918 }
919 self.stop_target_and_settle(session_id, &session, executor, &persist)
920 }
921
922 fn stop_target_and_settle(
927 &mut self,
928 session_id: &str,
929 session: &SessionRecord,
930 executor: &impl CommandExecutor,
931 persist: &impl Fn(&SessionRecord) -> Result<()>,
932 ) -> Result<bool> {
933 let mut deferred = false;
934 if let Some(locator) = &session.target {
935 let backend = backend_locator(locator, session, &self.config)?;
936 if self.state.subagents.contains_key(session_id) {
940 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
941 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
942 plan.execute(executor)?;
943 deferred = true;
944 } else {
945 execute_target_cleanup(&backend, session_id, executor)?;
946 }
947 }
948 if let Some(worktree) = &session.managed_worktree {
949 retire_managed_worktree(executor, worktree)
950 .context("retire managed raw-session worktree after stopping the target")?;
951 }
952 let record = self.state.sessions.get_mut(session_id).unwrap();
953 record.state = SessionState::Stopped;
954 if !deferred {
955 record.target = None;
956 }
957 record.updated_at = now();
958 record.last_error = None;
959 record.last_checkpoint_error = None;
960 persist_session_record_transition_or_restore(
961 &mut self.state,
962 session_id,
963 session,
964 "persist stopped state after tearing down the current target",
965 persist,
966 )?;
967 Ok(deferred)
968 }
969
970 pub fn destroy_session_controlled(
974 &mut self,
975 session_id: &str,
976 executor: &impl CommandExecutor,
977 ) -> Result<()> {
978 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
979 }
980
981 pub fn destroy_session_controlled_with(
989 &mut self,
990 session_id: &str,
991 executor: &impl CommandExecutor,
992 branch: BranchDisposition,
993 ) -> Result<()> {
994 self.destroy_session_controlled_with_checkout(
995 session_id,
996 executor,
997 branch,
998 CheckoutDisposition::Remove,
999 )
1000 .map(|_| ())
1001 }
1002
1003 pub(crate) fn destroy_session_controlled_with_checkout(
1009 &mut self,
1010 session_id: &str,
1011 executor: &impl CommandExecutor,
1012 branch: BranchDisposition,
1013 checkout: CheckoutDisposition,
1014 ) -> Result<Option<PathBuf>> {
1015 let session = self
1016 .state
1017 .sessions
1018 .get(session_id)
1019 .with_context(|| format!("unknown session {session_id}"))?
1020 .clone();
1021 if session.state.is_active() {
1022 bail!("refusing to destroy active session {session_id}");
1023 }
1024 let mut retained_checkout = None;
1025 if let Some(worktree) = &session.managed_worktree {
1026 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1027 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1028 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1029 "check the managed raw-session worktree for uncommitted changes",
1030 )?);
1031 if keep {
1032 retained_checkout = Some(worktree.worktree_root.clone());
1036 } else {
1037 cleanup_managed_worktree(executor, worktree, branch)
1038 .context("remove managed raw-session worktree")?;
1039 }
1040 }
1041 if let Some(checkpoint) = &session.checkpoint
1042 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1043 && error.kind() != std::io::ErrorKind::NotFound
1044 {
1045 return Err(error).with_context(|| {
1046 format!(
1047 "remove session recovery archive {}",
1048 checkpoint.archive_path.display()
1049 )
1050 });
1051 }
1052 mj_core::attachment::AttachmentStore::controller(session_id)?
1053 .remove_session_data()
1054 .context("remove session image attachments")?;
1055 crate::database::delete_session(session_id)
1056 .context("destroy stopped session in database")?;
1057 self.state.subagents.remove(session_id);
1058 self.state.destroy_stopped_session(session_id)?;
1059 Ok(retained_checkout)
1060 }
1061
1062 pub fn force_destroy_session(
1074 &mut self,
1075 session_id: &str,
1076 executor: &impl CommandExecutor,
1077 branch: BranchDisposition,
1078 ) -> Result<()> {
1079 self.force_destroy_session_with(
1080 session_id,
1081 executor,
1082 branch,
1083 crate::database::delete_session,
1084 )
1085 }
1086
1087 fn force_destroy_session_with(
1088 &mut self,
1089 session_id: &str,
1090 executor: &impl CommandExecutor,
1091 branch: BranchDisposition,
1092 delete: impl Fn(&str) -> Result<()>,
1093 ) -> Result<()> {
1094 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1095 if result.is_err() {
1096 self.settle_failed_destroy(session_id);
1097 }
1098 result
1099 }
1100
1101 fn settle_failed_destroy(&mut self, session_id: &str) {
1107 let Some(record) = self.state.sessions.get(session_id) else {
1108 return;
1109 };
1110 if !record.state.is_active() || record.state == SessionState::Error {
1111 return;
1112 }
1113 let previous = record.clone();
1114 let record = self.state.sessions.get_mut(session_id).unwrap();
1115 record.state = SessionState::Error;
1116 record.updated_at = now();
1117 if let Err(error) = persist_session_record_transition_or_restore(
1118 &mut self.state,
1119 session_id,
1120 &previous,
1121 "persist the failed destroy",
1122 &crate::database::save_lifecycle_session,
1123 ) {
1124 tracing::warn!(
1125 session_id,
1126 error = format!("{error:#}"),
1127 "could not record a failed destroy"
1128 );
1129 }
1130 }
1131
1132 fn force_destroy_session_steps(
1133 &mut self,
1134 session_id: &str,
1135 executor: &impl CommandExecutor,
1136 branch: BranchDisposition,
1137 delete: impl Fn(&str) -> Result<()>,
1138 ) -> Result<()> {
1139 let session = self
1140 .state
1141 .sessions
1142 .get(session_id)
1143 .with_context(|| format!("unknown session {session_id}"))?
1144 .clone();
1145 if let Some(locator) = &session.target {
1149 let backend = backend_locator(locator, &session, &self.config)?;
1150 if self.state.subagents.contains_key(session_id) {
1151 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
1152 } else {
1153 execute_target_cleanup(&backend, session_id, executor)?;
1154 }
1155 }
1156 if let Some(worktree) = &session.managed_worktree {
1157 cleanup_managed_worktree(executor, worktree, branch)
1158 .context("remove managed raw-session worktree")?;
1159 }
1160 if let Some(checkpoint) = &session.checkpoint
1161 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1162 && error.kind() != std::io::ErrorKind::NotFound
1163 {
1164 return Err(error).with_context(|| {
1165 format!(
1166 "remove session recovery archive {}",
1167 checkpoint.archive_path.display()
1168 )
1169 });
1170 }
1171 mj_core::attachment::AttachmentStore::controller(session_id)?
1172 .remove_session_data()
1173 .context("remove session image attachments")?;
1174 delete(session_id).context("force destroy session in database")?;
1175 self.state.subagents.remove(session_id);
1176 self.state.destroy_session_force(session_id)?;
1177 Ok(())
1178 }
1179}
1180
1181pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1196 match session.state {
1197 SessionState::Provisioning | SessionState::Parked => true,
1198 SessionState::Error if subagent => true,
1199 SessionState::Closing
1200 | SessionState::Destroying
1201 | SessionState::Error
1202 | SessionState::Lost => session.target.is_none(),
1203 SessionState::Running
1204 | SessionState::Disconnected
1205 | SessionState::Checkpointing
1206 | SessionState::Stopped
1207 | SessionState::DestroyedWithDataLoss => false,
1208 }
1209}
1210
1211fn execute_target_cleanup(
1212 backend: &targets::TargetLocator,
1213 session_id: &str,
1214 executor: &impl CommandExecutor,
1215) -> Result<()> {
1216 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1217 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1218 Ok(true) => {
1219 tracing::warn!(
1220 session_id,
1221 error = format!("{cleanup_error:#}"),
1222 "target cleanup command failed, but the target was confirmed absent"
1223 );
1224 }
1225 Ok(false) => {
1226 tracing::error!(
1227 session_id,
1228 error = format!("{cleanup_error:#}"),
1229 "target cleanup failed and the target is still present"
1230 );
1231 return Err(cleanup_error);
1232 }
1233 Err(probe_error) => {
1234 tracing::error!(
1235 session_id,
1236 cleanup_error = format!("{cleanup_error:#}"),
1237 probe_error = format!("{probe_error:#}"),
1238 "target cleanup failed and exact absence could not be confirmed"
1239 );
1240 return Err(cleanup_error.context(format!(
1241 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1242 )));
1243 }
1244 }
1245 }
1246 Ok(())
1247}
1248
1249fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1250 record.state = SessionState::Closing;
1251 record.updated_at = updated_at;
1252 record.last_checkpoint_error = None;
1253}
1254
1255fn apply_close_checkpoint_failure(
1262 record: &mut SessionRecord,
1263 previous: &SessionRecord,
1264 error: &anyhow::Error,
1265 updated_at: String,
1266) {
1267 if WorkerRestartLeftNoWorker::marks(error) {
1268 record.state = SessionState::Error;
1269 record.last_error = Some(format!(
1270 "close failed and left the session without a live worker; retry the close, \
1271 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1272 ));
1273 } else {
1274 record.state = state_after_unsealed_close(previous);
1275 }
1276 record.last_checkpoint_error = Some(format!("{error:#}"));
1277 record.updated_at = updated_at;
1278}
1279
1280fn close_was_refused(error: &anyhow::Error) -> bool {
1288 error.chain().any(|cause| {
1289 cause
1290 .downcast_ref::<crate::worker_client::RelayRejected>()
1291 .is_some_and(|rejected| !rejected.is_retryable())
1292 })
1293}
1294
1295fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1296 if previous.state == SessionState::Closing {
1297 SessionState::Running
1298 } else {
1299 previous.state
1300 }
1301}
1302
1303fn apply_interrupted_close_error(
1304 record: &mut SessionRecord,
1305 error: &anyhow::Error,
1306 updated_at: &str,
1307) {
1308 let destroying = record.state == SessionState::Destroying;
1309 if !destroying {
1310 record.state = SessionState::Closing;
1311 }
1312 record.updated_at = updated_at.to_owned();
1313 record.last_error = Some(if destroying {
1314 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1315 } else {
1316 format!("close is safely resumable from its verified checkpoint: {error:#}")
1317 });
1318}
1319
1320#[cfg(test)]
1321mod tests;