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 started = std::time::Instant::now();
468 let mut phase = "wait for session actor";
469 tracing::info!(
470 session_id,
471 ?state,
472 checkpoint_frontier = ?verified.as_ref().map(|checkpoint| checkpoint.event_frontier),
473 "recovering suspension: asking the worker whether its relay was sealed"
474 );
475 let connection = async {
476 let handle = manager
477 .wait_for_session(session_id, Duration::from_secs(5))
478 .await?;
479 phase = "lease worker connection";
480 let mut lease = handle.lease_connection().await?;
481 phase = "read worker execution state";
482 let execution = lease.connection_mut().sync().await?.operational.execution;
483 anyhow::Ok((lease, execution))
484 }
485 .await;
486 let (mut lease, execution) = match connection {
487 Ok(connected) => connected,
488 Err(error) => {
489 tracing::warn!(
490 session_id,
491 phase,
492 elapsed_ms = started.elapsed().as_millis() as u64,
493 error = format!("{error:#}"),
494 "suspension could not query its worker; retaining target and checkpoint"
495 );
496 self.diagnose_suspension_worker(session_id, phase).await;
497 return Err(error.context(format!("{phase} while recovering suspension")));
498 }
499 };
500 tracing::info!(
501 session_id,
502 ?execution,
503 elapsed_ms = started.elapsed().as_millis() as u64,
504 "suspension recovery read the worker's execution state"
505 );
506 match execution {
507 RelayExecutionState::Closed => {}
508 RelayExecutionState::Closing => {
509 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
510 wait_for_relay_closed(lease.connection_mut()).await?;
511 }
512 RelayExecutionState::Idle | RelayExecutionState::Running => {
513 lease.release();
514 return self
515 .suspend_session_controlled_with_manager(
516 session_id,
517 executor,
518 Some(manager),
519 None,
520 SourceTargetDisposition::Destroy,
521 acknowledge_unpublished_work,
522 before_close,
523 None,
524 )
525 .await;
526 }
527 }
528 lease.release();
529 let verified = verified.context("closed relay has no verified checkpoint")?;
530 if let Some(before_close) = before_close {
532 before_close.await?;
533 }
534 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
535 }
536
537 async fn diagnose_suspension_worker(&self, session_id: &str, phase: &str) {
540 let placement = self.worker_placement(session_id);
541 let probe = match placement {
542 Ok((backend, worker_root)) => tokio::task::spawn_blocking(move || {
543 super::worker_binary::probe_worker(
544 &targets::BoundedProcessExecutor::new(Duration::from_secs(5)),
545 &backend,
546 &worker_root,
547 )
548 })
549 .await
550 .context("join the suspension worker diagnostic probe")
551 .and_then(std::convert::identity),
552 Err(error) => Err(error.context("resolve the retained suspension worker")),
553 };
554 match probe {
555 Ok(probe) => tracing::warn!(
556 session_id,
557 phase,
558 worker_pids = ?probe.pids,
559 worker_startup_step = ?probe.step(),
560 worker_state = %probe,
561 "suspension failure worker diagnostic; no recovery mutation was attempted"
562 ),
563 Err(error) => tracing::warn!(
564 session_id,
565 phase,
566 error = format!("{error:#}"),
567 "suspension failure worker diagnostic was unavailable; worker state remains unknown"
568 ),
569 }
570 }
571
572 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
579 self.fail_interrupted_lifecycle_with(
580 session_id,
581 cause,
582 crate::database::save_lifecycle_session,
583 )
584 }
585
586 fn fail_interrupted_lifecycle_with(
587 &mut self,
588 session_id: &str,
589 cause: &str,
590 persist: impl Fn(&SessionRecord) -> Result<()>,
591 ) -> Result<bool> {
592 let Some(record) = self.state.sessions.get_mut(session_id) else {
593 return Ok(false);
594 };
595 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
596 return Ok(false);
597 }
598 let previous = record.clone();
599 record.state = SessionState::Error;
600 record.updated_at = now();
601 record.last_error = Some(cause.to_owned());
602 persist_session_record_transition_or_restore(
603 &mut self.state,
604 session_id,
605 &previous,
606 "persist the failure of an interrupted lifecycle state",
607 &persist,
608 )?;
609 Ok(true)
610 }
611
612 pub fn fail_unready_session(
621 &mut self,
622 session_id: &str,
623 cause: &str,
624 observed_updated_at: &str,
625 ) -> Result<bool> {
626 self.fail_unready_session_with(
627 session_id,
628 cause,
629 observed_updated_at,
630 crate::database::save_lifecycle_session,
631 )
632 }
633
634 fn fail_unready_session_with(
635 &mut self,
636 session_id: &str,
637 cause: &str,
638 observed_updated_at: &str,
639 persist: impl Fn(&SessionRecord) -> Result<()>,
640 ) -> Result<bool> {
641 let Some(record) = self.state.sessions.get_mut(session_id) else {
642 return Ok(false);
643 };
644 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
645 return Ok(false);
646 }
647 let previous = record.clone();
648 record.state = SessionState::Error;
649 record.updated_at = now();
650 record.last_error = Some(cause.to_owned());
651 persist_session_record_transition_or_restore(
652 &mut self.state,
653 session_id,
654 &previous,
655 "persist a session whose harness never became usable",
656 &persist,
657 )?;
658 Ok(true)
659 }
660
661 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
665 let Some(record) = self.state.sessions.get(session_id) else {
666 return Ok(false);
667 };
668 let previous = record.clone();
671 let record = self.state.sessions.get_mut(session_id).unwrap();
672 record.last_error = Some(cause.to_owned());
673 record.updated_at = now();
674 persist_session_record_transition_or_restore(
675 &mut self.state,
676 session_id,
677 &previous,
678 "persist the reason a close did not finish",
679 &crate::database::save_lifecycle_session,
680 )?;
681 Ok(true)
682 }
683
684 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
689 let Some(record) = self.state.sessions.get(session_id) else {
690 return Ok(false);
691 };
692 if record.public_error().is_none() {
693 return Ok(false);
694 }
695 let previous = record.clone();
696 let record = self.state.sessions.get_mut(session_id).unwrap();
697 record.last_error = None;
698 record.updated_at = now();
699 persist_session_record_transition_or_restore(
700 &mut self.state,
701 session_id,
702 &previous,
703 "clear the reason a close did not finish",
704 &crate::database::save_lifecycle_session,
705 )?;
706 Ok(true)
707 }
708
709 pub fn suspend_session_without_checkpoint(
722 &mut self,
723 session_id: &str,
724 executor: &impl CommandExecutor,
725 ) -> Result<bool> {
726 self.suspend_session_without_checkpoint_with(
727 session_id,
728 executor,
729 crate::database::save_lifecycle_session,
730 )
731 }
732
733 fn suspend_session_without_checkpoint_with(
734 &mut self,
735 session_id: &str,
736 executor: &impl CommandExecutor,
737 persist: impl Fn(&SessionRecord) -> Result<()>,
738 ) -> Result<bool> {
739 let session = self
740 .state
741 .sessions
742 .get(session_id)
743 .with_context(|| format!("unknown session {session_id}"))?
744 .clone();
745 ensure!(
746 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
747 "session {session_id} has a workspace to checkpoint; suspend it instead"
748 );
749 self.stop_target_and_settle(session_id, &session, executor, &persist)
750 }
751
752 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
753 let record = self.state.sessions.get_mut(session_id).unwrap();
754 apply_interrupted_close_error(record, error, &now());
755 self.persist_session_state(session_id)
756 }
757
758 fn destroy_after_verified_checkpoint(
761 &mut self,
762 session_id: &str,
763 verified: &CheckpointMetadata,
764 executor: &impl CommandExecutor,
765 ) -> Result<bool> {
766 self.destroy_after_verified_checkpoint_with(
767 session_id,
768 verified,
769 executor,
770 crate::database::save_lifecycle_session,
771 )
772 }
773
774 fn destroy_after_verified_checkpoint_with(
775 &mut self,
776 session_id: &str,
777 verified: &CheckpointMetadata,
778 executor: &impl CommandExecutor,
779 persist: impl Fn(&SessionRecord) -> Result<()>,
780 ) -> Result<bool> {
781 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
782 let _target_guard = target_mutex.lock().map_err(|_| {
783 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
784 })?;
785 let session = self
786 .state
787 .sessions
788 .get(session_id)
789 .with_context(|| format!("unknown session {session_id}"))?
790 .clone();
791 ensure!(
792 matches!(
793 session.state,
794 SessionState::Closing | SessionState::Destroying
795 ),
796 "refusing to destroy session {session_id}: it is not closing or destroying"
797 );
798 ensure!(
799 session.checkpoint.as_ref() == Some(verified),
800 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
801 );
802 if session.state == SessionState::Closing {
803 let record = self.state.sessions.get_mut(session_id).unwrap();
804 record.state = SessionState::Destroying;
805 record.updated_at = now();
806 record.last_error = None;
807 persist_session_record_transition_or_restore(
808 &mut self.state,
809 session_id,
810 &session,
811 "persist destroying state before target cleanup",
812 &persist,
813 )?;
814 }
815
816 let destroying = self
817 .state
818 .sessions
819 .get(session_id)
820 .expect("destroying session disappeared")
821 .clone();
822 {
823 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
824 verify_installed_checkpoint_gate(session_id, verified)?;
825 }
826 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
831 tracing::warn!(
832 session_id,
833 error = format!("{error:#}"),
834 "could not record that the second-opinion conversation ends with this target"
835 );
836 }
837 let locator = destroying
838 .target
839 .as_ref()
840 .context("session has no target")?;
841 let backend = backend_locator(locator, &destroying, &self.config)?;
842 let deferred = if self.state.subagents.contains_key(session_id) {
843 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
844 false
845 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
846 plan.execute(executor)?;
847 true
848 } else {
849 execute_target_cleanup(&backend, session_id, executor)?;
850 false
851 };
852 if let Some(worktree) = &destroying.managed_worktree {
853 retire_managed_worktree(executor, worktree)
854 .context("retire managed raw-session worktree after verified close")?;
855 }
856 let record = self.state.sessions.get_mut(session_id).unwrap();
857 record.state = SessionState::Stopped;
858 if !deferred {
859 record.target = None;
860 }
861 record.updated_at = now();
862 record.last_error = None;
863 persist_session_record_transition_or_restore(
864 &mut self.state,
865 session_id,
866 &destroying,
867 "persist stopped state after target cleanup",
868 &persist,
869 )?;
870 Ok(deferred)
871 }
872
873 pub fn cleanup_stopped_target(
877 &mut self,
878 session_id: &str,
879 executor: &impl CommandExecutor,
880 ) -> Result<()> {
881 self.cleanup_stopped_target_with(
882 session_id,
883 executor,
884 crate::database::save_lifecycle_session,
885 )
886 }
887
888 fn cleanup_stopped_target_with(
889 &mut self,
890 session_id: &str,
891 executor: &impl CommandExecutor,
892 persist: impl Fn(&SessionRecord) -> Result<()>,
893 ) -> Result<()> {
894 let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
895 let _target_guard = target_mutex.lock().map_err(|_| {
896 anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
897 })?;
898 let previous = self
899 .state
900 .sessions
901 .get(session_id)
902 .with_context(|| format!("unknown session {session_id}"))?
903 .clone();
904 ensure!(
905 previous.state == SessionState::Stopped,
906 "refusing deferred cleanup for active session {session_id}"
907 );
908 let Some(locator) = previous.target.as_ref() else {
909 return Ok(());
910 };
911 let backend = backend_locator(locator, &previous, &self.config)?;
912 ensure!(
913 targets::quiesce_plan(&backend, session_id)?.is_some(),
914 "session {session_id} retained a non-Podman target after stopping"
915 );
916 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
917 let record = self.state.sessions.get_mut(session_id).unwrap();
918 record.updated_at = now();
919 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
920 let persisted = persist_session_record_transition_or_restore(
921 &mut self.state,
922 session_id,
923 &previous,
924 "persist deferred target cleanup failure",
925 &persist,
926 );
927 return match persisted {
928 Ok(()) => Err(error),
929 Err(persist_error) => Err(error.context(format!(
930 "also failed to persist deferred target cleanup failure: {persist_error:#}"
931 ))),
932 };
933 }
934 let record = self.state.sessions.get_mut(session_id).unwrap();
935 record.target = None;
936 record.updated_at = now();
937 record.last_error = None;
938 persist_session_record_transition_or_restore(
939 &mut self.state,
940 session_id,
941 &previous,
942 "persist completion of deferred Podman target cleanup",
943 &persist,
944 )
945 }
946
947 pub fn force_stop(
950 &mut self,
951 session_id: &str,
952 executor: &impl CommandExecutor,
953 ) -> Result<bool> {
954 self.force_stop_with(
955 session_id,
956 executor,
957 crate::database::save_lifecycle_session,
958 )
959 }
960
961 fn force_stop_with(
962 &mut self,
963 session_id: &str,
964 executor: &impl CommandExecutor,
965 persist: impl Fn(&SessionRecord) -> Result<()>,
966 ) -> Result<bool> {
967 let session = self
968 .state
969 .sessions
970 .get(session_id)
971 .with_context(|| format!("unknown session {session_id}"))?
972 .clone();
973 ensure!(
974 session.state.is_active(),
975 "session {session_id} is already inactive"
976 );
977 let checkpoint = session
978 .checkpoint
979 .as_ref()
980 .context("force stop requires an existing recovery archive")?;
981 {
984 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
985 verify_installed_checkpoint_gate(session_id, checkpoint)
986 .context("verify the recovery archive before force stopping")?;
987 }
988 self.stop_target_and_settle(session_id, &session, executor, &persist)
989 }
990
991 fn stop_target_and_settle(
996 &mut self,
997 session_id: &str,
998 session: &SessionRecord,
999 executor: &impl CommandExecutor,
1000 persist: &impl Fn(&SessionRecord) -> Result<()>,
1001 ) -> Result<bool> {
1002 let mut deferred = false;
1003 if let Some(locator) = &session.target {
1004 let backend = backend_locator(locator, session, &self.config)?;
1005 if self.state.subagents.contains_key(session_id) {
1009 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
1010 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1011 plan.execute(executor)?;
1012 deferred = true;
1013 } else {
1014 execute_target_cleanup(&backend, session_id, executor)?;
1015 }
1016 }
1017 if let Some(worktree) = &session.managed_worktree {
1018 retire_managed_worktree(executor, worktree)
1019 .context("retire managed raw-session worktree after stopping the target")?;
1020 }
1021 let record = self.state.sessions.get_mut(session_id).unwrap();
1022 record.state = SessionState::Stopped;
1023 if !deferred {
1024 record.target = None;
1025 }
1026 record.updated_at = now();
1027 record.last_error = None;
1028 record.last_checkpoint_error = None;
1029 persist_session_record_transition_or_restore(
1030 &mut self.state,
1031 session_id,
1032 session,
1033 "persist stopped state after tearing down the current target",
1034 persist,
1035 )?;
1036 Ok(deferred)
1037 }
1038
1039 pub fn destroy_session_controlled(
1043 &mut self,
1044 session_id: &str,
1045 executor: &impl CommandExecutor,
1046 ) -> Result<()> {
1047 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1048 }
1049
1050 pub fn destroy_session_controlled_with(
1058 &mut self,
1059 session_id: &str,
1060 executor: &impl CommandExecutor,
1061 branch: BranchDisposition,
1062 ) -> Result<()> {
1063 self.destroy_session_controlled_with_checkout(
1064 session_id,
1065 executor,
1066 branch,
1067 CheckoutDisposition::Remove,
1068 )
1069 .map(|_| ())
1070 }
1071
1072 pub(crate) fn destroy_session_controlled_with_checkout(
1078 &mut self,
1079 session_id: &str,
1080 executor: &impl CommandExecutor,
1081 branch: BranchDisposition,
1082 checkout: CheckoutDisposition,
1083 ) -> Result<Option<PathBuf>> {
1084 let session = self
1085 .state
1086 .sessions
1087 .get(session_id)
1088 .with_context(|| format!("unknown session {session_id}"))?
1089 .clone();
1090 if session.state.is_active() {
1091 bail!("refusing to destroy active session {session_id}");
1092 }
1093 let mut retained_checkout = None;
1094 if let Some(worktree) = &session.managed_worktree {
1095 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1096 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1097 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1098 "check the managed raw-session worktree for uncommitted changes",
1099 )?);
1100 if keep {
1101 retained_checkout = Some(worktree.worktree_root.clone());
1105 } else {
1106 cleanup_managed_worktree(executor, worktree, branch)
1107 .context("remove managed raw-session worktree")?;
1108 }
1109 }
1110 if let Some(checkpoint) = &session.checkpoint
1111 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1112 && error.kind() != std::io::ErrorKind::NotFound
1113 {
1114 return Err(error).with_context(|| {
1115 format!(
1116 "remove session recovery archive {}",
1117 checkpoint.archive_path.display()
1118 )
1119 });
1120 }
1121 mj_core::attachment::AttachmentStore::controller(session_id)?
1122 .remove_session_data()
1123 .context("remove session image attachments")?;
1124 crate::database::delete_session(session_id)
1125 .context("destroy stopped session in database")?;
1126 self.state.subagents.remove(session_id);
1127 self.state.destroy_stopped_session(session_id)?;
1128 Ok(retained_checkout)
1129 }
1130
1131 pub fn force_destroy_session(
1143 &mut self,
1144 session_id: &str,
1145 executor: &impl CommandExecutor,
1146 branch: BranchDisposition,
1147 ) -> Result<()> {
1148 self.force_destroy_session_with(
1149 session_id,
1150 executor,
1151 branch,
1152 crate::database::delete_session,
1153 )
1154 }
1155
1156 fn force_destroy_session_with(
1157 &mut self,
1158 session_id: &str,
1159 executor: &impl CommandExecutor,
1160 branch: BranchDisposition,
1161 delete: impl Fn(&str) -> Result<()>,
1162 ) -> Result<()> {
1163 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1164 if result.is_err() {
1165 self.settle_failed_destroy(session_id);
1166 }
1167 result
1168 }
1169
1170 fn settle_failed_destroy(&mut self, session_id: &str) {
1176 let Some(record) = self.state.sessions.get(session_id) else {
1177 return;
1178 };
1179 if !record.state.is_active() || record.state == SessionState::Error {
1180 return;
1181 }
1182 let previous = record.clone();
1183 let record = self.state.sessions.get_mut(session_id).unwrap();
1184 record.state = SessionState::Error;
1185 record.updated_at = now();
1186 if let Err(error) = persist_session_record_transition_or_restore(
1187 &mut self.state,
1188 session_id,
1189 &previous,
1190 "persist the failed destroy",
1191 &crate::database::save_lifecycle_session,
1192 ) {
1193 tracing::warn!(
1194 session_id,
1195 error = format!("{error:#}"),
1196 "could not record a failed destroy"
1197 );
1198 }
1199 }
1200
1201 fn force_destroy_session_steps(
1202 &mut self,
1203 session_id: &str,
1204 executor: &impl CommandExecutor,
1205 branch: BranchDisposition,
1206 delete: impl Fn(&str) -> Result<()>,
1207 ) -> Result<()> {
1208 let session = self
1209 .state
1210 .sessions
1211 .get(session_id)
1212 .with_context(|| format!("unknown session {session_id}"))?
1213 .clone();
1214 if let Some(locator) = &session.target {
1218 let backend = backend_locator(locator, &session, &self.config)?;
1219 if self.state.subagents.contains_key(session_id) {
1220 targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
1221 } else {
1222 execute_target_cleanup(&backend, session_id, executor)?;
1223 }
1224 }
1225 if let Some(worktree) = &session.managed_worktree {
1226 cleanup_managed_worktree(executor, worktree, branch)
1227 .context("remove managed raw-session worktree")?;
1228 }
1229 if let Some(checkpoint) = &session.checkpoint
1230 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1231 && error.kind() != std::io::ErrorKind::NotFound
1232 {
1233 return Err(error).with_context(|| {
1234 format!(
1235 "remove session recovery archive {}",
1236 checkpoint.archive_path.display()
1237 )
1238 });
1239 }
1240 mj_core::attachment::AttachmentStore::controller(session_id)?
1241 .remove_session_data()
1242 .context("remove session image attachments")?;
1243 delete(session_id).context("force destroy session in database")?;
1244 self.state.subagents.remove(session_id);
1245 self.state.destroy_session_force(session_id)?;
1246 Ok(())
1247 }
1248}
1249
1250pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1265 match session.state {
1266 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1267 SessionState::Error if subagent => true,
1268 SessionState::Closing
1269 | SessionState::Destroying
1270 | SessionState::Error
1271 | SessionState::Lost => session.target.is_none(),
1272 SessionState::Running
1273 | SessionState::Disconnected
1274 | SessionState::Checkpointing
1275 | SessionState::Stopped
1276 | SessionState::DestroyedWithDataLoss => false,
1277 }
1278}
1279
1280fn execute_target_cleanup(
1281 backend: &targets::TargetLocator,
1282 session_id: &str,
1283 executor: &impl CommandExecutor,
1284) -> Result<()> {
1285 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1286 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1287 Ok(true) => {
1288 tracing::warn!(
1289 session_id,
1290 error = format!("{cleanup_error:#}"),
1291 "target cleanup command failed, but the target was confirmed absent"
1292 );
1293 }
1294 Ok(false) => {
1295 tracing::error!(
1296 session_id,
1297 error = format!("{cleanup_error:#}"),
1298 "target cleanup failed and the target is still present"
1299 );
1300 return Err(cleanup_error);
1301 }
1302 Err(probe_error) => {
1303 tracing::error!(
1304 session_id,
1305 cleanup_error = format!("{cleanup_error:#}"),
1306 probe_error = format!("{probe_error:#}"),
1307 "target cleanup failed and exact absence could not be confirmed"
1308 );
1309 return Err(cleanup_error.context(format!(
1310 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1311 )));
1312 }
1313 }
1314 }
1315 Ok(())
1316}
1317
1318fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1319 record.state = SessionState::Closing;
1320 record.updated_at = updated_at;
1321 record.last_checkpoint_error = None;
1322}
1323
1324fn apply_close_checkpoint_failure(
1331 record: &mut SessionRecord,
1332 previous: &SessionRecord,
1333 error: &anyhow::Error,
1334 updated_at: String,
1335) {
1336 if WorkerRestartLeftNoWorker::marks(error) {
1337 record.state = SessionState::Error;
1338 record.last_error = Some(format!(
1339 "close failed and left the session without a live worker; retry the close, \
1340 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1341 ));
1342 } else {
1343 record.state = state_after_unsealed_close(previous);
1344 }
1345 record.last_checkpoint_error = Some(format!("{error:#}"));
1346 record.updated_at = updated_at;
1347}
1348
1349fn close_was_refused(error: &anyhow::Error) -> bool {
1357 error.chain().any(|cause| {
1358 cause
1359 .downcast_ref::<crate::worker_client::RelayRejected>()
1360 .is_some_and(|rejected| !rejected.is_retryable())
1361 })
1362}
1363
1364fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1365 if previous.state == SessionState::Closing {
1366 SessionState::Running
1367 } else {
1368 previous.state
1369 }
1370}
1371
1372fn apply_interrupted_close_error(
1373 record: &mut SessionRecord,
1374 error: &anyhow::Error,
1375 updated_at: &str,
1376) {
1377 let destroying = record.state == SessionState::Destroying;
1378 if !destroying {
1379 record.state = SessionState::Closing;
1380 }
1381 record.updated_at = updated_at.to_owned();
1382 record.last_error = Some(if destroying {
1383 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1384 } else {
1385 format!("close is safely resumable from its verified checkpoint: {error:#}")
1386 });
1387}
1388
1389#[cfg(test)]
1390mod tests;