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::mbx::release::BuildStateRelease;
20use super::worker_restart::WorkerRestartLeftNoWorker;
21use super::worktree::{
22 cleanup_managed_worktree, managed_worktree_checkout_is_dirty, retire_managed_worktree,
23};
24use super::{Controller, now, persist_session_record_transition_or_restore};
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub enum BranchDisposition {
29 Delete,
32 Keep,
35 DeleteIfMerged,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45pub enum CheckoutDisposition {
46 Remove,
49 KeepWhenDirty,
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub(super) enum SourceTargetDisposition {
58 Destroy,
60 RetainForInPlaceSwap,
63}
64
65pub type BeforeClose = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>;
72
73impl Controller {
74 pub async fn suspend_session(&mut self, session_id: &str) -> Result<()> {
79 self.suspend_session_controlled(session_id, &ProcessExecutor)
80 .await
81 }
82
83 pub async fn suspend_session_controlled(
84 &mut self,
85 session_id: &str,
86 executor: &(impl CommandExecutor + Sync),
87 ) -> Result<()> {
88 if self
89 .suspend_session_controlled_with_manager(
90 session_id,
91 executor,
92 None,
93 None,
94 SourceTargetDisposition::Destroy,
95 CheckpointExportPolicy::ReuseUnchangedArchive,
96 true,
97 None,
98 None,
99 )
100 .await?
101 {
102 self.cleanup_stopped_target(session_id, executor)?;
103 }
104 Ok(())
105 }
106
107 pub async fn suspend_session_managed_controlled(
108 &mut self,
109 session_id: &str,
110 executor: &(impl CommandExecutor + Sync),
111 manager: &SessionManagerControl,
112 acknowledge_unpublished_work: bool,
113 before_close: Option<BeforeClose>,
114 ) -> Result<bool> {
115 self.suspend_session_controlled_with_manager(
116 session_id,
117 executor,
118 Some(manager),
119 None,
120 SourceTargetDisposition::Destroy,
121 CheckpointExportPolicy::ReuseUnchangedArchive,
122 acknowledge_unpublished_work,
123 before_close,
124 None,
125 )
126 .await
127 }
128
129 pub(crate) async fn suspend_session_for_restart(
132 &mut self,
133 session_id: &str,
134 executor: &(impl CommandExecutor + Sync),
135 manager: &SessionManagerControl,
136 before_close: Option<BeforeClose>,
137 ) -> Result<()> {
138 self.suspend_session_controlled_with_manager(
139 session_id,
140 executor,
141 Some(manager),
142 None,
143 SourceTargetDisposition::RetainForInPlaceSwap,
144 CheckpointExportPolicy::Always,
145 true,
146 before_close,
147 None,
148 )
149 .await?;
150 Ok(())
151 }
152
153 #[allow(clippy::too_many_arguments)]
154 pub(super) async fn suspend_session_for_move(
155 &mut self,
156 session_id: &str,
157 executor: &(impl CommandExecutor + Sync),
158 manager: &SessionManagerControl,
159 operation: &mut mj_core::state::MoveOperation,
160 preparation: Option<&mj_core::state::MovePreparation>,
161 disposition: SourceTargetDisposition,
162 mut source_relay: super::move_session::MoveSourceRelay,
163 ) -> Result<bool> {
164 let retained = source_relay.owner();
165 crate::worker_lifecycle::run_with_owner(
166 session_id,
167 "suspend session for move",
168 executor,
169 retained,
170 async {
171 self.prepare_move_source_checkpoint(
172 session_id,
173 executor,
174 manager,
175 operation,
176 &mut source_relay,
177 )
178 .await?;
179 if disposition == SourceTargetDisposition::RetainForInPlaceSwap
180 || operation.workspace_transfer.is_some()
181 {
182 self.seal_move_handoff(
183 session_id,
184 executor,
185 manager,
186 operation,
187 preparation,
188 source_relay.take(),
189 )
190 .await?;
191 return Ok(false);
192 }
193 self.suspend_session_controlled_with_manager(
194 session_id,
195 executor,
196 Some(manager),
197 Some((operation, preparation)),
198 disposition,
199 CheckpointExportPolicy::ReuseUnchangedArchive,
200 true,
201 None,
202 source_relay.take(),
203 )
204 .await
205 },
206 )
207 .await
208 }
209
210 #[allow(clippy::too_many_arguments)]
211 async fn suspend_session_controlled_with_manager(
212 &mut self,
213 session_id: &str,
214 executor: &(impl CommandExecutor + Sync),
215 manager: Option<&SessionManagerControl>,
216 move_intent: Option<(
217 &mut mj_core::state::MoveOperation,
218 Option<&mj_core::state::MovePreparation>,
219 )>,
220 disposition: SourceTargetDisposition,
221 checkpoint_export_policy: CheckpointExportPolicy,
222 acknowledge_unpublished_work: bool,
223 before_close: Option<BeforeClose>,
224 held_relay: Option<super::checkpoint::ControllerRelayLease>,
225 ) -> Result<bool> {
226 crate::worker_lifecycle::run(session_id, "suspend session controlled with manager", executor, async {
227 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
228 let previous = self
229 .state
230 .sessions
231 .get(session_id)
232 .with_context(|| format!("unknown session {session_id}"))?
233 .clone();
234 let record = self.state.sessions.get_mut(session_id).unwrap();
235 apply_close_checkpoint_started(record, now());
239 self.persist_session_transition_or_restore(
240 session_id,
241 &previous,
242 "persist closing state before checkpointing the session",
243 )?;
244
245 let mut latched = match self
248 .checkpoint_session_latched_for_operation(
249 session_id,
250 executor,
251 manager,
252 LatchExclusivity::HoldThroughClose,
253 checkpoint_export_policy,
254 None,
255 held_relay,
256 )
257 .await
258 {
259 Ok(latched) => latched,
260 Err(error) => {
261 let record = self.state.sessions.get_mut(session_id).unwrap();
262 apply_close_checkpoint_failure(record, &previous, &error, now());
266 return Err(
267 self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
268 );
269 }
270 };
271
272 let artifact = latched.artifact.clone();
273 let publication = if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
276 None
277 } else {
278 match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
279 Some(ManagedCheckoutKind::Clone) => Some(
280 super::publication::assess_clone_checkpoint(&previous, &artifact.metadata),
281 ),
282 Some(ManagedCheckoutKind::Worktree) => None,
283 None if previous.project_directory.is_none() => {
284 Some(super::publication::assess_network_checkpoint(
285 &previous,
286 &artifact.metadata,
287 &self.config,
288 ))
289 }
290 None => None,
291 }
292 };
293 if !acknowledge_unpublished_work
294 && publication
295 .as_ref()
296 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
297 {
298 latched.relay.cancel_abandoned_barrier().await?;
299 let mut restored = previous.clone();
303 restored.state = state_after_unsealed_close(&previous);
304 restored.updated_at = now();
305 self.state.sessions.insert(session_id.to_owned(), restored);
306 self.persist_session_transition_or_restore(
307 session_id,
308 &previous,
309 "restore session after unpublished-work preflight",
310 )?;
311 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
312 return Err(mj_core::refusal::Refusal::precondition(
313 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
314 ).into());
315 }
316 let record = self.state.sessions.get_mut(session_id).unwrap();
317 record.state = SessionState::Closing;
318 record.native_session_id = Some(artifact.native_session_id.clone());
319 record.checkpoint = Some(artifact.metadata.clone());
320 record.publication = publication;
321 record.updated_at = now();
322 record.last_error = None;
323 record.last_checkpoint_error = None;
324 self.persist_checkpoint_transition_or_restore(
325 session_id,
326 &previous,
327 "persist verified checkpoint and closing state before sealing the relay",
328 )?;
329 if let Some((operation, preparation)) = move_intent {
330 if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
333 let record = self.state.sessions.get_mut(session_id).unwrap();
334 record.state = previous.state;
335 record.last_error = Some(format!("{error:#}"));
336 self.persist_session_transition_or_restore(
337 session_id,
338 &previous,
339 "restore source after move preflight failure",
340 )?;
341 return Err(error);
342 }
343 operation.checkpoint = Some(artifact.metadata.clone());
344 operation.updated_at = now();
345 crate::database::save_move_operation(operation)?;
346 }
347 if let Some(before_close) = before_close
348 && let Err(error) = before_close.await
349 {
350 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
353 tracing::warn!(
354 session_id,
355 error = format!("{cancel:#}"),
356 "could not release the checkpoint barrier of a close that did not proceed"
357 );
358 }
359 let record = self.state.sessions.get_mut(session_id).unwrap();
360 record.state = state_after_unsealed_close(&previous);
361 record.last_error = Some(format!("{error:#}"));
362 record.updated_at = now();
363 self.persist_session_transition_or_restore(
364 session_id,
365 &previous,
366 "restore a session whose close did not proceed past its checkpoint",
367 )?;
368 return Err(error);
369 }
370 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
371 release_projection_behind_checkpoint(session_id, &artifact.metadata);
374
375 let close_command_id = new_command_id("close")?;
376 let barrier_command_id = latched.barrier_command_id.clone();
377 let close_result = {
378 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
379 latched
380 .relay
381 .connection_mut()
382 .submit(
383 close_command_id,
384 RelayCommand::Close {
385 barrier_command_id: barrier_command_id.clone(),
386 expected: latched.cursor.clone(),
387 },
388 )
389 .await
390 };
391 if let Err(error) = close_result {
392 let error = error.context("seal verified checkpoint for close");
393 if !close_was_refused(&error) {
394 self.record_interrupted_close(session_id, &error)?;
397 return Err(error);
398 }
399 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
404 tracing::warn!(
405 session_id,
406 error = format!("{cancel:#}"),
407 "could not release the checkpoint barrier of a refused close"
408 );
409 }
410 let record = self.state.sessions.get_mut(session_id).unwrap();
411 record.state = state_after_unsealed_close(&previous);
412 record.last_error = Some(format!("{error:#}"));
413 record.updated_at = now();
414 self.persist_session_transition_or_restore(
415 session_id,
416 &previous,
417 "restore a session whose worker refused its close",
418 )?;
419 return Err(error);
420 }
421 let close_result = {
422 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
423 latched
424 .relay
425 .connection_mut()
426 .submit(
427 new_command_id("checkpoint-complete")?,
428 RelayCommand::CompleteCheckpoint { barrier_command_id },
429 )
430 .await
431 };
432 if let Err(error) = close_result {
433 self.record_interrupted_close(session_id, &error)?;
434 return Err(error.context("release verified close checkpoint"));
435 }
436 let close_result = {
437 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
438 wait_for_relay_closed(latched.relay.connection_mut()).await
439 };
440 if let Err(error) = close_result {
441 self.record_interrupted_close(session_id, &error)?;
442 return Err(error);
443 }
444 latched.relay.release();
445
446 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
447 return Ok(false);
453 }
454 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
455 Ok(deferred) => Ok(deferred),
456 Err(error) => {
457 self.record_interrupted_close(session_id, &error)?;
458 Err(error)
459 }
460 }
461
462 }).await
463 }
464
465 pub async fn recover_interrupted_close_managed(
477 &mut self,
478 session_id: &str,
479 executor: &(impl CommandExecutor + Sync),
480 manager: &SessionManagerControl,
481 acknowledge_unpublished_work: bool,
482 before_close: Option<BeforeClose>,
483 ) -> Result<bool> {
484 crate::worker_lifecycle::run(session_id, "recover interrupted close managed", executor, async {
485 let (state, verified) = {
486 let session = self
487 .state
488 .sessions
489 .get(session_id)
490 .with_context(|| format!("unknown session {session_id}"))?;
491 ensure!(
492 matches!(
493 session.state,
494 SessionState::Closing | SessionState::Destroying
495 ),
496 "session {session_id} has no interrupted close to recover"
497 );
498 (session.state, session.checkpoint.clone())
499 };
500 if state == SessionState::Destroying {
501 let verified = verified.context("destroying session has no verified checkpoint")?;
502 if let Some(before_close) = before_close {
503 before_close.await?;
504 }
505 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
506 }
507 ensure!(
508 state == SessionState::Closing,
509 "session {session_id} has no relay close to recover"
510 );
511 let started = std::time::Instant::now();
512 let mut phase = "wait for session actor";
513 tracing::info!(
514 session_id,
515 ?state,
516 checkpoint_frontier = ?verified.as_ref().map(|checkpoint| checkpoint.event_frontier),
517 "recovering suspension: asking the worker whether its relay was sealed"
518 );
519 let connection = async {
520 let handle = manager
521 .wait_for_session(session_id, Duration::from_secs(5))
522 .await?;
523 phase = "lease worker connection";
524 let mut lease = handle.lease_connection().await?;
525 phase = "read worker execution state";
526 let execution = lease.connection_mut().sync().await?.operational.execution;
527 anyhow::Ok((lease, execution))
528 }
529 .await;
530 let (mut lease, execution) = match connection {
531 Ok(connected) => connected,
532 Err(error) => {
533 tracing::warn!(
534 session_id,
535 phase,
536 elapsed_ms = started.elapsed().as_millis() as u64,
537 error = format!("{error:#}"),
538 "suspension could not query its worker; retaining target and checkpoint"
539 );
540 self.diagnose_suspension_worker(session_id, phase).await;
541 return Err(error.context(format!("{phase} while recovering suspension")));
542 }
543 };
544 tracing::info!(
545 session_id,
546 ?execution,
547 elapsed_ms = started.elapsed().as_millis() as u64,
548 "suspension recovery read the worker's execution state"
549 );
550 match execution {
551 RelayExecutionState::Closed => {}
552 RelayExecutionState::Closing => {
553 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
554 wait_for_relay_closed(lease.connection_mut()).await?;
555 }
556 RelayExecutionState::Idle | RelayExecutionState::Running => {
557 lease.release();
558 return self
559 .suspend_session_controlled_with_manager(
560 session_id,
561 executor,
562 Some(manager),
563 None,
564 SourceTargetDisposition::Destroy,
565 CheckpointExportPolicy::ReuseUnchangedArchive,
566 acknowledge_unpublished_work,
567 before_close,
568 None,
569 )
570 .await;
571 }
572 }
573 lease.release();
574 let verified = verified.context("closed relay has no verified checkpoint")?;
575 if let Some(before_close) = before_close {
577 before_close.await?;
578 }
579 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
580
581 }).await
582 }
583
584 async fn diagnose_suspension_worker(&self, session_id: &str, phase: &str) {
587 let placement = self.worker_placement(session_id);
588 let probe = match placement {
589 Ok((backend, worker_root)) => tokio::task::spawn_blocking(move || {
590 super::worker_binary::probe_worker(
591 &targets::BoundedProcessExecutor::new(Duration::from_secs(5)),
592 &backend,
593 &worker_root,
594 )
595 })
596 .await
597 .context("join the suspension worker diagnostic probe")
598 .and_then(std::convert::identity),
599 Err(error) => Err(error.context("resolve the retained suspension worker")),
600 };
601 match probe {
602 Ok(probe) => tracing::warn!(
603 session_id,
604 phase,
605 worker_pids = ?probe.pids,
606 worker_startup_step = ?probe.step(),
607 worker_state = %probe,
608 "suspension failure worker diagnostic; no recovery mutation was attempted"
609 ),
610 Err(error) => tracing::warn!(
611 session_id,
612 phase,
613 error = format!("{error:#}"),
614 "suspension failure worker diagnostic was unavailable; worker state remains unknown"
615 ),
616 }
617 }
618
619 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
626 self.fail_interrupted_lifecycle_with(
627 session_id,
628 cause,
629 crate::database::save_lifecycle_session,
630 )
631 }
632
633 fn fail_interrupted_lifecycle_with(
634 &mut self,
635 session_id: &str,
636 cause: &str,
637 persist: impl Fn(&SessionRecord) -> Result<()>,
638 ) -> Result<bool> {
639 let Some(record) = self.state.sessions.get_mut(session_id) else {
640 return Ok(false);
641 };
642 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
643 return Ok(false);
644 }
645 let previous = record.clone();
646 record.state = SessionState::Error;
647 record.updated_at = now();
648 record.last_error = Some(cause.to_owned());
649 persist_session_record_transition_or_restore(
650 &mut self.state,
651 session_id,
652 &previous,
653 "persist the failure of an interrupted lifecycle state",
654 &persist,
655 )?;
656 Ok(true)
657 }
658
659 pub fn fail_unready_session(
668 &mut self,
669 session_id: &str,
670 cause: &str,
671 observed_updated_at: &str,
672 ) -> Result<bool> {
673 self.fail_unready_session_with(
674 session_id,
675 cause,
676 observed_updated_at,
677 crate::database::save_lifecycle_session,
678 )
679 }
680
681 fn fail_unready_session_with(
682 &mut self,
683 session_id: &str,
684 cause: &str,
685 observed_updated_at: &str,
686 persist: impl Fn(&SessionRecord) -> Result<()>,
687 ) -> Result<bool> {
688 let Some(record) = self.state.sessions.get_mut(session_id) else {
689 return Ok(false);
690 };
691 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
692 return Ok(false);
693 }
694 let previous = record.clone();
695 record.state = SessionState::Error;
696 record.updated_at = now();
697 record.last_error = Some(cause.to_owned());
698 persist_session_record_transition_or_restore(
699 &mut self.state,
700 session_id,
701 &previous,
702 "persist a session whose harness never became usable",
703 &persist,
704 )?;
705 Ok(true)
706 }
707
708 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
712 let Some(record) = self.state.sessions.get(session_id) else {
713 return Ok(false);
714 };
715 let previous = record.clone();
718 let record = self.state.sessions.get_mut(session_id).unwrap();
719 record.last_error = Some(cause.to_owned());
720 record.updated_at = now();
721 persist_session_record_transition_or_restore(
722 &mut self.state,
723 session_id,
724 &previous,
725 "persist the reason a close did not finish",
726 &crate::database::save_lifecycle_session,
727 )?;
728 Ok(true)
729 }
730
731 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
736 let Some(record) = self.state.sessions.get(session_id) else {
737 return Ok(false);
738 };
739 if record.public_error().is_none() {
740 return Ok(false);
741 }
742 let previous = record.clone();
743 let record = self.state.sessions.get_mut(session_id).unwrap();
744 record.last_error = None;
745 record.updated_at = now();
746 persist_session_record_transition_or_restore(
747 &mut self.state,
748 session_id,
749 &previous,
750 "clear the reason a close did not finish",
751 &crate::database::save_lifecycle_session,
752 )?;
753 Ok(true)
754 }
755
756 pub fn suspend_session_without_checkpoint(
769 &mut self,
770 session_id: &str,
771 executor: &impl CommandExecutor,
772 ) -> Result<bool> {
773 self.suspend_session_without_checkpoint_with(
774 session_id,
775 executor,
776 crate::database::save_lifecycle_session,
777 )
778 }
779
780 fn suspend_session_without_checkpoint_with(
781 &mut self,
782 session_id: &str,
783 executor: &impl CommandExecutor,
784 persist: impl Fn(&SessionRecord) -> Result<()>,
785 ) -> Result<bool> {
786 let session = self
787 .state
788 .sessions
789 .get(session_id)
790 .with_context(|| format!("unknown session {session_id}"))?
791 .clone();
792 ensure!(
793 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
794 "session {session_id} has a workspace to checkpoint; suspend it instead"
795 );
796 self.stop_target_and_settle(session_id, &session, executor, &persist)
797 }
798
799 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
800 let record = self.state.sessions.get_mut(session_id).unwrap();
801 apply_interrupted_close_error(record, error, &now());
802 self.persist_session_state(session_id)
803 }
804
805 fn destroy_after_verified_checkpoint(
808 &mut self,
809 session_id: &str,
810 verified: &CheckpointMetadata,
811 executor: &impl CommandExecutor,
812 ) -> Result<bool> {
813 self.destroy_after_verified_checkpoint_with(
814 session_id,
815 verified,
816 executor,
817 crate::database::save_lifecycle_session,
818 )
819 }
820
821 fn destroy_after_verified_checkpoint_with(
822 &mut self,
823 session_id: &str,
824 verified: &CheckpointMetadata,
825 executor: &impl CommandExecutor,
826 persist: impl Fn(&SessionRecord) -> Result<()>,
827 ) -> Result<bool> {
828 crate::worker_lifecycle::run_blocking(
829 session_id,
830 "destroy after verified checkpoint with",
831 executor,
832 || {
833 let session = self
834 .state
835 .sessions
836 .get(session_id)
837 .with_context(|| format!("unknown session {session_id}"))?
838 .clone();
839 ensure!(
840 matches!(
841 session.state,
842 SessionState::Closing | SessionState::Destroying
843 ),
844 "refusing to destroy session {session_id}: it is not closing or destroying"
845 );
846 ensure!(
847 session.checkpoint.as_ref() == Some(verified),
848 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
849 );
850 if session.state == SessionState::Closing {
851 let record = self.state.sessions.get_mut(session_id).unwrap();
852 record.state = SessionState::Destroying;
853 record.updated_at = now();
854 record.last_error = None;
855 persist_session_record_transition_or_restore(
856 &mut self.state,
857 session_id,
858 &session,
859 "persist destroying state before target cleanup",
860 &persist,
861 )?;
862 }
863
864 let destroying = self
865 .state
866 .sessions
867 .get(session_id)
868 .expect("destroying session disappeared")
869 .clone();
870 {
871 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
872 verify_installed_checkpoint_gate(session_id, verified)?;
873 }
874 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
879 tracing::warn!(
880 session_id,
881 error = format!("{error:#}"),
882 "could not record that the second-opinion conversation ends with this target"
883 );
884 }
885 let locator = destroying
886 .target
887 .as_ref()
888 .context("session has no target")?;
889 let backend = backend_locator(locator, &destroying, &self.config)?;
890 let deferred = if self.state.subagents.contains_key(session_id) {
891 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
892 .execute(executor)?;
893 false
894 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
895 plan.execute(executor)?;
896 true
897 } else {
898 execute_target_cleanup(&backend, &destroying, &self.config, executor)?;
899 false
900 };
901 if let Some(worktree) = &destroying.managed_worktree {
902 retire_managed_worktree(executor, worktree)
903 .context("retire managed raw-session worktree after verified close")?;
904 }
905 let record = self.state.sessions.get_mut(session_id).unwrap();
906 record.state = SessionState::Stopped;
907 if !deferred {
908 record.target = None;
909 }
910 record.updated_at = now();
911 record.last_error = None;
912 persist_session_record_transition_or_restore(
913 &mut self.state,
914 session_id,
915 &destroying,
916 "persist stopped state after target cleanup",
917 &persist,
918 )?;
919 Ok(deferred)
920 },
921 )
922 }
923
924 pub fn cleanup_stopped_target(
928 &mut self,
929 session_id: &str,
930 executor: &impl CommandExecutor,
931 ) -> Result<()> {
932 self.cleanup_stopped_target_with(
933 session_id,
934 executor,
935 crate::database::save_lifecycle_session,
936 )
937 }
938
939 fn cleanup_stopped_target_with(
940 &mut self,
941 session_id: &str,
942 executor: &impl CommandExecutor,
943 persist: impl Fn(&SessionRecord) -> Result<()>,
944 ) -> Result<()> {
945 crate::worker_lifecycle::run_blocking(
946 session_id,
947 "cleanup stopped target with",
948 executor,
949 || {
950 let previous = self
951 .state
952 .sessions
953 .get(session_id)
954 .with_context(|| format!("unknown session {session_id}"))?
955 .clone();
956 ensure!(
957 previous.state == SessionState::Stopped,
958 "refusing deferred cleanup for active session {session_id}"
959 );
960 let Some(locator) = previous.target.as_ref() else {
961 return Ok(());
962 };
963 let backend = backend_locator(locator, &previous, &self.config)?;
964 ensure!(
965 targets::quiesce_plan(&backend, session_id)?.is_some(),
966 "session {session_id} retained a non-Podman target after stopping"
967 );
968 if let Err(error) =
969 execute_target_cleanup(&backend, &previous, &self.config, executor)
970 {
971 let record = self.state.sessions.get_mut(session_id).unwrap();
972 record.updated_at = now();
973 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
974 let persisted = persist_session_record_transition_or_restore(
975 &mut self.state,
976 session_id,
977 &previous,
978 "persist deferred target cleanup failure",
979 &persist,
980 );
981 return match persisted {
982 Ok(()) => Err(error),
983 Err(persist_error) => Err(error.context(format!(
984 "also failed to persist deferred target cleanup failure: {persist_error:#}"
985 ))),
986 };
987 }
988 let record = self.state.sessions.get_mut(session_id).unwrap();
989 record.target = None;
990 record.updated_at = now();
991 record.last_error = None;
992 persist_session_record_transition_or_restore(
993 &mut self.state,
994 session_id,
995 &previous,
996 "persist completion of deferred Podman target cleanup",
997 &persist,
998 )
999 },
1000 )
1001 }
1002
1003 pub fn force_stop(
1006 &mut self,
1007 session_id: &str,
1008 executor: &impl CommandExecutor,
1009 ) -> Result<bool> {
1010 self.force_stop_with(
1011 session_id,
1012 executor,
1013 crate::database::save_lifecycle_session,
1014 )
1015 }
1016
1017 fn force_stop_with(
1018 &mut self,
1019 session_id: &str,
1020 executor: &impl CommandExecutor,
1021 persist: impl Fn(&SessionRecord) -> Result<()>,
1022 ) -> Result<bool> {
1023 let session = self
1024 .state
1025 .sessions
1026 .get(session_id)
1027 .with_context(|| format!("unknown session {session_id}"))?
1028 .clone();
1029 ensure!(
1030 session.state.is_active(),
1031 "session {session_id} is already inactive"
1032 );
1033 let checkpoint = session
1034 .checkpoint
1035 .as_ref()
1036 .context("force stop requires an existing recovery archive")?;
1037 {
1040 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1041 verify_installed_checkpoint_gate(session_id, checkpoint)
1042 .context("verify the recovery archive before force stopping")?;
1043 }
1044 self.stop_target_and_settle(session_id, &session, executor, &persist)
1045 }
1046
1047 fn stop_target_and_settle(
1052 &mut self,
1053 session_id: &str,
1054 session: &SessionRecord,
1055 executor: &impl CommandExecutor,
1056 persist: &impl Fn(&SessionRecord) -> Result<()>,
1057 ) -> Result<bool> {
1058 crate::worker_lifecycle::run_blocking(
1059 session_id,
1060 "stop target and settle",
1061 executor,
1062 || {
1063 let mut deferred = false;
1064 if let Some(locator) = &session.target {
1065 let backend = backend_locator(locator, session, &self.config)?;
1066 if self.state.subagents.contains_key(session_id) {
1070 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1071 .execute(executor)?;
1072 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1073 plan.execute(executor)?;
1074 deferred = true;
1075 } else {
1076 execute_target_cleanup(&backend, session, &self.config, executor)?;
1077 }
1078 }
1079 if let Some(worktree) = &session.managed_worktree {
1080 retire_managed_worktree(executor, worktree)
1081 .context("retire managed raw-session worktree after stopping the target")?;
1082 }
1083 let record = self.state.sessions.get_mut(session_id).unwrap();
1084 record.state = SessionState::Stopped;
1085 if !deferred {
1086 record.target = None;
1087 }
1088 record.updated_at = now();
1089 record.last_error = None;
1090 record.last_checkpoint_error = None;
1091 persist_session_record_transition_or_restore(
1092 &mut self.state,
1093 session_id,
1094 session,
1095 "persist stopped state after tearing down the current target",
1096 persist,
1097 )?;
1098 Ok(deferred)
1099 },
1100 )
1101 }
1102
1103 pub fn destroy_session_controlled(
1107 &mut self,
1108 session_id: &str,
1109 executor: &impl CommandExecutor,
1110 ) -> Result<()> {
1111 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1112 }
1113
1114 pub fn destroy_session_controlled_with(
1122 &mut self,
1123 session_id: &str,
1124 executor: &impl CommandExecutor,
1125 branch: BranchDisposition,
1126 ) -> Result<()> {
1127 self.destroy_session_controlled_with_checkout(
1128 session_id,
1129 executor,
1130 branch,
1131 CheckoutDisposition::Remove,
1132 )
1133 .map(|_| ())
1134 }
1135
1136 pub(crate) fn destroy_session_controlled_with_checkout(
1142 &mut self,
1143 session_id: &str,
1144 executor: &impl CommandExecutor,
1145 branch: BranchDisposition,
1146 checkout: CheckoutDisposition,
1147 ) -> Result<Option<PathBuf>> {
1148 let session = self
1149 .state
1150 .sessions
1151 .get(session_id)
1152 .with_context(|| format!("unknown session {session_id}"))?
1153 .clone();
1154 if session.state.is_active() {
1155 bail!("refusing to destroy active session {session_id}");
1156 }
1157 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1158 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1159 }
1160 let mut retained_checkout = None;
1161 if let Some(worktree) = &session.managed_worktree {
1162 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1163 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1164 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1165 "check the managed raw-session worktree for uncommitted changes",
1166 )?);
1167 if keep {
1168 retained_checkout = Some(worktree.worktree_root.clone());
1172 } else {
1173 cleanup_managed_worktree(executor, worktree, branch)
1174 .context("remove managed raw-session worktree")?;
1175 }
1176 }
1177 if let Some(checkpoint) = &session.checkpoint
1178 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1179 && error.kind() != std::io::ErrorKind::NotFound
1180 {
1181 return Err(error).with_context(|| {
1182 format!(
1183 "remove session recovery archive {}",
1184 checkpoint.archive_path.display()
1185 )
1186 });
1187 }
1188 mj_core::attachment::AttachmentStore::controller(session_id)?
1189 .remove_session_data()
1190 .context("remove session image attachments")?;
1191 crate::database::delete_session(session_id)
1192 .context("destroy stopped session in database")?;
1193 self.state.subagents.remove(session_id);
1194 self.state.destroy_stopped_session(session_id)?;
1195 Ok(retained_checkout)
1196 }
1197
1198 pub fn force_destroy_session(
1210 &mut self,
1211 session_id: &str,
1212 executor: &impl CommandExecutor,
1213 branch: BranchDisposition,
1214 ) -> Result<()> {
1215 self.force_destroy_session_with(
1216 session_id,
1217 executor,
1218 branch,
1219 crate::database::delete_session,
1220 )
1221 }
1222
1223 fn force_destroy_session_with(
1224 &mut self,
1225 session_id: &str,
1226 executor: &impl CommandExecutor,
1227 branch: BranchDisposition,
1228 delete: impl Fn(&str) -> Result<()>,
1229 ) -> Result<()> {
1230 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1231 if result.is_err() {
1232 self.settle_failed_destroy(session_id);
1233 }
1234 result
1235 }
1236
1237 fn settle_failed_destroy(&mut self, session_id: &str) {
1243 let Some(record) = self.state.sessions.get(session_id) else {
1244 return;
1245 };
1246 if !record.state.is_active() || record.state == SessionState::Error {
1247 return;
1248 }
1249 let previous = record.clone();
1250 let record = self.state.sessions.get_mut(session_id).unwrap();
1251 record.state = SessionState::Error;
1252 record.updated_at = now();
1253 if let Err(error) = persist_session_record_transition_or_restore(
1254 &mut self.state,
1255 session_id,
1256 &previous,
1257 "persist the failed destroy",
1258 &crate::database::save_lifecycle_session,
1259 ) {
1260 tracing::warn!(
1261 session_id,
1262 error = format!("{error:#}"),
1263 "could not record a failed destroy"
1264 );
1265 }
1266 }
1267
1268 fn force_destroy_session_steps(
1269 &mut self,
1270 session_id: &str,
1271 executor: &impl CommandExecutor,
1272 branch: BranchDisposition,
1273 delete: impl Fn(&str) -> Result<()>,
1274 ) -> Result<()> {
1275 crate::worker_lifecycle::run_blocking(
1276 session_id,
1277 "force destroy session steps",
1278 executor,
1279 || {
1280 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1281 let session = self
1282 .state
1283 .sessions
1284 .get(session_id)
1285 .with_context(|| format!("unknown session {session_id}"))?
1286 .clone();
1287 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1288 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1289 }
1290 if let Some(locator) = &session.target {
1294 let backend = backend_locator(locator, &session, &self.config)?;
1295 if self.state.subagents.contains_key(session_id) {
1296 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1297 .execute(executor)?;
1298 } else {
1299 execute_target_cleanup(&backend, &session, &self.config, executor)?;
1300 }
1301 }
1302 if let Some(worktree) = &session.managed_worktree {
1303 cleanup_managed_worktree(executor, worktree, branch)
1304 .context("remove managed raw-session worktree")?;
1305 }
1306 if let Some(checkpoint) = &session.checkpoint
1307 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1308 && error.kind() != std::io::ErrorKind::NotFound
1309 {
1310 return Err(error).with_context(|| {
1311 format!(
1312 "remove session recovery archive {}",
1313 checkpoint.archive_path.display()
1314 )
1315 });
1316 }
1317 mj_core::attachment::AttachmentStore::controller(session_id)?
1318 .remove_session_data()
1319 .context("remove session image attachments")?;
1320 delete(session_id).context("force destroy session in database")?;
1321 self.state.subagents.remove(session_id);
1322 self.state.destroy_session_force(session_id)?;
1323 Ok(())
1324 },
1325 )
1326 }
1327}
1328
1329pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1344 match session.state {
1345 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1346 SessionState::Error if subagent => true,
1347 SessionState::Closing
1348 | SessionState::Destroying
1349 | SessionState::Error
1350 | SessionState::Lost => session.target.is_none(),
1351 SessionState::Running
1352 | SessionState::Disconnected
1353 | SessionState::Checkpointing
1354 | SessionState::Stopped
1355 | SessionState::DestroyedWithDataLoss => false,
1356 }
1357}
1358
1359fn execute_target_cleanup(
1363 backend: &targets::TargetLocator,
1364 session: &SessionRecord,
1365 config: &mj_core::config::Config,
1366 executor: &impl CommandExecutor,
1367) -> Result<()> {
1368 let session_id = session.id.as_str();
1369 let _owner = crate::worker_lifecycle::require(session_id)?;
1370 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1371 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1372 Ok(true) => {
1373 tracing::warn!(
1374 session_id,
1375 error = format!("{cleanup_error:#}"),
1376 "target cleanup command failed, but the target was confirmed absent"
1377 );
1378 }
1379 Ok(false) => {
1380 tracing::error!(
1381 session_id,
1382 error = format!("{cleanup_error:#}"),
1383 "target cleanup failed and the target is still present"
1384 );
1385 return Err(cleanup_error);
1386 }
1387 Err(probe_error) => {
1388 tracing::error!(
1389 session_id,
1390 cleanup_error = format!("{cleanup_error:#}"),
1391 probe_error = format!("{probe_error:#}"),
1392 "target cleanup failed and exact absence could not be confirmed"
1393 );
1394 return Err(cleanup_error.context(format!(
1395 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1396 )));
1397 }
1398 }
1399 }
1400 if let Some(release) = BuildStateRelease::for_target(session, backend, config) {
1401 release.run(executor);
1402 }
1403 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
1406 crate::database::finish_worker_restart(session_id, &intent.operation_id)?;
1407 }
1408 Ok(())
1409}
1410
1411fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1412 record.state = SessionState::Closing;
1413 record.updated_at = updated_at;
1414 record.last_checkpoint_error = None;
1415}
1416
1417fn apply_close_checkpoint_failure(
1424 record: &mut SessionRecord,
1425 previous: &SessionRecord,
1426 error: &anyhow::Error,
1427 updated_at: String,
1428) {
1429 if WorkerRestartLeftNoWorker::marks(error) {
1430 record.state = SessionState::Error;
1431 record.last_error = Some(format!(
1432 "close failed and left the session without a live worker; retry the close, \
1433 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1434 ));
1435 } else {
1436 record.state = state_after_unsealed_close(previous);
1437 }
1438 record.last_checkpoint_error = Some(format!("{error:#}"));
1439 record.updated_at = updated_at;
1440}
1441
1442fn close_was_refused(error: &anyhow::Error) -> bool {
1450 error.chain().any(|cause| {
1451 cause
1452 .downcast_ref::<crate::worker_client::RelayRejected>()
1453 .is_some_and(|rejected| !rejected.is_retryable())
1454 })
1455}
1456
1457fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1458 if previous.state == SessionState::Closing {
1459 SessionState::Running
1460 } else {
1461 previous.state
1462 }
1463}
1464
1465fn apply_interrupted_close_error(
1466 record: &mut SessionRecord,
1467 error: &anyhow::Error,
1468 updated_at: &str,
1469) {
1470 let destroying = record.state == SessionState::Destroying;
1471 if !destroying {
1472 record.state = SessionState::Closing;
1473 }
1474 record.updated_at = updated_at.to_owned();
1475 record.last_error = Some(if destroying {
1476 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1477 } else {
1478 format!("close is safely resumable from its verified checkpoint: {error:#}")
1479 });
1480}
1481
1482#[cfg(test)]
1483mod tests;