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 let checkout = self.state.checkout(session_id)?;
279 match &checkout {
280 mj_core::state::Checkout::ManagedWorktree { worktree, .. }
281 if worktree.kind == ManagedCheckoutKind::Clone =>
282 {
283 Some(super::publication::assess_clone_checkpoint_for_worktree(
284 worktree,
285 &artifact.metadata,
286 ))
287 }
288 mj_core::state::Checkout::ManagedWorkspace => {
289 Some(super::publication::assess_network_checkpoint(
290 &previous,
291 &artifact.metadata,
292 &self.config,
293 ))
294 }
295 mj_core::state::Checkout::Borrowed { .. }
296 if checkout.project_directory().is_none() =>
297 {
298 Some(super::publication::assess_network_checkpoint(
299 &previous,
300 &artifact.metadata,
301 &self.config,
302 ))
303 }
304 mj_core::state::Checkout::Attached { .. }
305 | mj_core::state::Checkout::ManagedWorktree { .. }
306 | mj_core::state::Checkout::Borrowed { .. } => None,
307 }
308 };
309 if !acknowledge_unpublished_work
310 && publication
311 .as_ref()
312 .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
313 {
314 latched.relay.cancel_abandoned_barrier().await?;
315 let mut restored = previous.clone();
319 restored.state = state_after_unsealed_close(&previous);
320 restored.updated_at = now();
321 self.state.sessions.insert(session_id.to_owned(), restored);
322 self.persist_session_transition_or_restore(
323 session_id,
324 &previous,
325 "restore session after unpublished-work preflight",
326 )?;
327 let _ = std::fs::remove_file(&artifact.metadata.archive_path);
328 return Err(mj_core::refusal::Refusal::precondition(
329 "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
330 ).into());
331 }
332 let record = self.state.sessions.get_mut(session_id).unwrap();
333 record.state = SessionState::Closing;
334 record.native_session_id = Some(artifact.native_session_id.clone());
335 record.checkpoint = Some(artifact.metadata.clone());
336 record.publication = publication;
337 record.updated_at = now();
338 record.last_error = None;
339 record.last_checkpoint_error = None;
340 self.persist_checkpoint_transition_or_restore(
341 session_id,
342 &previous,
343 "persist verified checkpoint and closing state before sealing the relay",
344 )?;
345 if let Some((operation, preparation)) = move_intent {
346 let validation = match self
349 .github_token_for_repository_preflight(session_id)
350 .await
351 {
352 Ok(token) => self.validate_move_checkpoint(
353 operation,
354 preparation,
355 token.as_deref(),
356 executor,
357 ),
358 Err(error) => Err(error),
359 };
360 if let Err(error) = validation {
361 let record = self.state.sessions.get_mut(session_id).unwrap();
362 record.state = previous.state;
363 record.last_error = Some(format!("{error:#}"));
364 self.persist_session_transition_or_restore(
365 session_id,
366 &previous,
367 "restore source after move preflight failure",
368 )?;
369 return Err(error);
370 }
371 operation.checkpoint = Some(artifact.metadata.clone());
372 operation.updated_at = now();
373 crate::database::save_move_operation(operation)?;
374 }
375 if let Some(before_close) = before_close
376 && let Err(error) = before_close.await
377 {
378 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
381 tracing::warn!(
382 session_id,
383 error = format!("{cancel:#}"),
384 "could not release the checkpoint barrier of a close that did not proceed"
385 );
386 }
387 let record = self.state.sessions.get_mut(session_id).unwrap();
388 record.state = state_after_unsealed_close(&previous);
389 record.last_error = Some(format!("{error:#}"));
390 record.updated_at = now();
391 self.persist_session_transition_or_restore(
392 session_id,
393 &previous,
394 "restore a session whose close did not proceed past its checkpoint",
395 )?;
396 return Err(error);
397 }
398 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
399 release_projection_behind_checkpoint(session_id, &artifact.metadata);
402
403 let close_command_id = new_command_id("close")?;
404 let barrier_command_id = latched.barrier_command_id.clone();
405 let close_result = {
406 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
407 latched
408 .relay
409 .connection_mut()
410 .submit(
411 close_command_id,
412 RelayCommand::Close {
413 barrier_command_id: barrier_command_id.clone(),
414 expected: latched.cursor.clone(),
415 },
416 )
417 .await
418 };
419 if let Err(error) = close_result {
420 let error = error.context("seal verified checkpoint for close");
421 if !close_was_refused(&error) {
422 self.record_interrupted_close(session_id, &error)?;
425 return Err(error);
426 }
427 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
432 tracing::warn!(
433 session_id,
434 error = format!("{cancel:#}"),
435 "could not release the checkpoint barrier of a refused close"
436 );
437 }
438 let record = self.state.sessions.get_mut(session_id).unwrap();
439 record.state = state_after_unsealed_close(&previous);
440 record.last_error = Some(format!("{error:#}"));
441 record.updated_at = now();
442 self.persist_session_transition_or_restore(
443 session_id,
444 &previous,
445 "restore a session whose worker refused its close",
446 )?;
447 return Err(error);
448 }
449 let close_result = {
450 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
451 latched
452 .relay
453 .connection_mut()
454 .submit(
455 new_command_id("checkpoint-complete")?,
456 RelayCommand::CompleteCheckpoint { barrier_command_id },
457 )
458 .await
459 };
460 if let Err(error) = close_result {
461 self.record_interrupted_close(session_id, &error)?;
462 return Err(error.context("release verified close checkpoint"));
463 }
464 let close_result = {
465 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
466 wait_for_relay_closed(latched.relay.connection_mut()).await
467 };
468 if let Err(error) = close_result {
469 self.record_interrupted_close(session_id, &error)?;
470 return Err(error);
471 }
472 latched.relay.release();
473
474 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
475 return Ok(false);
481 }
482 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
483 Ok(deferred) => Ok(deferred),
484 Err(error) => {
485 self.record_interrupted_close(session_id, &error)?;
486 Err(error)
487 }
488 }
489
490 }).await
491 }
492
493 pub async fn recover_interrupted_close_managed(
505 &mut self,
506 session_id: &str,
507 executor: &(impl CommandExecutor + Sync),
508 manager: &SessionManagerControl,
509 acknowledge_unpublished_work: bool,
510 before_close: Option<BeforeClose>,
511 ) -> Result<bool> {
512 crate::worker_lifecycle::run(session_id, "recover interrupted close managed", executor, async {
513 let (state, verified) = {
514 let session = self
515 .state
516 .sessions
517 .get(session_id)
518 .with_context(|| format!("unknown session {session_id}"))?;
519 ensure!(
520 matches!(
521 session.state,
522 SessionState::Closing | SessionState::Destroying
523 ),
524 "session {session_id} has no interrupted close to recover"
525 );
526 (session.state, session.checkpoint.clone())
527 };
528 if state == SessionState::Destroying {
529 let verified = verified.context("destroying session has no verified checkpoint")?;
530 if let Some(before_close) = before_close {
531 before_close.await?;
532 }
533 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
534 }
535 ensure!(
536 state == SessionState::Closing,
537 "session {session_id} has no relay close to recover"
538 );
539 let started = std::time::Instant::now();
540 let mut phase = "wait for session actor";
541 tracing::info!(
542 session_id,
543 ?state,
544 checkpoint_frontier = ?verified.as_ref().map(|checkpoint| checkpoint.event_frontier),
545 "recovering suspension: asking the worker whether its relay was sealed"
546 );
547 let connection = async {
548 let handle = manager
549 .wait_for_session(session_id, Duration::from_secs(5))
550 .await?;
551 phase = "lease worker connection";
552 let mut lease = handle.lease_connection().await?;
553 phase = "read worker execution state";
554 let execution = lease.connection_mut().sync().await?.operational.execution;
555 anyhow::Ok((lease, execution))
556 }
557 .await;
558 let (mut lease, execution) = match connection {
559 Ok(connected) => connected,
560 Err(error) => {
561 tracing::warn!(
562 session_id,
563 phase,
564 elapsed_ms = started.elapsed().as_millis() as u64,
565 error = format!("{error:#}"),
566 "suspension could not query its worker; retaining target and checkpoint"
567 );
568 self.diagnose_suspension_worker(session_id, phase).await;
569 return Err(error.context(format!("{phase} while recovering suspension")));
570 }
571 };
572 tracing::info!(
573 session_id,
574 ?execution,
575 elapsed_ms = started.elapsed().as_millis() as u64,
576 "suspension recovery read the worker's execution state"
577 );
578 match execution {
579 RelayExecutionState::Closed => {}
580 RelayExecutionState::Closing => {
581 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
582 wait_for_relay_closed(lease.connection_mut()).await?;
583 }
584 RelayExecutionState::Idle | RelayExecutionState::Running => {
585 lease.release();
586 return self
587 .suspend_session_controlled_with_manager(
588 session_id,
589 executor,
590 Some(manager),
591 None,
592 SourceTargetDisposition::Destroy,
593 CheckpointExportPolicy::ReuseUnchangedArchive,
594 acknowledge_unpublished_work,
595 before_close,
596 None,
597 )
598 .await;
599 }
600 }
601 lease.release();
602 let verified = verified.context("closed relay has no verified checkpoint")?;
603 if let Some(before_close) = before_close {
605 before_close.await?;
606 }
607 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
608
609 }).await
610 }
611
612 async fn diagnose_suspension_worker(&self, session_id: &str, phase: &str) {
615 let placement = self.worker_placement(session_id);
616 let probe = match placement {
617 Ok((backend, worker_root)) => tokio::task::spawn_blocking(move || {
618 super::worker_binary::probe_worker(
619 &targets::BoundedProcessExecutor::new(Duration::from_secs(5)),
620 &backend,
621 &worker_root,
622 )
623 })
624 .await
625 .context("join the suspension worker diagnostic probe")
626 .and_then(std::convert::identity),
627 Err(error) => Err(error.context("resolve the retained suspension worker")),
628 };
629 match probe {
630 Ok(probe) => tracing::warn!(
631 session_id,
632 phase,
633 worker_pids = ?probe.pids,
634 worker_startup_step = ?probe.step(),
635 worker_state = %probe,
636 "suspension failure worker diagnostic; no recovery mutation was attempted"
637 ),
638 Err(error) => tracing::warn!(
639 session_id,
640 phase,
641 error = format!("{error:#}"),
642 "suspension failure worker diagnostic was unavailable; worker state remains unknown"
643 ),
644 }
645 }
646
647 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
654 self.fail_interrupted_lifecycle_with(
655 session_id,
656 cause,
657 crate::database::save_lifecycle_session,
658 )
659 }
660
661 fn fail_interrupted_lifecycle_with(
662 &mut self,
663 session_id: &str,
664 cause: &str,
665 persist: impl Fn(&SessionRecord) -> Result<()>,
666 ) -> Result<bool> {
667 let Some(record) = self.state.sessions.get_mut(session_id) else {
668 return Ok(false);
669 };
670 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
671 return Ok(false);
672 }
673 let previous = record.clone();
674 record.state = SessionState::Error;
675 record.updated_at = now();
676 record.last_error = Some(cause.to_owned());
677 persist_session_record_transition_or_restore(
678 &mut self.state,
679 session_id,
680 &previous,
681 "persist the failure of an interrupted lifecycle state",
682 &persist,
683 )?;
684 Ok(true)
685 }
686
687 pub fn fail_unready_session(
696 &mut self,
697 session_id: &str,
698 cause: &str,
699 observed_updated_at: &str,
700 ) -> Result<bool> {
701 self.fail_unready_session_with(
702 session_id,
703 cause,
704 observed_updated_at,
705 crate::database::save_lifecycle_session,
706 )
707 }
708
709 pub(crate) async fn persist_harness_preparation_failure(
713 &self,
714 session_id: &str,
715 cause: &str,
716 observed_updated_at: &str,
717 ) -> Result<()> {
718 let session_id = session_id.to_owned();
719 let cause = cause.to_owned();
720 let observed_updated_at = observed_updated_at.to_owned();
721 tokio::task::spawn_blocking(move || {
722 let mut controller = Controller::load()?;
723 controller.fail_unready_session(&session_id, &cause, &observed_updated_at)
724 })
725 .await
726 .context("record harness preparation failure task panicked")??;
727 Ok(())
728 }
729
730 fn fail_unready_session_with(
731 &mut self,
732 session_id: &str,
733 cause: &str,
734 observed_updated_at: &str,
735 persist: impl Fn(&SessionRecord) -> Result<()>,
736 ) -> Result<bool> {
737 let Some(record) = self.state.sessions.get_mut(session_id) else {
738 return Ok(false);
739 };
740 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
741 return Ok(false);
742 }
743 let previous = record.clone();
744 record.state = SessionState::Error;
745 record.updated_at = now();
746 record.last_error = Some(cause.to_owned());
747 persist_session_record_transition_or_restore(
748 &mut self.state,
749 session_id,
750 &previous,
751 "persist a session whose harness never became usable",
752 &persist,
753 )?;
754 Ok(true)
755 }
756
757 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
761 let Some(record) = self.state.sessions.get(session_id) else {
762 return Ok(false);
763 };
764 let previous = record.clone();
767 let record = self.state.sessions.get_mut(session_id).unwrap();
768 record.last_error = Some(cause.to_owned());
769 record.updated_at = now();
770 persist_session_record_transition_or_restore(
771 &mut self.state,
772 session_id,
773 &previous,
774 "persist the reason a close did not finish",
775 &crate::database::save_lifecycle_session,
776 )?;
777 Ok(true)
778 }
779
780 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
785 let Some(record) = self.state.sessions.get(session_id) else {
786 return Ok(false);
787 };
788 if record.public_error().is_none() {
789 return Ok(false);
790 }
791 let previous = record.clone();
792 let record = self.state.sessions.get_mut(session_id).unwrap();
793 record.last_error = None;
794 record.updated_at = now();
795 persist_session_record_transition_or_restore(
796 &mut self.state,
797 session_id,
798 &previous,
799 "clear the reason a close did not finish",
800 &crate::database::save_lifecycle_session,
801 )?;
802 Ok(true)
803 }
804
805 pub fn suspend_session_without_checkpoint(
818 &mut self,
819 session_id: &str,
820 executor: &impl CommandExecutor,
821 ) -> Result<bool> {
822 self.suspend_session_without_checkpoint_with(
823 session_id,
824 executor,
825 crate::database::save_lifecycle_session,
826 )
827 }
828
829 fn suspend_session_without_checkpoint_with(
830 &mut self,
831 session_id: &str,
832 executor: &impl CommandExecutor,
833 persist: impl Fn(&SessionRecord) -> Result<()>,
834 ) -> Result<bool> {
835 let session = self
836 .state
837 .sessions
838 .get(session_id)
839 .with_context(|| format!("unknown session {session_id}"))?
840 .clone();
841 ensure!(
842 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
843 "session {session_id} has a workspace to checkpoint; suspend it instead"
844 );
845 self.stop_target_and_settle(session_id, &session, executor, &persist)
846 }
847
848 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
849 let record = self.state.sessions.get_mut(session_id).unwrap();
850 apply_interrupted_close_error(record, error, &now());
851 self.persist_session_state(session_id)
852 }
853
854 fn destroy_after_verified_checkpoint(
857 &mut self,
858 session_id: &str,
859 verified: &CheckpointMetadata,
860 executor: &impl CommandExecutor,
861 ) -> Result<bool> {
862 self.destroy_after_verified_checkpoint_with(
863 session_id,
864 verified,
865 executor,
866 crate::database::save_lifecycle_session,
867 )
868 }
869
870 fn destroy_after_verified_checkpoint_with(
871 &mut self,
872 session_id: &str,
873 verified: &CheckpointMetadata,
874 executor: &impl CommandExecutor,
875 persist: impl Fn(&SessionRecord) -> Result<()>,
876 ) -> Result<bool> {
877 crate::worker_lifecycle::run_blocking(
878 session_id,
879 "destroy after verified checkpoint with",
880 executor,
881 || {
882 let session = self
883 .state
884 .sessions
885 .get(session_id)
886 .with_context(|| format!("unknown session {session_id}"))?
887 .clone();
888 ensure!(
889 matches!(
890 session.state,
891 SessionState::Closing | SessionState::Destroying
892 ),
893 "refusing to destroy session {session_id}: it is not closing or destroying"
894 );
895 ensure!(
896 session.checkpoint.as_ref() == Some(verified),
897 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
898 );
899 if session.state == SessionState::Closing {
900 let record = self.state.sessions.get_mut(session_id).unwrap();
901 record.state = SessionState::Destroying;
902 record.updated_at = now();
903 record.last_error = None;
904 persist_session_record_transition_or_restore(
905 &mut self.state,
906 session_id,
907 &session,
908 "persist destroying state before target cleanup",
909 &persist,
910 )?;
911 }
912
913 let destroying = self
914 .state
915 .sessions
916 .get(session_id)
917 .expect("destroying session disappeared")
918 .clone();
919 {
920 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
921 verify_installed_checkpoint_gate(session_id, verified)?;
922 }
923 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
928 tracing::warn!(
929 session_id,
930 error = format!("{error:#}"),
931 "could not record that the second-opinion conversation ends with this target"
932 );
933 }
934 let locator = destroying
935 .target
936 .as_ref()
937 .context("session has no target")?;
938 let backend = backend_locator(locator, &destroying, &self.config)?;
939 let deferred = if self.state.subagents.contains_key(session_id) {
940 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
941 .execute(executor)?;
942 false
943 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
944 plan.execute(executor)?;
945 true
946 } else {
947 execute_target_cleanup(&backend, &destroying, &self.config, executor)?;
948 false
949 };
950 if let mj_core::state::Checkout::ManagedWorktree { worktree, .. } =
951 self.state.checkout(session_id)?
952 {
953 retire_managed_worktree(executor, worktree)
954 .context("retire managed raw-session worktree after verified close")?;
955 }
956 let record = self.state.sessions.get_mut(session_id).unwrap();
957 record.state = SessionState::Stopped;
958 if !deferred {
959 record.target = None;
960 }
961 record.updated_at = now();
962 record.last_error = None;
963 persist_session_record_transition_or_restore(
964 &mut self.state,
965 session_id,
966 &destroying,
967 "persist stopped state after target cleanup",
968 &persist,
969 )?;
970 Ok(deferred)
971 },
972 )
973 }
974
975 pub fn cleanup_stopped_target(
979 &mut self,
980 session_id: &str,
981 executor: &impl CommandExecutor,
982 ) -> Result<()> {
983 self.cleanup_stopped_target_with(
984 session_id,
985 executor,
986 crate::database::save_lifecycle_session,
987 )
988 }
989
990 fn cleanup_stopped_target_with(
991 &mut self,
992 session_id: &str,
993 executor: &impl CommandExecutor,
994 persist: impl Fn(&SessionRecord) -> Result<()>,
995 ) -> Result<()> {
996 crate::worker_lifecycle::run_blocking(
997 session_id,
998 "cleanup stopped target with",
999 executor,
1000 || {
1001 let previous = self
1002 .state
1003 .sessions
1004 .get(session_id)
1005 .with_context(|| format!("unknown session {session_id}"))?
1006 .clone();
1007 ensure!(
1008 previous.state == SessionState::Stopped,
1009 "refusing deferred cleanup for active session {session_id}"
1010 );
1011 let Some(locator) = previous.target.as_ref() else {
1012 return Ok(());
1013 };
1014 let backend = backend_locator(locator, &previous, &self.config)?;
1015 ensure!(
1016 targets::quiesce_plan(&backend, session_id)?.is_some(),
1017 "session {session_id} retained a non-Podman target after stopping"
1018 );
1019 if let Err(error) =
1020 execute_target_cleanup(&backend, &previous, &self.config, executor)
1021 {
1022 let record = self.state.sessions.get_mut(session_id).unwrap();
1023 record.updated_at = now();
1024 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
1025 let persisted = persist_session_record_transition_or_restore(
1026 &mut self.state,
1027 session_id,
1028 &previous,
1029 "persist deferred target cleanup failure",
1030 &persist,
1031 );
1032 return match persisted {
1033 Ok(()) => Err(error),
1034 Err(persist_error) => Err(error.context(format!(
1035 "also failed to persist deferred target cleanup failure: {persist_error:#}"
1036 ))),
1037 };
1038 }
1039 let record = self.state.sessions.get_mut(session_id).unwrap();
1040 record.target = None;
1041 record.updated_at = now();
1042 record.last_error = None;
1043 persist_session_record_transition_or_restore(
1044 &mut self.state,
1045 session_id,
1046 &previous,
1047 "persist completion of deferred Podman target cleanup",
1048 &persist,
1049 )
1050 },
1051 )
1052 }
1053
1054 pub fn force_stop(
1057 &mut self,
1058 session_id: &str,
1059 executor: &impl CommandExecutor,
1060 ) -> Result<bool> {
1061 self.force_stop_with(
1062 session_id,
1063 executor,
1064 crate::database::save_lifecycle_session,
1065 )
1066 }
1067
1068 fn force_stop_with(
1069 &mut self,
1070 session_id: &str,
1071 executor: &impl CommandExecutor,
1072 persist: impl Fn(&SessionRecord) -> Result<()>,
1073 ) -> Result<bool> {
1074 let session = self
1075 .state
1076 .sessions
1077 .get(session_id)
1078 .with_context(|| format!("unknown session {session_id}"))?
1079 .clone();
1080 ensure!(
1081 session.state.is_active(),
1082 "session {session_id} is already inactive"
1083 );
1084 let checkpoint = session
1085 .checkpoint
1086 .as_ref()
1087 .context("force stop requires an existing recovery archive")?;
1088 {
1091 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1092 verify_installed_checkpoint_gate(session_id, checkpoint)
1093 .context("verify the recovery archive before force stopping")?;
1094 }
1095 self.stop_target_and_settle(session_id, &session, executor, &persist)
1096 }
1097
1098 fn stop_target_and_settle(
1103 &mut self,
1104 session_id: &str,
1105 session: &SessionRecord,
1106 executor: &impl CommandExecutor,
1107 persist: &impl Fn(&SessionRecord) -> Result<()>,
1108 ) -> Result<bool> {
1109 crate::worker_lifecycle::run_blocking(
1110 session_id,
1111 "stop target and settle",
1112 executor,
1113 || {
1114 let mut deferred = false;
1115 let checkout = self.state.checkout(session_id)?;
1116 if let Some(locator) = &session.target {
1117 let backend = backend_locator(locator, session, &self.config)?;
1118 if self.state.subagents.contains_key(session_id) {
1122 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1123 .execute(executor)?;
1124 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1125 plan.execute(executor)?;
1126 deferred = true;
1127 } else {
1128 execute_target_cleanup(&backend, session, &self.config, executor)?;
1129 }
1130 }
1131 if let mj_core::state::Checkout::ManagedWorktree { worktree, .. } = checkout {
1132 retire_managed_worktree(executor, worktree)
1133 .context("retire managed raw-session worktree after stopping the target")?;
1134 }
1135 let record = self.state.sessions.get_mut(session_id).unwrap();
1136 record.state = SessionState::Stopped;
1137 if !deferred {
1138 record.target = None;
1139 }
1140 record.updated_at = now();
1141 record.last_error = None;
1142 record.last_checkpoint_error = None;
1143 persist_session_record_transition_or_restore(
1144 &mut self.state,
1145 session_id,
1146 session,
1147 "persist stopped state after tearing down the current target",
1148 persist,
1149 )?;
1150 Ok(deferred)
1151 },
1152 )
1153 }
1154
1155 pub fn destroy_session_controlled(
1159 &mut self,
1160 session_id: &str,
1161 executor: &impl CommandExecutor,
1162 ) -> Result<()> {
1163 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1164 }
1165
1166 pub fn destroy_session_controlled_with(
1174 &mut self,
1175 session_id: &str,
1176 executor: &impl CommandExecutor,
1177 branch: BranchDisposition,
1178 ) -> Result<()> {
1179 self.destroy_session_controlled_with_checkout(
1180 session_id,
1181 executor,
1182 branch,
1183 CheckoutDisposition::Remove,
1184 )
1185 .map(|_| ())
1186 }
1187
1188 pub(crate) fn destroy_session_controlled_with_checkout(
1194 &mut self,
1195 session_id: &str,
1196 executor: &impl CommandExecutor,
1197 branch: BranchDisposition,
1198 checkout: CheckoutDisposition,
1199 ) -> Result<Option<PathBuf>> {
1200 let session = self
1201 .state
1202 .sessions
1203 .get(session_id)
1204 .with_context(|| format!("unknown session {session_id}"))?
1205 .clone();
1206 if session.state.is_active() {
1207 bail!("refusing to destroy active session {session_id}");
1208 }
1209 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1210 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1211 }
1212 let mut retained_checkout = None;
1213 if let mj_core::state::Checkout::ManagedWorktree { worktree, .. } =
1214 self.state.checkout(session_id)?
1215 {
1216 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1217 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1218 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1219 "check the managed raw-session worktree for uncommitted changes",
1220 )?);
1221 if keep {
1222 retained_checkout = Some(worktree.worktree_root.clone());
1226 } else {
1227 cleanup_managed_worktree(executor, worktree, branch)
1228 .context("remove managed raw-session worktree")?;
1229 }
1230 }
1231 if let Some(checkpoint) = &session.checkpoint
1232 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1233 && error.kind() != std::io::ErrorKind::NotFound
1234 {
1235 return Err(error).with_context(|| {
1236 format!(
1237 "remove session recovery archive {}",
1238 checkpoint.archive_path.display()
1239 )
1240 });
1241 }
1242 mj_core::attachment::AttachmentStore::controller(session_id)?
1243 .remove_session_data()
1244 .context("remove session image attachments")?;
1245 crate::database::delete_session(session_id)
1246 .context("destroy stopped session in database")?;
1247 self.state.subagents.remove(session_id);
1248 self.state.destroy_stopped_session(session_id)?;
1249 Ok(retained_checkout)
1250 }
1251
1252 pub fn force_destroy_session(
1264 &mut self,
1265 session_id: &str,
1266 executor: &impl CommandExecutor,
1267 branch: BranchDisposition,
1268 ) -> Result<()> {
1269 self.force_destroy_session_with(
1270 session_id,
1271 executor,
1272 branch,
1273 crate::database::delete_session,
1274 )
1275 }
1276
1277 fn force_destroy_session_with(
1278 &mut self,
1279 session_id: &str,
1280 executor: &impl CommandExecutor,
1281 branch: BranchDisposition,
1282 delete: impl Fn(&str) -> Result<()>,
1283 ) -> Result<()> {
1284 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1285 if result.is_err() {
1286 self.settle_failed_destroy(session_id);
1287 }
1288 result
1289 }
1290
1291 fn settle_failed_destroy(&mut self, session_id: &str) {
1297 let Some(record) = self.state.sessions.get(session_id) else {
1298 return;
1299 };
1300 if !record.state.is_active() || record.state == SessionState::Error {
1301 return;
1302 }
1303 let previous = record.clone();
1304 let record = self.state.sessions.get_mut(session_id).unwrap();
1305 record.state = SessionState::Error;
1306 record.updated_at = now();
1307 if let Err(error) = persist_session_record_transition_or_restore(
1308 &mut self.state,
1309 session_id,
1310 &previous,
1311 "persist the failed destroy",
1312 &crate::database::save_lifecycle_session,
1313 ) {
1314 tracing::warn!(
1315 session_id,
1316 error = format!("{error:#}"),
1317 "could not record a failed destroy"
1318 );
1319 }
1320 }
1321
1322 fn force_destroy_session_steps(
1323 &mut self,
1324 session_id: &str,
1325 executor: &impl CommandExecutor,
1326 branch: BranchDisposition,
1327 delete: impl Fn(&str) -> Result<()>,
1328 ) -> Result<()> {
1329 crate::worker_lifecycle::run_blocking(
1330 session_id,
1331 "force destroy session steps",
1332 executor,
1333 || {
1334 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1335 let session = self
1336 .state
1337 .sessions
1338 .get(session_id)
1339 .with_context(|| format!("unknown session {session_id}"))?
1340 .clone();
1341 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1342 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1343 }
1344 if let Some(locator) = &session.target {
1348 let backend = backend_locator(locator, &session, &self.config)?;
1349 if self.state.subagents.contains_key(session_id) {
1350 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1351 .execute(executor)?;
1352 } else {
1353 execute_target_cleanup(&backend, &session, &self.config, executor)?;
1354 }
1355 }
1356 if let mj_core::state::Checkout::ManagedWorktree { worktree, .. } =
1357 self.state.checkout(session_id)?
1358 {
1359 cleanup_managed_worktree(executor, worktree, branch)
1360 .context("remove managed raw-session worktree")?;
1361 }
1362 if let Some(checkpoint) = &session.checkpoint
1363 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1364 && error.kind() != std::io::ErrorKind::NotFound
1365 {
1366 return Err(error).with_context(|| {
1367 format!(
1368 "remove session recovery archive {}",
1369 checkpoint.archive_path.display()
1370 )
1371 });
1372 }
1373 mj_core::attachment::AttachmentStore::controller(session_id)?
1374 .remove_session_data()
1375 .context("remove session image attachments")?;
1376 delete(session_id).context("force destroy session in database")?;
1377 self.state.subagents.remove(session_id);
1378 self.state.destroy_session_force(session_id)?;
1379 Ok(())
1380 },
1381 )
1382 }
1383}
1384
1385pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1400 match session.state {
1401 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1402 SessionState::Error if subagent => true,
1403 SessionState::Closing
1404 | SessionState::Destroying
1405 | SessionState::Error
1406 | SessionState::Lost => session.target.is_none(),
1407 SessionState::Running
1408 | SessionState::Disconnected
1409 | SessionState::Checkpointing
1410 | SessionState::Stopped
1411 | SessionState::DestroyedWithDataLoss => false,
1412 }
1413}
1414
1415fn execute_target_cleanup(
1419 backend: &targets::TargetLocator,
1420 session: &SessionRecord,
1421 config: &mj_core::config::Config,
1422 executor: &impl CommandExecutor,
1423) -> Result<()> {
1424 let session_id = session.id.as_str();
1425 let _owner = crate::worker_lifecycle::require(session_id)?;
1426 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1427 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1428 Ok(true) => {
1429 tracing::warn!(
1430 session_id,
1431 error = format!("{cleanup_error:#}"),
1432 "target cleanup command failed, but the target was confirmed absent"
1433 );
1434 }
1435 Ok(false) => {
1436 tracing::error!(
1437 session_id,
1438 error = format!("{cleanup_error:#}"),
1439 "target cleanup failed and the target is still present"
1440 );
1441 return Err(cleanup_error);
1442 }
1443 Err(probe_error) => {
1444 tracing::error!(
1445 session_id,
1446 cleanup_error = format!("{cleanup_error:#}"),
1447 probe_error = format!("{probe_error:#}"),
1448 "target cleanup failed and exact absence could not be confirmed"
1449 );
1450 return Err(cleanup_error.context(format!(
1451 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1452 )));
1453 }
1454 }
1455 }
1456 if let Some(release) = BuildStateRelease::for_target(session, backend, config) {
1457 release.run(executor);
1458 }
1459 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
1462 crate::database::finish_worker_restart(session_id, &intent.operation_id)?;
1463 }
1464 Ok(())
1465}
1466
1467fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1468 record.state = SessionState::Closing;
1469 record.updated_at = updated_at;
1470 record.last_checkpoint_error = None;
1471}
1472
1473fn apply_close_checkpoint_failure(
1480 record: &mut SessionRecord,
1481 previous: &SessionRecord,
1482 error: &anyhow::Error,
1483 updated_at: String,
1484) {
1485 if WorkerRestartLeftNoWorker::marks(error) {
1486 record.state = SessionState::Error;
1487 record.last_error = Some(format!(
1488 "close failed and left the session without a live worker; retry the close, \
1489 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1490 ));
1491 } else {
1492 record.state = state_after_unsealed_close(previous);
1493 }
1494 record.last_checkpoint_error = Some(format!("{error:#}"));
1495 record.updated_at = updated_at;
1496}
1497
1498fn close_was_refused(error: &anyhow::Error) -> bool {
1506 error.chain().any(|cause| {
1507 cause
1508 .downcast_ref::<crate::worker_client::RelayRejected>()
1509 .is_some_and(|rejected| !rejected.is_retryable())
1510 })
1511}
1512
1513fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1514 if previous.state == SessionState::Closing {
1515 SessionState::Running
1516 } else {
1517 previous.state
1518 }
1519}
1520
1521fn apply_interrupted_close_error(
1522 record: &mut SessionRecord,
1523 error: &anyhow::Error,
1524 updated_at: &str,
1525) {
1526 let destroying = record.state == SessionState::Destroying;
1527 if !destroying {
1528 record.state = SessionState::Closing;
1529 }
1530 record.updated_at = updated_at.to_owned();
1531 record.last_error = Some(if destroying {
1532 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1533 } else {
1534 format!("close is safely resumable from its verified checkpoint: {error:#}")
1535 });
1536}
1537
1538#[cfg(test)]
1539mod tests;