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 let retained = source_relay.owner();
138 crate::worker_lifecycle::run_with_owner(
139 session_id,
140 "suspend session for move",
141 executor,
142 retained,
143 async {
144 self.prepare_move_source_checkpoint(
145 session_id,
146 executor,
147 manager,
148 operation,
149 &mut source_relay,
150 )
151 .await?;
152 if disposition == SourceTargetDisposition::RetainForInPlaceSwap
153 || operation.workspace_transfer.is_some()
154 {
155 self.seal_move_handoff(
156 session_id,
157 executor,
158 manager,
159 operation,
160 preparation,
161 source_relay.take(),
162 )
163 .await?;
164 return Ok(false);
165 }
166 self.suspend_session_controlled_with_manager(
167 session_id,
168 executor,
169 Some(manager),
170 Some((operation, preparation)),
171 disposition,
172 true,
173 None,
174 source_relay.take(),
175 )
176 .await
177 },
178 )
179 .await
180 }
181
182 #[allow(clippy::too_many_arguments)]
183 async fn suspend_session_controlled_with_manager(
184 &mut self,
185 session_id: &str,
186 executor: &(impl CommandExecutor + Sync),
187 manager: Option<&SessionManagerControl>,
188 move_intent: Option<(
189 &mut mj_core::state::MoveOperation,
190 Option<&mj_core::state::MovePreparation>,
191 )>,
192 disposition: SourceTargetDisposition,
193 acknowledge_unpublished_work: bool,
194 before_close: Option<BeforeClose>,
195 held_relay: Option<super::checkpoint::ControllerRelayLease>,
196 ) -> Result<bool> {
197 crate::worker_lifecycle::run(session_id, "suspend session controlled with manager", executor, async {
198 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
199 let previous = self
200 .state
201 .sessions
202 .get(session_id)
203 .with_context(|| format!("unknown session {session_id}"))?
204 .clone();
205 let record = self.state.sessions.get_mut(session_id).unwrap();
206 apply_close_checkpoint_started(record, now());
210 self.persist_session_transition_or_restore(
211 session_id,
212 &previous,
213 "persist closing state before checkpointing the session",
214 )?;
215
216 let mut latched = match self
219 .checkpoint_session_latched_for_operation(
220 session_id,
221 executor,
222 manager,
223 LatchExclusivity::HoldThroughClose,
224 CheckpointExportPolicy::ReuseUnchangedArchive,
225 None,
226 held_relay,
227 )
228 .await
229 {
230 Ok(latched) => latched,
231 Err(error) => {
232 let record = self.state.sessions.get_mut(session_id).unwrap();
233 apply_close_checkpoint_failure(record, &previous, &error, now());
237 return Err(
238 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
239 );
240 }
241 };
242
243 let artifact = latched.artifact.clone();
244 let publication = if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
247 None
248 } else {
249 match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
250 Some(ManagedCheckoutKind::Clone) => Some(
251 super::publication::assess_clone_checkpoint(&previous, &artifact.metadata),
252 ),
253 Some(ManagedCheckoutKind::Worktree) => None,
254 None if previous.project_directory.is_none() => {
255 Some(super::publication::assess_network_checkpoint(
256 &previous,
257 &artifact.metadata,
258 &self.config,
259 ))
260 }
261 None => None,
262 }
263 };
264 if !acknowledge_unpublished_work
265 && publication
266 .as_ref()
267 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
268 {
269 latched.relay.cancel_abandoned_barrier().await?;
270 let mut restored = previous.clone();
274 restored.state = state_after_unsealed_close(&previous);
275 restored.updated_at = now();
276 self.state.sessions.insert(session_id.to_owned(), restored);
277 self.persist_session_transition_or_restore(
278 session_id,
279 &previous,
280 "restore session after unpublished-work preflight",
281 )?;
282 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
283 return Err(mj_core::refusal::Refusal::precondition(
284 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
285 ).into());
286 }
287 let record = self.state.sessions.get_mut(session_id).unwrap();
288 record.state = SessionState::Closing;
289 record.native_session_id = Some(artifact.native_session_id.clone());
290 record.checkpoint = Some(artifact.metadata.clone());
291 record.publication = publication;
292 record.updated_at = now();
293 record.last_error = None;
294 record.last_checkpoint_error = None;
295 self.persist_checkpoint_transition_or_restore(
296 session_id,
297 &previous,
298 "persist verified checkpoint and closing state before sealing the relay",
299 )?;
300 if let Some((operation, preparation)) = move_intent {
301 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
304 let record = self.state.sessions.get_mut(session_id).unwrap();
305 record.state = previous.state;
306 record.last_error = Some(format!("{error:#}"));
307 self.persist_session_transition_or_restore(
308 session_id,
309 &previous,
310 "restore source after move preflight failure",
311 )?;
312 return Err(error);
313 }
314 operation.checkpoint = Some(artifact.metadata.clone());
315 operation.updated_at = now();
316 crate::database::save_move_operation(operation)?;
317 }
318 if let Some(before_close) = before_close
319 && let Err(error) = before_close.await
320 {
321 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
324 tracing::warn!(
325 session_id,
326 error = format!("{cancel:#}"),
327 "could not release the checkpoint barrier of a close that did not proceed"
328 );
329 }
330 let record = self.state.sessions.get_mut(session_id).unwrap();
331 record.state = state_after_unsealed_close(&previous);
332 record.last_error = Some(format!("{error:#}"));
333 record.updated_at = now();
334 self.persist_session_transition_or_restore(
335 session_id,
336 &previous,
337 "restore a session whose close did not proceed past its checkpoint",
338 )?;
339 return Err(error);
340 }
341 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
342 release_projection_behind_checkpoint(session_id, &artifact.metadata);
345
346 let close_command_id = new_command_id("close")?;
347 let barrier_command_id = latched.barrier_command_id.clone();
348 let close_result = {
349 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
350 latched
351 .relay
352 .connection_mut()
353 .submit(
354 close_command_id,
355 RelayCommand::Close {
356 barrier_command_id: barrier_command_id.clone(),
357 expected: latched.cursor.clone(),
358 },
359 )
360 .await
361 };
362 if let Err(error) = close_result {
363 let error = error.context("seal verified checkpoint for close");
364 if !close_was_refused(&error) {
365 self.record_interrupted_close(session_id, &error)?;
368 return Err(error);
369 }
370 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
375 tracing::warn!(
376 session_id,
377 error = format!("{cancel:#}"),
378 "could not release the checkpoint barrier of a refused close"
379 );
380 }
381 let record = self.state.sessions.get_mut(session_id).unwrap();
382 record.state = state_after_unsealed_close(&previous);
383 record.last_error = Some(format!("{error:#}"));
384 record.updated_at = now();
385 self.persist_session_transition_or_restore(
386 session_id,
387 &previous,
388 "restore a session whose worker refused its close",
389 )?;
390 return Err(error);
391 }
392 let close_result = {
393 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
394 latched
395 .relay
396 .connection_mut()
397 .submit(
398 new_command_id("checkpoint-complete")?,
399 RelayCommand::CompleteCheckpoint { barrier_command_id },
400 )
401 .await
402 };
403 if let Err(error) = close_result {
404 self.record_interrupted_close(session_id, &error)?;
405 return Err(error.context("release verified close checkpoint"));
406 }
407 let close_result = {
408 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
409 wait_for_relay_closed(latched.relay.connection_mut()).await
410 };
411 if let Err(error) = close_result {
412 self.record_interrupted_close(session_id, &error)?;
413 return Err(error);
414 }
415 latched.relay.release();
416
417 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
418 return Ok(false);
424 }
425 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
426 Ok(deferred) => Ok(deferred),
427 Err(error) => {
428 self.record_interrupted_close(session_id, &error)?;
429 Err(error)
430 }
431 }
432
433 }).await
434 }
435
436 pub async fn recover_interrupted_close_managed(
448 &mut self,
449 session_id: &str,
450 executor: &(impl CommandExecutor + Sync),
451 manager: &SessionManagerControl,
452 acknowledge_unpublished_work: bool,
453 before_close: Option<BeforeClose>,
454 ) -> Result<bool> {
455 crate::worker_lifecycle::run(session_id, "recover interrupted close managed", executor, async {
456 let (state, verified) = {
457 let session = self
458 .state
459 .sessions
460 .get(session_id)
461 .with_context(|| format!("unknown session {session_id}"))?;
462 ensure!(
463 matches!(
464 session.state,
465 SessionState::Closing | SessionState::Destroying
466 ),
467 "session {session_id} has no interrupted close to recover"
468 );
469 (session.state, session.checkpoint.clone())
470 };
471 if state == SessionState::Destroying {
472 let verified = verified.context("destroying session has no verified checkpoint")?;
473 if let Some(before_close) = before_close {
474 before_close.await?;
475 }
476 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
477 }
478 ensure!(
479 state == SessionState::Closing,
480 "session {session_id} has no relay close to recover"
481 );
482 let started = std::time::Instant::now();
483 let mut phase = "wait for session actor";
484 tracing::info!(
485 session_id,
486 ?state,
487 checkpoint_frontier = ?verified.as_ref().map(|checkpoint| checkpoint.event_frontier),
488 "recovering suspension: asking the worker whether its relay was sealed"
489 );
490 let connection = async {
491 let handle = manager
492 .wait_for_session(session_id, Duration::from_secs(5))
493 .await?;
494 phase = "lease worker connection";
495 let mut lease = handle.lease_connection().await?;
496 phase = "read worker execution state";
497 let execution = lease.connection_mut().sync().await?.operational.execution;
498 anyhow::Ok((lease, execution))
499 }
500 .await;
501 let (mut lease, execution) = match connection {
502 Ok(connected) => connected,
503 Err(error) => {
504 tracing::warn!(
505 session_id,
506 phase,
507 elapsed_ms = started.elapsed().as_millis() as u64,
508 error = format!("{error:#}"),
509 "suspension could not query its worker; retaining target and checkpoint"
510 );
511 self.diagnose_suspension_worker(session_id, phase).await;
512 return Err(error.context(format!("{phase} while recovering suspension")));
513 }
514 };
515 tracing::info!(
516 session_id,
517 ?execution,
518 elapsed_ms = started.elapsed().as_millis() as u64,
519 "suspension recovery read the worker's execution state"
520 );
521 match execution {
522 RelayExecutionState::Closed => {}
523 RelayExecutionState::Closing => {
524 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
525 wait_for_relay_closed(lease.connection_mut()).await?;
526 }
527 RelayExecutionState::Idle | RelayExecutionState::Running => {
528 lease.release();
529 return self
530 .suspend_session_controlled_with_manager(
531 session_id,
532 executor,
533 Some(manager),
534 None,
535 SourceTargetDisposition::Destroy,
536 acknowledge_unpublished_work,
537 before_close,
538 None,
539 )
540 .await;
541 }
542 }
543 lease.release();
544 let verified = verified.context("closed relay has no verified checkpoint")?;
545 if let Some(before_close) = before_close {
547 before_close.await?;
548 }
549 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
550
551 }).await
552 }
553
554 async fn diagnose_suspension_worker(&self, session_id: &str, phase: &str) {
557 let placement = self.worker_placement(session_id);
558 let probe = match placement {
559 Ok((backend, worker_root)) => tokio::task::spawn_blocking(move || {
560 super::worker_binary::probe_worker(
561 &targets::BoundedProcessExecutor::new(Duration::from_secs(5)),
562 &backend,
563 &worker_root,
564 )
565 })
566 .await
567 .context("join the suspension worker diagnostic probe")
568 .and_then(std::convert::identity),
569 Err(error) => Err(error.context("resolve the retained suspension worker")),
570 };
571 match probe {
572 Ok(probe) => tracing::warn!(
573 session_id,
574 phase,
575 worker_pids = ?probe.pids,
576 worker_startup_step = ?probe.step(),
577 worker_state = %probe,
578 "suspension failure worker diagnostic; no recovery mutation was attempted"
579 ),
580 Err(error) => tracing::warn!(
581 session_id,
582 phase,
583 error = format!("{error:#}"),
584 "suspension failure worker diagnostic was unavailable; worker state remains unknown"
585 ),
586 }
587 }
588
589 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
596 self.fail_interrupted_lifecycle_with(
597 session_id,
598 cause,
599 crate::database::save_lifecycle_session,
600 )
601 }
602
603 fn fail_interrupted_lifecycle_with(
604 &mut self,
605 session_id: &str,
606 cause: &str,
607 persist: impl Fn(&SessionRecord) -> Result<()>,
608 ) -> Result<bool> {
609 let Some(record) = self.state.sessions.get_mut(session_id) else {
610 return Ok(false);
611 };
612 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
613 return Ok(false);
614 }
615 let previous = record.clone();
616 record.state = SessionState::Error;
617 record.updated_at = now();
618 record.last_error = Some(cause.to_owned());
619 persist_session_record_transition_or_restore(
620 &mut self.state,
621 session_id,
622 &previous,
623 "persist the failure of an interrupted lifecycle state",
624 &persist,
625 )?;
626 Ok(true)
627 }
628
629 pub fn fail_unready_session(
638 &mut self,
639 session_id: &str,
640 cause: &str,
641 observed_updated_at: &str,
642 ) -> Result<bool> {
643 self.fail_unready_session_with(
644 session_id,
645 cause,
646 observed_updated_at,
647 crate::database::save_lifecycle_session,
648 )
649 }
650
651 fn fail_unready_session_with(
652 &mut self,
653 session_id: &str,
654 cause: &str,
655 observed_updated_at: &str,
656 persist: impl Fn(&SessionRecord) -> Result<()>,
657 ) -> Result<bool> {
658 let Some(record) = self.state.sessions.get_mut(session_id) else {
659 return Ok(false);
660 };
661 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
662 return Ok(false);
663 }
664 let previous = record.clone();
665 record.state = SessionState::Error;
666 record.updated_at = now();
667 record.last_error = Some(cause.to_owned());
668 persist_session_record_transition_or_restore(
669 &mut self.state,
670 session_id,
671 &previous,
672 "persist a session whose harness never became usable",
673 &persist,
674 )?;
675 Ok(true)
676 }
677
678 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
682 let Some(record) = self.state.sessions.get(session_id) else {
683 return Ok(false);
684 };
685 let previous = record.clone();
688 let record = self.state.sessions.get_mut(session_id).unwrap();
689 record.last_error = Some(cause.to_owned());
690 record.updated_at = now();
691 persist_session_record_transition_or_restore(
692 &mut self.state,
693 session_id,
694 &previous,
695 "persist the reason a close did not finish",
696 &crate::database::save_lifecycle_session,
697 )?;
698 Ok(true)
699 }
700
701 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
706 let Some(record) = self.state.sessions.get(session_id) else {
707 return Ok(false);
708 };
709 if record.public_error().is_none() {
710 return Ok(false);
711 }
712 let previous = record.clone();
713 let record = self.state.sessions.get_mut(session_id).unwrap();
714 record.last_error = None;
715 record.updated_at = now();
716 persist_session_record_transition_or_restore(
717 &mut self.state,
718 session_id,
719 &previous,
720 "clear the reason a close did not finish",
721 &crate::database::save_lifecycle_session,
722 )?;
723 Ok(true)
724 }
725
726 pub fn suspend_session_without_checkpoint(
739 &mut self,
740 session_id: &str,
741 executor: &impl CommandExecutor,
742 ) -> Result<bool> {
743 self.suspend_session_without_checkpoint_with(
744 session_id,
745 executor,
746 crate::database::save_lifecycle_session,
747 )
748 }
749
750 fn suspend_session_without_checkpoint_with(
751 &mut self,
752 session_id: &str,
753 executor: &impl CommandExecutor,
754 persist: impl Fn(&SessionRecord) -> Result<()>,
755 ) -> Result<bool> {
756 let session = self
757 .state
758 .sessions
759 .get(session_id)
760 .with_context(|| format!("unknown session {session_id}"))?
761 .clone();
762 ensure!(
763 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
764 "session {session_id} has a workspace to checkpoint; suspend it instead"
765 );
766 self.stop_target_and_settle(session_id, &session, executor, &persist)
767 }
768
769 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
770 let record = self.state.sessions.get_mut(session_id).unwrap();
771 apply_interrupted_close_error(record, error, &now());
772 self.persist_session_state(session_id)
773 }
774
775 fn destroy_after_verified_checkpoint(
778 &mut self,
779 session_id: &str,
780 verified: &CheckpointMetadata,
781 executor: &impl CommandExecutor,
782 ) -> Result<bool> {
783 self.destroy_after_verified_checkpoint_with(
784 session_id,
785 verified,
786 executor,
787 crate::database::save_lifecycle_session,
788 )
789 }
790
791 fn destroy_after_verified_checkpoint_with(
792 &mut self,
793 session_id: &str,
794 verified: &CheckpointMetadata,
795 executor: &impl CommandExecutor,
796 persist: impl Fn(&SessionRecord) -> Result<()>,
797 ) -> Result<bool> {
798 crate::worker_lifecycle::run_blocking(
799 session_id,
800 "destroy after verified checkpoint with",
801 executor,
802 || {
803 let session = self
804 .state
805 .sessions
806 .get(session_id)
807 .with_context(|| format!("unknown session {session_id}"))?
808 .clone();
809 ensure!(
810 matches!(
811 session.state,
812 SessionState::Closing | SessionState::Destroying
813 ),
814 "refusing to destroy session {session_id}: it is not closing or destroying"
815 );
816 ensure!(
817 session.checkpoint.as_ref() == Some(verified),
818 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
819 );
820 if session.state == SessionState::Closing {
821 let record = self.state.sessions.get_mut(session_id).unwrap();
822 record.state = SessionState::Destroying;
823 record.updated_at = now();
824 record.last_error = None;
825 persist_session_record_transition_or_restore(
826 &mut self.state,
827 session_id,
828 &session,
829 "persist destroying state before target cleanup",
830 &persist,
831 )?;
832 }
833
834 let destroying = self
835 .state
836 .sessions
837 .get(session_id)
838 .expect("destroying session disappeared")
839 .clone();
840 {
841 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
842 verify_installed_checkpoint_gate(session_id, verified)?;
843 }
844 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
849 tracing::warn!(
850 session_id,
851 error = format!("{error:#}"),
852 "could not record that the second-opinion conversation ends with this target"
853 );
854 }
855 let locator = destroying
856 .target
857 .as_ref()
858 .context("session has no target")?;
859 let backend = backend_locator(locator, &destroying, &self.config)?;
860 let deferred = if self.state.subagents.contains_key(session_id) {
861 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
862 .execute(executor)?;
863 false
864 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
865 plan.execute(executor)?;
866 true
867 } else {
868 execute_target_cleanup(&backend, session_id, executor)?;
869 false
870 };
871 if let Some(worktree) = &destroying.managed_worktree {
872 retire_managed_worktree(executor, worktree)
873 .context("retire managed raw-session worktree after verified close")?;
874 }
875 let record = self.state.sessions.get_mut(session_id).unwrap();
876 record.state = SessionState::Stopped;
877 if !deferred {
878 record.target = None;
879 }
880 record.updated_at = now();
881 record.last_error = None;
882 persist_session_record_transition_or_restore(
883 &mut self.state,
884 session_id,
885 &destroying,
886 "persist stopped state after target cleanup",
887 &persist,
888 )?;
889 Ok(deferred)
890 },
891 )
892 }
893
894 pub fn cleanup_stopped_target(
898 &mut self,
899 session_id: &str,
900 executor: &impl CommandExecutor,
901 ) -> Result<()> {
902 self.cleanup_stopped_target_with(
903 session_id,
904 executor,
905 crate::database::save_lifecycle_session,
906 )
907 }
908
909 fn cleanup_stopped_target_with(
910 &mut self,
911 session_id: &str,
912 executor: &impl CommandExecutor,
913 persist: impl Fn(&SessionRecord) -> Result<()>,
914 ) -> Result<()> {
915 crate::worker_lifecycle::run_blocking(
916 session_id,
917 "cleanup stopped target with",
918 executor,
919 || {
920 let previous = self
921 .state
922 .sessions
923 .get(session_id)
924 .with_context(|| format!("unknown session {session_id}"))?
925 .clone();
926 ensure!(
927 previous.state == SessionState::Stopped,
928 "refusing deferred cleanup for active session {session_id}"
929 );
930 let Some(locator) = previous.target.as_ref() else {
931 return Ok(());
932 };
933 let backend = backend_locator(locator, &previous, &self.config)?;
934 ensure!(
935 targets::quiesce_plan(&backend, session_id)?.is_some(),
936 "session {session_id} retained a non-Podman target after stopping"
937 );
938 if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
939 let record = self.state.sessions.get_mut(session_id).unwrap();
940 record.updated_at = now();
941 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
942 let persisted = persist_session_record_transition_or_restore(
943 &mut self.state,
944 session_id,
945 &previous,
946 "persist deferred target cleanup failure",
947 &persist,
948 );
949 return match persisted {
950 Ok(()) => Err(error),
951 Err(persist_error) => Err(error.context(format!(
952 "also failed to persist deferred target cleanup failure: {persist_error:#}"
953 ))),
954 };
955 }
956 let record = self.state.sessions.get_mut(session_id).unwrap();
957 record.target = None;
958 record.updated_at = now();
959 record.last_error = None;
960 persist_session_record_transition_or_restore(
961 &mut self.state,
962 session_id,
963 &previous,
964 "persist completion of deferred Podman target cleanup",
965 &persist,
966 )
967 },
968 )
969 }
970
971 pub fn force_stop(
974 &mut self,
975 session_id: &str,
976 executor: &impl CommandExecutor,
977 ) -> Result<bool> {
978 self.force_stop_with(
979 session_id,
980 executor,
981 crate::database::save_lifecycle_session,
982 )
983 }
984
985 fn force_stop_with(
986 &mut self,
987 session_id: &str,
988 executor: &impl CommandExecutor,
989 persist: impl Fn(&SessionRecord) -> Result<()>,
990 ) -> Result<bool> {
991 let session = self
992 .state
993 .sessions
994 .get(session_id)
995 .with_context(|| format!("unknown session {session_id}"))?
996 .clone();
997 ensure!(
998 session.state.is_active(),
999 "session {session_id} is already inactive"
1000 );
1001 let checkpoint = session
1002 .checkpoint
1003 .as_ref()
1004 .context("force stop requires an existing recovery archive")?;
1005 {
1008 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1009 verify_installed_checkpoint_gate(session_id, checkpoint)
1010 .context("verify the recovery archive before force stopping")?;
1011 }
1012 self.stop_target_and_settle(session_id, &session, executor, &persist)
1013 }
1014
1015 fn stop_target_and_settle(
1020 &mut self,
1021 session_id: &str,
1022 session: &SessionRecord,
1023 executor: &impl CommandExecutor,
1024 persist: &impl Fn(&SessionRecord) -> Result<()>,
1025 ) -> Result<bool> {
1026 crate::worker_lifecycle::run_blocking(
1027 session_id,
1028 "stop target and settle",
1029 executor,
1030 || {
1031 let mut deferred = false;
1032 if let Some(locator) = &session.target {
1033 let backend = backend_locator(locator, session, &self.config)?;
1034 if self.state.subagents.contains_key(session_id) {
1038 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1039 .execute(executor)?;
1040 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1041 plan.execute(executor)?;
1042 deferred = true;
1043 } else {
1044 execute_target_cleanup(&backend, session_id, executor)?;
1045 }
1046 }
1047 if let Some(worktree) = &session.managed_worktree {
1048 retire_managed_worktree(executor, worktree)
1049 .context("retire managed raw-session worktree after stopping the target")?;
1050 }
1051 let record = self.state.sessions.get_mut(session_id).unwrap();
1052 record.state = SessionState::Stopped;
1053 if !deferred {
1054 record.target = None;
1055 }
1056 record.updated_at = now();
1057 record.last_error = None;
1058 record.last_checkpoint_error = None;
1059 persist_session_record_transition_or_restore(
1060 &mut self.state,
1061 session_id,
1062 session,
1063 "persist stopped state after tearing down the current target",
1064 persist,
1065 )?;
1066 Ok(deferred)
1067 },
1068 )
1069 }
1070
1071 pub fn destroy_session_controlled(
1075 &mut self,
1076 session_id: &str,
1077 executor: &impl CommandExecutor,
1078 ) -> Result<()> {
1079 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1080 }
1081
1082 pub fn destroy_session_controlled_with(
1090 &mut self,
1091 session_id: &str,
1092 executor: &impl CommandExecutor,
1093 branch: BranchDisposition,
1094 ) -> Result<()> {
1095 self.destroy_session_controlled_with_checkout(
1096 session_id,
1097 executor,
1098 branch,
1099 CheckoutDisposition::Remove,
1100 )
1101 .map(|_| ())
1102 }
1103
1104 pub(crate) fn destroy_session_controlled_with_checkout(
1110 &mut self,
1111 session_id: &str,
1112 executor: &impl CommandExecutor,
1113 branch: BranchDisposition,
1114 checkout: CheckoutDisposition,
1115 ) -> Result<Option<PathBuf>> {
1116 let session = self
1117 .state
1118 .sessions
1119 .get(session_id)
1120 .with_context(|| format!("unknown session {session_id}"))?
1121 .clone();
1122 if session.state.is_active() {
1123 bail!("refusing to destroy active session {session_id}");
1124 }
1125 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1126 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1127 }
1128 let mut retained_checkout = None;
1129 if let Some(worktree) = &session.managed_worktree {
1130 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1131 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1132 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1133 "check the managed raw-session worktree for uncommitted changes",
1134 )?);
1135 if keep {
1136 retained_checkout = Some(worktree.worktree_root.clone());
1140 } else {
1141 cleanup_managed_worktree(executor, worktree, branch)
1142 .context("remove managed raw-session worktree")?;
1143 }
1144 }
1145 if let Some(checkpoint) = &session.checkpoint
1146 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1147 && error.kind() != std::io::ErrorKind::NotFound
1148 {
1149 return Err(error).with_context(|| {
1150 format!(
1151 "remove session recovery archive {}",
1152 checkpoint.archive_path.display()
1153 )
1154 });
1155 }
1156 mj_core::attachment::AttachmentStore::controller(session_id)?
1157 .remove_session_data()
1158 .context("remove session image attachments")?;
1159 crate::database::delete_session(session_id)
1160 .context("destroy stopped session in database")?;
1161 self.state.subagents.remove(session_id);
1162 self.state.destroy_stopped_session(session_id)?;
1163 Ok(retained_checkout)
1164 }
1165
1166 pub fn force_destroy_session(
1178 &mut self,
1179 session_id: &str,
1180 executor: &impl CommandExecutor,
1181 branch: BranchDisposition,
1182 ) -> Result<()> {
1183 self.force_destroy_session_with(
1184 session_id,
1185 executor,
1186 branch,
1187 crate::database::delete_session,
1188 )
1189 }
1190
1191 fn force_destroy_session_with(
1192 &mut self,
1193 session_id: &str,
1194 executor: &impl CommandExecutor,
1195 branch: BranchDisposition,
1196 delete: impl Fn(&str) -> Result<()>,
1197 ) -> Result<()> {
1198 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1199 if result.is_err() {
1200 self.settle_failed_destroy(session_id);
1201 }
1202 result
1203 }
1204
1205 fn settle_failed_destroy(&mut self, session_id: &str) {
1211 let Some(record) = self.state.sessions.get(session_id) else {
1212 return;
1213 };
1214 if !record.state.is_active() || record.state == SessionState::Error {
1215 return;
1216 }
1217 let previous = record.clone();
1218 let record = self.state.sessions.get_mut(session_id).unwrap();
1219 record.state = SessionState::Error;
1220 record.updated_at = now();
1221 if let Err(error) = persist_session_record_transition_or_restore(
1222 &mut self.state,
1223 session_id,
1224 &previous,
1225 "persist the failed destroy",
1226 &crate::database::save_lifecycle_session,
1227 ) {
1228 tracing::warn!(
1229 session_id,
1230 error = format!("{error:#}"),
1231 "could not record a failed destroy"
1232 );
1233 }
1234 }
1235
1236 fn force_destroy_session_steps(
1237 &mut self,
1238 session_id: &str,
1239 executor: &impl CommandExecutor,
1240 branch: BranchDisposition,
1241 delete: impl Fn(&str) -> Result<()>,
1242 ) -> Result<()> {
1243 crate::worker_lifecycle::run_blocking(
1244 session_id,
1245 "force destroy session steps",
1246 executor,
1247 || {
1248 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1249 let session = self
1250 .state
1251 .sessions
1252 .get(session_id)
1253 .with_context(|| format!("unknown session {session_id}"))?
1254 .clone();
1255 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1256 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1257 }
1258 if let Some(locator) = &session.target {
1262 let backend = backend_locator(locator, &session, &self.config)?;
1263 if self.state.subagents.contains_key(session_id) {
1264 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1265 .execute(executor)?;
1266 } else {
1267 execute_target_cleanup(&backend, session_id, executor)?;
1268 }
1269 }
1270 if let Some(worktree) = &session.managed_worktree {
1271 cleanup_managed_worktree(executor, worktree, branch)
1272 .context("remove managed raw-session worktree")?;
1273 }
1274 if let Some(checkpoint) = &session.checkpoint
1275 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1276 && error.kind() != std::io::ErrorKind::NotFound
1277 {
1278 return Err(error).with_context(|| {
1279 format!(
1280 "remove session recovery archive {}",
1281 checkpoint.archive_path.display()
1282 )
1283 });
1284 }
1285 mj_core::attachment::AttachmentStore::controller(session_id)?
1286 .remove_session_data()
1287 .context("remove session image attachments")?;
1288 delete(session_id).context("force destroy session in database")?;
1289 self.state.subagents.remove(session_id);
1290 self.state.destroy_session_force(session_id)?;
1291 Ok(())
1292 },
1293 )
1294 }
1295}
1296
1297pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1312 match session.state {
1313 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1314 SessionState::Error if subagent => true,
1315 SessionState::Closing
1316 | SessionState::Destroying
1317 | SessionState::Error
1318 | SessionState::Lost => session.target.is_none(),
1319 SessionState::Running
1320 | SessionState::Disconnected
1321 | SessionState::Checkpointing
1322 | SessionState::Stopped
1323 | SessionState::DestroyedWithDataLoss => false,
1324 }
1325}
1326
1327fn execute_target_cleanup(
1328 backend: &targets::TargetLocator,
1329 session_id: &str,
1330 executor: &impl CommandExecutor,
1331) -> Result<()> {
1332 let _owner = crate::worker_lifecycle::require(session_id)?;
1333 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1334 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1335 Ok(true) => {
1336 tracing::warn!(
1337 session_id,
1338 error = format!("{cleanup_error:#}"),
1339 "target cleanup command failed, but the target was confirmed absent"
1340 );
1341 }
1342 Ok(false) => {
1343 tracing::error!(
1344 session_id,
1345 error = format!("{cleanup_error:#}"),
1346 "target cleanup failed and the target is still present"
1347 );
1348 return Err(cleanup_error);
1349 }
1350 Err(probe_error) => {
1351 tracing::error!(
1352 session_id,
1353 cleanup_error = format!("{cleanup_error:#}"),
1354 probe_error = format!("{probe_error:#}"),
1355 "target cleanup failed and exact absence could not be confirmed"
1356 );
1357 return Err(cleanup_error.context(format!(
1358 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1359 )));
1360 }
1361 }
1362 }
1363 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
1366 crate::database::finish_worker_restart(session_id, &intent.operation_id)?;
1367 }
1368 Ok(())
1369}
1370
1371fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1372 record.state = SessionState::Closing;
1373 record.updated_at = updated_at;
1374 record.last_checkpoint_error = None;
1375}
1376
1377fn apply_close_checkpoint_failure(
1384 record: &mut SessionRecord,
1385 previous: &SessionRecord,
1386 error: &anyhow::Error,
1387 updated_at: String,
1388) {
1389 if WorkerRestartLeftNoWorker::marks(error) {
1390 record.state = SessionState::Error;
1391 record.last_error = Some(format!(
1392 "close failed and left the session without a live worker; retry the close, \
1393 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1394 ));
1395 } else {
1396 record.state = state_after_unsealed_close(previous);
1397 }
1398 record.last_checkpoint_error = Some(format!("{error:#}"));
1399 record.updated_at = updated_at;
1400}
1401
1402fn close_was_refused(error: &anyhow::Error) -> bool {
1410 error.chain().any(|cause| {
1411 cause
1412 .downcast_ref::<crate::worker_client::RelayRejected>()
1413 .is_some_and(|rejected| !rejected.is_retryable())
1414 })
1415}
1416
1417fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1418 if previous.state == SessionState::Closing {
1419 SessionState::Running
1420 } else {
1421 previous.state
1422 }
1423}
1424
1425fn apply_interrupted_close_error(
1426 record: &mut SessionRecord,
1427 error: &anyhow::Error,
1428 updated_at: &str,
1429) {
1430 let destroying = record.state == SessionState::Destroying;
1431 if !destroying {
1432 record.state = SessionState::Closing;
1433 }
1434 record.updated_at = updated_at.to_owned();
1435 record.last_error = Some(if destroying {
1436 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1437 } else {
1438 format!("close is safely resumable from its verified checkpoint: {error:#}")
1439 });
1440}
1441
1442#[cfg(test)]
1443mod tests;