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
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 true,
87 )
88 .await?
89 {
90 self.cleanup_stopped_target(session_id, executor)?;
91 }
92 Ok(())
93 }
94
95 pub async fn suspend_session_managed_controlled(
96 &mut self,
97 session_id: &str,
98 executor: &(impl CommandExecutor + Sync),
99 manager: &SessionManagerControl,
100 acknowledge_unpublished_work: bool,
101 ) -> Result<bool> {
102 self.suspend_session_controlled_with_manager(
103 session_id,
104 executor,
105 Some(manager),
106 None,
107 SourceTargetDisposition::Destroy,
108 acknowledge_unpublished_work,
109 )
110 .await
111 }
112
113 pub(super) async fn suspend_session_for_move(
114 &mut self,
115 session_id: &str,
116 executor: &(impl CommandExecutor + Sync),
117 manager: &SessionManagerControl,
118 operation: &mut mj_core::state::MoveOperation,
119 preparation: Option<&mj_core::state::MovePreparation>,
120 disposition: SourceTargetDisposition,
121 ) -> Result<bool> {
122 self.prepare_move_source_checkpoint(session_id, executor, manager, operation)
123 .await?;
124 self.suspend_session_controlled_with_manager(
125 session_id,
126 executor,
127 Some(manager),
128 Some((operation, preparation)),
129 disposition,
130 true,
131 )
132 .await
133 }
134
135 async fn suspend_session_controlled_with_manager(
136 &mut self,
137 session_id: &str,
138 executor: &(impl CommandExecutor + Sync),
139 manager: Option<&SessionManagerControl>,
140 move_intent: Option<(
141 &mut mj_core::state::MoveOperation,
142 Option<&mj_core::state::MovePreparation>,
143 )>,
144 disposition: SourceTargetDisposition,
145 acknowledge_unpublished_work: bool,
146 ) -> Result<bool> {
147 let previous = self
148 .state
149 .sessions
150 .get(session_id)
151 .with_context(|| format!("unknown session {session_id}"))?
152 .clone();
153 let record = self.state.sessions.get_mut(session_id).unwrap();
154 apply_close_checkpoint_started(record, now());
158 self.persist_session_transition_or_restore(
159 session_id,
160 &previous,
161 "persist closing state before checkpointing the session",
162 )?;
163
164 let mut latched = match self
167 .checkpoint_session_latched(
168 session_id,
169 executor,
170 manager,
171 LatchExclusivity::HoldThroughClose,
172 CheckpointExportPolicy::ReuseUnchangedArchive,
173 )
174 .await
175 {
176 Ok(latched) => latched,
177 Err(error) => {
178 let record = self.state.sessions.get_mut(session_id).unwrap();
179 apply_close_checkpoint_failure(record, &previous, &error, now());
183 return Err(
184 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
185 );
186 }
187 };
188
189 let artifact = latched.artifact.clone();
190 let publication = match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
191 Some(ManagedCheckoutKind::Clone) => Some(super::publication::assess_clone_checkpoint(
192 &previous,
193 &artifact.metadata,
194 )),
195 Some(ManagedCheckoutKind::Worktree) => None,
196 None if previous.project_directory.is_none() => {
197 Some(super::publication::assess_network_checkpoint(
198 &previous,
199 &artifact.metadata,
200 &self.config,
201 ))
202 }
203 None => None,
204 };
205 if !acknowledge_unpublished_work
206 && publication
207 .as_ref()
208 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
209 {
210 latched.relay.cancel_abandoned_barrier().await?;
211 self.state
212 .sessions
213 .insert(session_id.to_owned(), previous.clone());
214 self.persist_session_transition_or_restore(
215 session_id,
216 &previous,
217 "restore session after unpublished-work preflight",
218 )?;
219 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
220 return Err(mj_core::refusal::Refusal::precondition(
221 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
222 ).into());
223 }
224 let record = self.state.sessions.get_mut(session_id).unwrap();
225 record.state = SessionState::Closing;
226 record.native_session_id = Some(artifact.native_session_id.clone());
227 record.checkpoint = Some(artifact.metadata.clone());
228 record.publication = publication;
229 record.updated_at = now();
230 record.last_error = None;
231 record.last_checkpoint_error = None;
232 self.persist_checkpoint_transition_or_restore(
233 session_id,
234 &previous,
235 "persist verified checkpoint and closing state before sealing the relay",
236 )?;
237 if let Some((operation, preparation)) = move_intent {
238 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
241 let record = self.state.sessions.get_mut(session_id).unwrap();
242 record.state = previous.state;
243 record.last_error = Some(format!("{error:#}"));
244 self.persist_session_transition_or_restore(
245 session_id,
246 &previous,
247 "restore source after move preflight failure",
248 )?;
249 return Err(error);
250 }
251 operation.checkpoint = Some(artifact.metadata.clone());
252 operation.updated_at = now();
253 crate::database::save_move_operation(operation)?;
254 }
255 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
256 release_projection_behind_checkpoint(session_id, &artifact.metadata);
259
260 let close_command_id = new_command_id("close")?;
261 let barrier_command_id = latched.barrier_command_id.clone();
262 let close_result = {
263 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
264 latched
265 .relay
266 .connection_mut()
267 .submit(
268 close_command_id,
269 RelayCommand::Close {
270 barrier_command_id: barrier_command_id.clone(),
271 expected: latched.cursor.clone(),
272 },
273 )
274 .await
275 };
276 if let Err(error) = close_result {
277 self.record_interrupted_close(session_id, &error)?;
278 return Err(error.context("seal verified checkpoint for close"));
279 }
280 let close_result = {
281 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
282 latched
283 .relay
284 .connection_mut()
285 .submit(
286 new_command_id("checkpoint-complete")?,
287 RelayCommand::CompleteCheckpoint { barrier_command_id },
288 )
289 .await
290 };
291 if let Err(error) = close_result {
292 self.record_interrupted_close(session_id, &error)?;
293 return Err(error.context("release verified close checkpoint"));
294 }
295 let close_result = {
296 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
297 wait_for_relay_closed(latched.relay.connection_mut()).await
298 };
299 if let Err(error) = close_result {
300 self.record_interrupted_close(session_id, &error)?;
301 return Err(error);
302 }
303 latched.relay.release();
304
305 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
306 return Ok(false);
312 }
313 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
314 Ok(deferred) => Ok(deferred),
315 Err(error) => {
316 self.record_interrupted_close(session_id, &error)?;
317 Err(error)
318 }
319 }
320 }
321
322 pub async fn recover_interrupted_close_managed(
328 &mut self,
329 session_id: &str,
330 executor: &(impl CommandExecutor + Sync),
331 manager: &SessionManagerControl,
332 ) -> Result<bool> {
333 let (state, verified) = {
334 let session = self
335 .state
336 .sessions
337 .get(session_id)
338 .with_context(|| format!("unknown session {session_id}"))?;
339 ensure!(
340 matches!(
341 session.state,
342 SessionState::Closing | SessionState::Destroying
343 ),
344 "session {session_id} has no interrupted close to recover"
345 );
346 (session.state, session.checkpoint.clone())
347 };
348 if state == SessionState::Destroying {
349 let verified = verified.context("destroying session has no verified checkpoint")?;
350 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
351 }
352 ensure!(
353 state == SessionState::Closing,
354 "session {session_id} has no relay close to recover"
355 );
356 let handle = manager
357 .wait_for_session(session_id, Duration::from_secs(5))
358 .await?;
359 let mut lease = handle.lease_connection().await?;
360 let execution = lease.connection_mut().sync().await?.operational.execution;
361 match execution {
362 RelayExecutionState::Closed => {}
363 RelayExecutionState::Closing => {
364 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
365 wait_for_relay_closed(lease.connection_mut()).await?;
366 }
367 RelayExecutionState::Idle | RelayExecutionState::Running => {
368 lease.release();
369 return self
370 .suspend_session_controlled_with_manager(
371 session_id,
372 executor,
373 Some(manager),
374 None,
375 SourceTargetDisposition::Destroy,
376 true,
377 )
378 .await;
379 }
380 }
381 lease.release();
382 let verified = verified.context("closed relay has no verified checkpoint")?;
383 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
384 }
385
386 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
393 self.fail_interrupted_lifecycle_with(
394 session_id,
395 cause,
396 crate::database::save_lifecycle_session,
397 )
398 }
399
400 fn fail_interrupted_lifecycle_with(
401 &mut self,
402 session_id: &str,
403 cause: &str,
404 persist: impl Fn(&SessionRecord) -> Result<()>,
405 ) -> Result<bool> {
406 let Some(record) = self.state.sessions.get_mut(session_id) else {
407 return Ok(false);
408 };
409 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
410 return Ok(false);
411 }
412 let previous = record.clone();
413 record.state = SessionState::Error;
414 record.updated_at = now();
415 record.last_error = Some(cause.to_owned());
416 persist_session_record_transition_or_restore(
417 &mut self.state,
418 session_id,
419 &previous,
420 "persist the failure of an interrupted lifecycle state",
421 &persist,
422 )?;
423 Ok(true)
424 }
425
426 pub fn fail_unready_session(
435 &mut self,
436 session_id: &str,
437 cause: &str,
438 observed_updated_at: &str,
439 ) -> Result<bool> {
440 self.fail_unready_session_with(
441 session_id,
442 cause,
443 observed_updated_at,
444 crate::database::save_lifecycle_session,
445 )
446 }
447
448 fn fail_unready_session_with(
449 &mut self,
450 session_id: &str,
451 cause: &str,
452 observed_updated_at: &str,
453 persist: impl Fn(&SessionRecord) -> Result<()>,
454 ) -> Result<bool> {
455 let Some(record) = self.state.sessions.get_mut(session_id) else {
456 return Ok(false);
457 };
458 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
459 return Ok(false);
460 }
461 let previous = record.clone();
462 record.state = SessionState::Error;
463 record.updated_at = now();
464 record.last_error = Some(cause.to_owned());
465 persist_session_record_transition_or_restore(
466 &mut self.state,
467 session_id,
468 &previous,
469 "persist a session whose harness never became usable",
470 &persist,
471 )?;
472 Ok(true)
473 }
474
475 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
479 let Some(record) = self.state.sessions.get(session_id) else {
480 return Ok(false);
481 };
482 let previous = record.clone();
485 let record = self.state.sessions.get_mut(session_id).unwrap();
486 record.last_error = Some(cause.to_owned());
487 record.updated_at = now();
488 persist_session_record_transition_or_restore(
489 &mut self.state,
490 session_id,
491 &previous,
492 "persist the reason a close did not finish",
493 &crate::database::save_lifecycle_session,
494 )?;
495 Ok(true)
496 }
497
498 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
503 let Some(record) = self.state.sessions.get(session_id) else {
504 return Ok(false);
505 };
506 if record.public_error().is_none() {
507 return Ok(false);
508 }
509 let previous = record.clone();
510 let record = self.state.sessions.get_mut(session_id).unwrap();
511 record.last_error = None;
512 record.updated_at = now();
513 persist_session_record_transition_or_restore(
514 &mut self.state,
515 session_id,
516 &previous,
517 "clear the reason a close did not finish",
518 &crate::database::save_lifecycle_session,
519 )?;
520 Ok(true)
521 }
522
523 pub fn suspend_session_without_checkpoint(
536 &mut self,
537 session_id: &str,
538 executor: &impl CommandExecutor,
539 ) -> Result<bool> {
540 self.suspend_session_without_checkpoint_with(
541 session_id,
542 executor,
543 crate::database::save_lifecycle_session,
544 )
545 }
546
547 fn suspend_session_without_checkpoint_with(
548 &mut self,
549 session_id: &str,
550 executor: &impl CommandExecutor,
551 persist: impl Fn(&SessionRecord) -> Result<()>,
552 ) -> Result<bool> {
553 let session = self
554 .state
555 .sessions
556 .get(session_id)
557 .with_context(|| format!("unknown session {session_id}"))?
558 .clone();
559 ensure!(
560 has_nothing_to_checkpoint(&session),
561 "session {session_id} has a workspace to checkpoint; suspend it instead"
562 );
563 self.stop_target_and_settle(session_id, &session, executor, &persist)
564 }
565
566 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
567 let record = self.state.sessions.get_mut(session_id).unwrap();
568 apply_interrupted_close_error(record, error, &now());
569 self.persist_session_state(session_id)
570 }
571
572 fn destroy_after_verified_checkpoint(
575 &mut self,
576 session_id: &str,
577 verified: &CheckpointMetadata,
578 executor: &impl CommandExecutor,
579 ) -> Result<bool> {
580 self.destroy_after_verified_checkpoint_with(
581 session_id,
582 verified,
583 executor,
584 crate::database::save_lifecycle_session,
585 )
586 }
587
588 fn destroy_after_verified_checkpoint_with(
589 &mut self,
590 session_id: &str,
591 verified: &CheckpointMetadata,
592 executor: &impl CommandExecutor,
593 persist: impl Fn(&SessionRecord) -> Result<()>,
594 ) -> Result<bool> {
595 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
596 let _target_guard = target_mutex.lock().map_err(|_| {
597 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
598 })?;
599 let session = self
600 .state
601 .sessions
602 .get(session_id)
603 .with_context(|| format!("unknown session {session_id}"))?
604 .clone();
605 ensure!(
606 matches!(
607 session.state,
608 SessionState::Closing | SessionState::Destroying
609 ),
610 "refusing to destroy session {session_id}: it is not closing or destroying"
611 );
612 ensure!(
613 session.checkpoint.as_ref() == Some(verified),
614 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
615 );
616 if session.state == SessionState::Closing {
617 let record = self.state.sessions.get_mut(session_id).unwrap();
618 record.state = SessionState::Destroying;
619 record.updated_at = now();
620 record.last_error = None;
621 persist_session_record_transition_or_restore(
622 &mut self.state,
623 session_id,
624 &session,
625 "persist destroying state before target cleanup",
626 &persist,
627 )?;
628 }
629
630 let destroying = self
631 .state
632 .sessions
633 .get(session_id)
634 .expect("destroying session disappeared")
635 .clone();
636 {
637 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
638 verify_installed_checkpoint_gate(session_id, verified)?;
639 }
640 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
645 tracing::warn!(
646 session_id,
647 error = format!("{error:#}"),
648 "could not record that the second-opinion conversation ends with this target"
649 );
650 }
651 let locator = destroying
652 .target
653 .as_ref()
654 .context("session has no target")?;
655 let backend = backend_locator(locator, &destroying, &self.config)?;
656 let deferred = if self.state.subagents.contains_key(session_id) {
657 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
658 false
659 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
660 plan.execute(executor)?;
661 true
662 } else {
663 execute_target_cleanup(&backend, session_id, executor)?;
664 false
665 };
666 if let Some(worktree) = &destroying.managed_worktree {
667 retire_managed_worktree(executor, worktree)
668 .context("retire managed raw-session worktree after verified close")?;
669 }
670 let record = self.state.sessions.get_mut(session_id).unwrap();
671 record.state = SessionState::Stopped;
672 if !deferred {
673 record.target = None;
674 }
675 record.updated_at = now();
676 record.last_error = None;
677 persist_session_record_transition_or_restore(
678 &mut self.state,
679 session_id,
680 &destroying,
681 "persist stopped state after target cleanup",
682 &persist,
683 )?;
684 Ok(deferred)
685 }
686
687 pub fn cleanup_stopped_target(
691 &mut self,
692 session_id: &str,
693 executor: &impl CommandExecutor,
694 ) -> Result<()> {
695 self.cleanup_stopped_target_with(
696 session_id,
697 executor,
698 crate::database::save_lifecycle_session,
699 )
700 }
701
702 fn cleanup_stopped_target_with(
703 &mut self,
704 session_id: &str,
705 executor: &impl CommandExecutor,
706 persist: impl Fn(&SessionRecord) -> Result<()>,
707 ) -> Result<()> {
708 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
709 let _target_guard = target_mutex.lock().map_err(|_| {
710 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
711 })?;
712 let previous = self
713 .state
714 .sessions
715 .get(session_id)
716 .with_context(|| format!("unknown session {session_id}"))?
717 .clone();
718 ensure!(
719 previous.state == SessionState::Stopped,
720 "refusing deferred cleanup for active session {session_id}"
721 );
722 let Some(locator) = previous.target.as_ref() else {
723 return Ok(());
724 };
725 let backend = backend_locator(locator, &previous, &self.config)?;
726 ensure!(
727 targets::quiesce_plan(&backend, session_id)?.is_some(),
728 "session {session_id} retained a non-Podman target after stopping"
729 );
730 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
731 let record = self.state.sessions.get_mut(session_id).unwrap();
732 record.updated_at = now();
733 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
734 let persisted = persist_session_record_transition_or_restore(
735 &mut self.state,
736 session_id,
737 &previous,
738 "persist deferred target cleanup failure",
739 &persist,
740 );
741 return match persisted {
742 Ok(()) => Err(error),
743 Err(persist_error) => Err(error.context(format!(
744 "also failed to persist deferred target cleanup failure: {persist_error:#}"
745 ))),
746 };
747 }
748 let record = self.state.sessions.get_mut(session_id).unwrap();
749 record.target = None;
750 record.updated_at = now();
751 record.last_error = None;
752 persist_session_record_transition_or_restore(
753 &mut self.state,
754 session_id,
755 &previous,
756 "persist completion of deferred Podman target cleanup",
757 &persist,
758 )
759 }
760
761 pub fn force_stop(
764 &mut self,
765 session_id: &str,
766 executor: &impl CommandExecutor,
767 ) -> Result<bool> {
768 self.force_stop_with(
769 session_id,
770 executor,
771 crate::database::save_lifecycle_session,
772 )
773 }
774
775 fn force_stop_with(
776 &mut self,
777 session_id: &str,
778 executor: &impl CommandExecutor,
779 persist: impl Fn(&SessionRecord) -> Result<()>,
780 ) -> Result<bool> {
781 let session = self
782 .state
783 .sessions
784 .get(session_id)
785 .with_context(|| format!("unknown session {session_id}"))?
786 .clone();
787 ensure!(
788 session.state.is_active(),
789 "session {session_id} is already inactive"
790 );
791 let checkpoint = session
792 .checkpoint
793 .as_ref()
794 .context("force stop requires an existing recovery archive")?;
795 {
798 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
799 verify_installed_checkpoint_gate(session_id, checkpoint)
800 .context("verify the recovery archive before force stopping")?;
801 }
802 self.stop_target_and_settle(session_id, &session, executor, &persist)
803 }
804
805 fn stop_target_and_settle(
810 &mut self,
811 session_id: &str,
812 session: &SessionRecord,
813 executor: &impl CommandExecutor,
814 persist: &impl Fn(&SessionRecord) -> Result<()>,
815 ) -> Result<bool> {
816 let mut deferred = false;
817 if let Some(locator) = &session.target {
818 let backend = backend_locator(locator, session, &self.config)?;
819 if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
820 plan.execute(executor)?;
821 deferred = true;
822 } else {
823 execute_target_cleanup(&backend, session_id, executor)?;
824 }
825 }
826 if let Some(worktree) = &session.managed_worktree {
827 retire_managed_worktree(executor, worktree)
828 .context("retire managed raw-session worktree after stopping the target")?;
829 }
830 let record = self.state.sessions.get_mut(session_id).unwrap();
831 record.state = SessionState::Stopped;
832 if !deferred {
833 record.target = None;
834 }
835 record.updated_at = now();
836 record.last_error = None;
837 record.last_checkpoint_error = None;
838 persist_session_record_transition_or_restore(
839 &mut self.state,
840 session_id,
841 session,
842 "persist stopped state after tearing down the current target",
843 persist,
844 )?;
845 Ok(deferred)
846 }
847
848 pub fn destroy_session_controlled(
852 &mut self,
853 session_id: &str,
854 executor: &impl CommandExecutor,
855 ) -> Result<()> {
856 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
857 }
858
859 pub fn destroy_session_controlled_with(
867 &mut self,
868 session_id: &str,
869 executor: &impl CommandExecutor,
870 branch: BranchDisposition,
871 ) -> Result<()> {
872 self.destroy_session_controlled_with_checkout(
873 session_id,
874 executor,
875 branch,
876 CheckoutDisposition::Remove,
877 )
878 .map(|_| ())
879 }
880
881 pub(crate) fn destroy_session_controlled_with_checkout(
887 &mut self,
888 session_id: &str,
889 executor: &impl CommandExecutor,
890 branch: BranchDisposition,
891 checkout: CheckoutDisposition,
892 ) -> Result<Option<PathBuf>> {
893 let session = self
894 .state
895 .sessions
896 .get(session_id)
897 .with_context(|| format!("unknown session {session_id}"))?
898 .clone();
899 if session.state.is_active() {
900 bail!("refusing to destroy active session {session_id}");
901 }
902 let mut retained_checkout = None;
903 if let Some(worktree) = &session.managed_worktree {
904 let keep = checkout == CheckoutDisposition::KeepWhenDirty
905 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
906 || managed_worktree_checkout_is_dirty(executor, worktree).context(
907 "check the managed raw-session worktree for uncommitted changes",
908 )?);
909 if keep {
910 retained_checkout = Some(worktree.worktree_root.clone());
914 } else {
915 cleanup_managed_worktree(executor, worktree, branch)
916 .context("remove managed raw-session worktree")?;
917 }
918 }
919 if let Some(checkpoint) = &session.checkpoint
920 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
921 && error.kind() != std::io::ErrorKind::NotFound
922 {
923 return Err(error).with_context(|| {
924 format!(
925 "remove session recovery archive {}",
926 checkpoint.archive_path.display()
927 )
928 });
929 }
930 mj_core::attachment::AttachmentStore::controller(session_id)?
931 .remove_session_data()
932 .context("remove session image attachments")?;
933 crate::database::delete_session(session_id)
934 .context("destroy stopped session in database")?;
935 self.state.subagents.remove(session_id);
936 self.state.destroy_stopped_session(session_id)?;
937 Ok(retained_checkout)
938 }
939
940 pub fn force_destroy_session(
952 &mut self,
953 session_id: &str,
954 executor: &impl CommandExecutor,
955 branch: BranchDisposition,
956 ) -> Result<()> {
957 self.force_destroy_session_with(
958 session_id,
959 executor,
960 branch,
961 crate::database::delete_session,
962 )
963 }
964
965 fn force_destroy_session_with(
966 &mut self,
967 session_id: &str,
968 executor: &impl CommandExecutor,
969 branch: BranchDisposition,
970 delete: impl Fn(&str) -> Result<()>,
971 ) -> Result<()> {
972 let session = self
973 .state
974 .sessions
975 .get(session_id)
976 .with_context(|| format!("unknown session {session_id}"))?
977 .clone();
978 if let Some(locator) = &session.target {
982 let backend = backend_locator(locator, &session, &self.config)?;
983 if self.state.subagents.contains_key(session_id) {
984 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
985 } else {
986 execute_target_cleanup(&backend, session_id, executor)?;
987 }
988 }
989 if let Some(worktree) = &session.managed_worktree {
990 cleanup_managed_worktree(executor, worktree, branch)
991 .context("remove managed raw-session worktree")?;
992 }
993 if let Some(checkpoint) = &session.checkpoint
994 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
995 && error.kind() != std::io::ErrorKind::NotFound
996 {
997 return Err(error).with_context(|| {
998 format!(
999 "remove session recovery archive {}",
1000 checkpoint.archive_path.display()
1001 )
1002 });
1003 }
1004 mj_core::attachment::AttachmentStore::controller(session_id)?
1005 .remove_session_data()
1006 .context("remove session image attachments")?;
1007 delete(session_id).context("force destroy session in database")?;
1008 self.state.subagents.remove(session_id);
1009 self.state.destroy_session_force(session_id)?;
1010 Ok(())
1011 }
1012}
1013
1014pub fn has_nothing_to_checkpoint(session: &SessionRecord) -> bool {
1023 match session.state {
1024 SessionState::Provisioning => true,
1025 SessionState::Closing
1026 | SessionState::Destroying
1027 | SessionState::Error
1028 | SessionState::Lost => session.target.is_none(),
1029 SessionState::Running
1030 | SessionState::Disconnected
1031 | SessionState::Checkpointing
1032 | SessionState::Stopped
1033 | SessionState::DestroyedWithDataLoss => false,
1034 }
1035}
1036
1037fn execute_target_cleanup(
1038 backend: &targets::TargetLocator,
1039 session_id: &str,
1040 executor: &impl CommandExecutor,
1041) -> Result<()> {
1042 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1043 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1044 Ok(true) => {
1045 tracing::warn!(
1046 session_id,
1047 error = format!("{cleanup_error:#}"),
1048 "target cleanup command failed, but the target was confirmed absent"
1049 );
1050 }
1051 Ok(false) => {
1052 tracing::error!(
1053 session_id,
1054 error = format!("{cleanup_error:#}"),
1055 "target cleanup failed and the target is still present"
1056 );
1057 return Err(cleanup_error);
1058 }
1059 Err(probe_error) => {
1060 tracing::error!(
1061 session_id,
1062 cleanup_error = format!("{cleanup_error:#}"),
1063 probe_error = format!("{probe_error:#}"),
1064 "target cleanup failed and exact absence could not be confirmed"
1065 );
1066 return Err(cleanup_error.context(format!(
1067 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1068 )));
1069 }
1070 }
1071 }
1072 Ok(())
1073}
1074
1075fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1076 record.state = SessionState::Closing;
1077 record.updated_at = updated_at;
1078 record.last_checkpoint_error = None;
1079}
1080
1081fn apply_close_checkpoint_failure(
1088 record: &mut SessionRecord,
1089 previous: &SessionRecord,
1090 error: &anyhow::Error,
1091 updated_at: String,
1092) {
1093 if WorkerRestartLeftNoWorker::marks(error) {
1094 record.state = SessionState::Error;
1095 record.last_error = Some(format!(
1096 "close failed and left the session without a live worker; retry the close, \
1097 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1098 ));
1099 } else {
1100 record.state = if previous.state == SessionState::Closing {
1101 SessionState::Running
1102 } else {
1103 previous.state
1104 };
1105 }
1106 record.last_checkpoint_error = Some(format!("{error:#}"));
1107 record.updated_at = updated_at;
1108}
1109
1110fn apply_interrupted_close_error(
1111 record: &mut SessionRecord,
1112 error: &anyhow::Error,
1113 updated_at: &str,
1114) {
1115 let destroying = record.state == SessionState::Destroying;
1116 if !destroying {
1117 record.state = SessionState::Closing;
1118 }
1119 record.updated_at = updated_at.to_owned();
1120 record.last_error = Some(if destroying {
1121 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1122 } else {
1123 format!("close is safely resumable from its verified checkpoint: {error:#}")
1124 });
1125}
1126
1127#[cfg(test)]
1128mod tests;