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 let validation = match self
333 .github_token_for_repository_preflight(session_id)
334 .await
335 {
336 Ok(token) => self.validate_move_checkpoint(
337 operation,
338 preparation,
339 token.as_deref(),
340 executor,
341 ),
342 Err(error) => Err(error),
343 };
344 if let Err(error) = validation {
345 let record = self.state.sessions.get_mut(session_id).unwrap();
346 record.state = previous.state;
347 record.last_error = Some(format!("{error:#}"));
348 self.persist_session_transition_or_restore(
349 session_id,
350 &previous,
351 "restore source after move preflight failure",
352 )?;
353 return Err(error);
354 }
355 operation.checkpoint = Some(artifact.metadata.clone());
356 operation.updated_at = now();
357 crate::database::save_move_operation(operation)?;
358 }
359 if let Some(before_close) = before_close
360 && let Err(error) = before_close.await
361 {
362 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
365 tracing::warn!(
366 session_id,
367 error = format!("{cancel:#}"),
368 "could not release the checkpoint barrier of a close that did not proceed"
369 );
370 }
371 let record = self.state.sessions.get_mut(session_id).unwrap();
372 record.state = state_after_unsealed_close(&previous);
373 record.last_error = Some(format!("{error:#}"));
374 record.updated_at = now();
375 self.persist_session_transition_or_restore(
376 session_id,
377 &previous,
378 "restore a session whose close did not proceed past its checkpoint",
379 )?;
380 return Err(error);
381 }
382 prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
383 release_projection_behind_checkpoint(session_id, &artifact.metadata);
386
387 let close_command_id = new_command_id("close")?;
388 let barrier_command_id = latched.barrier_command_id.clone();
389 let close_result = {
390 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
391 latched
392 .relay
393 .connection_mut()
394 .submit(
395 close_command_id,
396 RelayCommand::Close {
397 barrier_command_id: barrier_command_id.clone(),
398 expected: latched.cursor.clone(),
399 },
400 )
401 .await
402 };
403 if let Err(error) = close_result {
404 let error = error.context("seal verified checkpoint for close");
405 if !close_was_refused(&error) {
406 self.record_interrupted_close(session_id, &error)?;
409 return Err(error);
410 }
411 if let Err(cancel) = latched.relay.cancel_abandoned_barrier().await {
416 tracing::warn!(
417 session_id,
418 error = format!("{cancel:#}"),
419 "could not release the checkpoint barrier of a refused close"
420 );
421 }
422 let record = self.state.sessions.get_mut(session_id).unwrap();
423 record.state = state_after_unsealed_close(&previous);
424 record.last_error = Some(format!("{error:#}"));
425 record.updated_at = now();
426 self.persist_session_transition_or_restore(
427 session_id,
428 &previous,
429 "restore a session whose worker refused its close",
430 )?;
431 return Err(error);
432 }
433 let close_result = {
434 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
435 latched
436 .relay
437 .connection_mut()
438 .submit(
439 new_command_id("checkpoint-complete")?,
440 RelayCommand::CompleteCheckpoint { barrier_command_id },
441 )
442 .await
443 };
444 if let Err(error) = close_result {
445 self.record_interrupted_close(session_id, &error)?;
446 return Err(error.context("release verified close checkpoint"));
447 }
448 let close_result = {
449 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
450 wait_for_relay_closed(latched.relay.connection_mut()).await
451 };
452 if let Err(error) = close_result {
453 self.record_interrupted_close(session_id, &error)?;
454 return Err(error);
455 }
456 latched.relay.release();
457
458 if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
459 return Ok(false);
465 }
466 match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
467 Ok(deferred) => Ok(deferred),
468 Err(error) => {
469 self.record_interrupted_close(session_id, &error)?;
470 Err(error)
471 }
472 }
473
474 }).await
475 }
476
477 pub async fn recover_interrupted_close_managed(
489 &mut self,
490 session_id: &str,
491 executor: &(impl CommandExecutor + Sync),
492 manager: &SessionManagerControl,
493 acknowledge_unpublished_work: bool,
494 before_close: Option<BeforeClose>,
495 ) -> Result<bool> {
496 crate::worker_lifecycle::run(session_id, "recover interrupted close managed", executor, async {
497 let (state, verified) = {
498 let session = self
499 .state
500 .sessions
501 .get(session_id)
502 .with_context(|| format!("unknown session {session_id}"))?;
503 ensure!(
504 matches!(
505 session.state,
506 SessionState::Closing | SessionState::Destroying
507 ),
508 "session {session_id} has no interrupted close to recover"
509 );
510 (session.state, session.checkpoint.clone())
511 };
512 if state == SessionState::Destroying {
513 let verified = verified.context("destroying session has no verified checkpoint")?;
514 if let Some(before_close) = before_close {
515 before_close.await?;
516 }
517 return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
518 }
519 ensure!(
520 state == SessionState::Closing,
521 "session {session_id} has no relay close to recover"
522 );
523 let started = std::time::Instant::now();
524 let mut phase = "wait for session actor";
525 tracing::info!(
526 session_id,
527 ?state,
528 checkpoint_frontier = ?verified.as_ref().map(|checkpoint| checkpoint.event_frontier),
529 "recovering suspension: asking the worker whether its relay was sealed"
530 );
531 let connection = async {
532 let handle = manager
533 .wait_for_session(session_id, Duration::from_secs(5))
534 .await?;
535 phase = "lease worker connection";
536 let mut lease = handle.lease_connection().await?;
537 phase = "read worker execution state";
538 let execution = lease.connection_mut().sync().await?.operational.execution;
539 anyhow::Ok((lease, execution))
540 }
541 .await;
542 let (mut lease, execution) = match connection {
543 Ok(connected) => connected,
544 Err(error) => {
545 tracing::warn!(
546 session_id,
547 phase,
548 elapsed_ms = started.elapsed().as_millis() as u64,
549 error = format!("{error:#}"),
550 "suspension could not query its worker; retaining target and checkpoint"
551 );
552 self.diagnose_suspension_worker(session_id, phase).await;
553 return Err(error.context(format!("{phase} while recovering suspension")));
554 }
555 };
556 tracing::info!(
557 session_id,
558 ?execution,
559 elapsed_ms = started.elapsed().as_millis() as u64,
560 "suspension recovery read the worker's execution state"
561 );
562 match execution {
563 RelayExecutionState::Closed => {}
564 RelayExecutionState::Closing => {
565 let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
566 wait_for_relay_closed(lease.connection_mut()).await?;
567 }
568 RelayExecutionState::Idle | RelayExecutionState::Running => {
569 lease.release();
570 return self
571 .suspend_session_controlled_with_manager(
572 session_id,
573 executor,
574 Some(manager),
575 None,
576 SourceTargetDisposition::Destroy,
577 CheckpointExportPolicy::ReuseUnchangedArchive,
578 acknowledge_unpublished_work,
579 before_close,
580 None,
581 )
582 .await;
583 }
584 }
585 lease.release();
586 let verified = verified.context("closed relay has no verified checkpoint")?;
587 if let Some(before_close) = before_close {
589 before_close.await?;
590 }
591 self.destroy_after_verified_checkpoint(session_id, &verified, executor)
592
593 }).await
594 }
595
596 async fn diagnose_suspension_worker(&self, session_id: &str, phase: &str) {
599 let placement = self.worker_placement(session_id);
600 let probe = match placement {
601 Ok((backend, worker_root)) => tokio::task::spawn_blocking(move || {
602 super::worker_binary::probe_worker(
603 &targets::BoundedProcessExecutor::new(Duration::from_secs(5)),
604 &backend,
605 &worker_root,
606 )
607 })
608 .await
609 .context("join the suspension worker diagnostic probe")
610 .and_then(std::convert::identity),
611 Err(error) => Err(error.context("resolve the retained suspension worker")),
612 };
613 match probe {
614 Ok(probe) => tracing::warn!(
615 session_id,
616 phase,
617 worker_pids = ?probe.pids,
618 worker_startup_step = ?probe.step(),
619 worker_state = %probe,
620 "suspension failure worker diagnostic; no recovery mutation was attempted"
621 ),
622 Err(error) => tracing::warn!(
623 session_id,
624 phase,
625 error = format!("{error:#}"),
626 "suspension failure worker diagnostic was unavailable; worker state remains unknown"
627 ),
628 }
629 }
630
631 pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
638 self.fail_interrupted_lifecycle_with(
639 session_id,
640 cause,
641 crate::database::save_lifecycle_session,
642 )
643 }
644
645 fn fail_interrupted_lifecycle_with(
646 &mut self,
647 session_id: &str,
648 cause: &str,
649 persist: impl Fn(&SessionRecord) -> Result<()>,
650 ) -> Result<bool> {
651 let Some(record) = self.state.sessions.get_mut(session_id) else {
652 return Ok(false);
653 };
654 if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
655 return Ok(false);
656 }
657 let previous = record.clone();
658 record.state = SessionState::Error;
659 record.updated_at = now();
660 record.last_error = Some(cause.to_owned());
661 persist_session_record_transition_or_restore(
662 &mut self.state,
663 session_id,
664 &previous,
665 "persist the failure of an interrupted lifecycle state",
666 &persist,
667 )?;
668 Ok(true)
669 }
670
671 pub fn fail_unready_session(
680 &mut self,
681 session_id: &str,
682 cause: &str,
683 observed_updated_at: &str,
684 ) -> Result<bool> {
685 self.fail_unready_session_with(
686 session_id,
687 cause,
688 observed_updated_at,
689 crate::database::save_lifecycle_session,
690 )
691 }
692
693 pub(crate) async fn persist_harness_preparation_failure(
697 &self,
698 session_id: &str,
699 cause: &str,
700 observed_updated_at: &str,
701 ) -> Result<()> {
702 let session_id = session_id.to_owned();
703 let cause = cause.to_owned();
704 let observed_updated_at = observed_updated_at.to_owned();
705 tokio::task::spawn_blocking(move || {
706 let mut controller = Controller::load()?;
707 controller.fail_unready_session(&session_id, &cause, &observed_updated_at)
708 })
709 .await
710 .context("record harness preparation failure task panicked")??;
711 Ok(())
712 }
713
714 fn fail_unready_session_with(
715 &mut self,
716 session_id: &str,
717 cause: &str,
718 observed_updated_at: &str,
719 persist: impl Fn(&SessionRecord) -> Result<()>,
720 ) -> Result<bool> {
721 let Some(record) = self.state.sessions.get_mut(session_id) else {
722 return Ok(false);
723 };
724 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
725 return Ok(false);
726 }
727 let previous = record.clone();
728 record.state = SessionState::Error;
729 record.updated_at = now();
730 record.last_error = Some(cause.to_owned());
731 persist_session_record_transition_or_restore(
732 &mut self.state,
733 session_id,
734 &previous,
735 "persist a session whose harness never became usable",
736 &persist,
737 )?;
738 Ok(true)
739 }
740
741 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
745 let Some(record) = self.state.sessions.get(session_id) else {
746 return Ok(false);
747 };
748 let previous = record.clone();
751 let record = self.state.sessions.get_mut(session_id).unwrap();
752 record.last_error = Some(cause.to_owned());
753 record.updated_at = now();
754 persist_session_record_transition_or_restore(
755 &mut self.state,
756 session_id,
757 &previous,
758 "persist the reason a close did not finish",
759 &crate::database::save_lifecycle_session,
760 )?;
761 Ok(true)
762 }
763
764 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
769 let Some(record) = self.state.sessions.get(session_id) else {
770 return Ok(false);
771 };
772 if record.public_error().is_none() {
773 return Ok(false);
774 }
775 let previous = record.clone();
776 let record = self.state.sessions.get_mut(session_id).unwrap();
777 record.last_error = None;
778 record.updated_at = now();
779 persist_session_record_transition_or_restore(
780 &mut self.state,
781 session_id,
782 &previous,
783 "clear the reason a close did not finish",
784 &crate::database::save_lifecycle_session,
785 )?;
786 Ok(true)
787 }
788
789 pub fn suspend_session_without_checkpoint(
802 &mut self,
803 session_id: &str,
804 executor: &impl CommandExecutor,
805 ) -> Result<bool> {
806 self.suspend_session_without_checkpoint_with(
807 session_id,
808 executor,
809 crate::database::save_lifecycle_session,
810 )
811 }
812
813 fn suspend_session_without_checkpoint_with(
814 &mut self,
815 session_id: &str,
816 executor: &impl CommandExecutor,
817 persist: impl Fn(&SessionRecord) -> Result<()>,
818 ) -> Result<bool> {
819 let session = self
820 .state
821 .sessions
822 .get(session_id)
823 .with_context(|| format!("unknown session {session_id}"))?
824 .clone();
825 ensure!(
826 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
827 "session {session_id} has a workspace to checkpoint; suspend it instead"
828 );
829 self.stop_target_and_settle(session_id, &session, executor, &persist)
830 }
831
832 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
833 let record = self.state.sessions.get_mut(session_id).unwrap();
834 apply_interrupted_close_error(record, error, &now());
835 self.persist_session_state(session_id)
836 }
837
838 fn destroy_after_verified_checkpoint(
841 &mut self,
842 session_id: &str,
843 verified: &CheckpointMetadata,
844 executor: &impl CommandExecutor,
845 ) -> Result<bool> {
846 self.destroy_after_verified_checkpoint_with(
847 session_id,
848 verified,
849 executor,
850 crate::database::save_lifecycle_session,
851 )
852 }
853
854 fn destroy_after_verified_checkpoint_with(
855 &mut self,
856 session_id: &str,
857 verified: &CheckpointMetadata,
858 executor: &impl CommandExecutor,
859 persist: impl Fn(&SessionRecord) -> Result<()>,
860 ) -> Result<bool> {
861 crate::worker_lifecycle::run_blocking(
862 session_id,
863 "destroy after verified checkpoint with",
864 executor,
865 || {
866 let session = self
867 .state
868 .sessions
869 .get(session_id)
870 .with_context(|| format!("unknown session {session_id}"))?
871 .clone();
872 ensure!(
873 matches!(
874 session.state,
875 SessionState::Closing | SessionState::Destroying
876 ),
877 "refusing to destroy session {session_id}: it is not closing or destroying"
878 );
879 ensure!(
880 session.checkpoint.as_ref() == Some(verified),
881 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
882 );
883 if session.state == SessionState::Closing {
884 let record = self.state.sessions.get_mut(session_id).unwrap();
885 record.state = SessionState::Destroying;
886 record.updated_at = now();
887 record.last_error = None;
888 persist_session_record_transition_or_restore(
889 &mut self.state,
890 session_id,
891 &session,
892 "persist destroying state before target cleanup",
893 &persist,
894 )?;
895 }
896
897 let destroying = self
898 .state
899 .sessions
900 .get(session_id)
901 .expect("destroying session disappeared")
902 .clone();
903 {
904 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
905 verify_installed_checkpoint_gate(session_id, verified)?;
906 }
907 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
912 tracing::warn!(
913 session_id,
914 error = format!("{error:#}"),
915 "could not record that the second-opinion conversation ends with this target"
916 );
917 }
918 let locator = destroying
919 .target
920 .as_ref()
921 .context("session has no target")?;
922 let backend = backend_locator(locator, &destroying, &self.config)?;
923 let deferred = if self.state.subagents.contains_key(session_id) {
924 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
925 .execute(executor)?;
926 false
927 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
928 plan.execute(executor)?;
929 true
930 } else {
931 execute_target_cleanup(&backend, &destroying, &self.config, executor)?;
932 false
933 };
934 if let Some(worktree) = &destroying.managed_worktree {
935 retire_managed_worktree(executor, worktree)
936 .context("retire managed raw-session worktree after verified close")?;
937 }
938 let record = self.state.sessions.get_mut(session_id).unwrap();
939 record.state = SessionState::Stopped;
940 if !deferred {
941 record.target = None;
942 }
943 record.updated_at = now();
944 record.last_error = None;
945 persist_session_record_transition_or_restore(
946 &mut self.state,
947 session_id,
948 &destroying,
949 "persist stopped state after target cleanup",
950 &persist,
951 )?;
952 Ok(deferred)
953 },
954 )
955 }
956
957 pub fn cleanup_stopped_target(
961 &mut self,
962 session_id: &str,
963 executor: &impl CommandExecutor,
964 ) -> Result<()> {
965 self.cleanup_stopped_target_with(
966 session_id,
967 executor,
968 crate::database::save_lifecycle_session,
969 )
970 }
971
972 fn cleanup_stopped_target_with(
973 &mut self,
974 session_id: &str,
975 executor: &impl CommandExecutor,
976 persist: impl Fn(&SessionRecord) -> Result<()>,
977 ) -> Result<()> {
978 crate::worker_lifecycle::run_blocking(
979 session_id,
980 "cleanup stopped target with",
981 executor,
982 || {
983 let previous = self
984 .state
985 .sessions
986 .get(session_id)
987 .with_context(|| format!("unknown session {session_id}"))?
988 .clone();
989 ensure!(
990 previous.state == SessionState::Stopped,
991 "refusing deferred cleanup for active session {session_id}"
992 );
993 let Some(locator) = previous.target.as_ref() else {
994 return Ok(());
995 };
996 let backend = backend_locator(locator, &previous, &self.config)?;
997 ensure!(
998 targets::quiesce_plan(&backend, session_id)?.is_some(),
999 "session {session_id} retained a non-Podman target after stopping"
1000 );
1001 if let Err(error) =
1002 execute_target_cleanup(&backend, &previous, &self.config, executor)
1003 {
1004 let record = self.state.sessions.get_mut(session_id).unwrap();
1005 record.updated_at = now();
1006 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
1007 let persisted = persist_session_record_transition_or_restore(
1008 &mut self.state,
1009 session_id,
1010 &previous,
1011 "persist deferred target cleanup failure",
1012 &persist,
1013 );
1014 return match persisted {
1015 Ok(()) => Err(error),
1016 Err(persist_error) => Err(error.context(format!(
1017 "also failed to persist deferred target cleanup failure: {persist_error:#}"
1018 ))),
1019 };
1020 }
1021 let record = self.state.sessions.get_mut(session_id).unwrap();
1022 record.target = None;
1023 record.updated_at = now();
1024 record.last_error = None;
1025 persist_session_record_transition_or_restore(
1026 &mut self.state,
1027 session_id,
1028 &previous,
1029 "persist completion of deferred Podman target cleanup",
1030 &persist,
1031 )
1032 },
1033 )
1034 }
1035
1036 pub fn force_stop(
1039 &mut self,
1040 session_id: &str,
1041 executor: &impl CommandExecutor,
1042 ) -> Result<bool> {
1043 self.force_stop_with(
1044 session_id,
1045 executor,
1046 crate::database::save_lifecycle_session,
1047 )
1048 }
1049
1050 fn force_stop_with(
1051 &mut self,
1052 session_id: &str,
1053 executor: &impl CommandExecutor,
1054 persist: impl Fn(&SessionRecord) -> Result<()>,
1055 ) -> Result<bool> {
1056 let session = self
1057 .state
1058 .sessions
1059 .get(session_id)
1060 .with_context(|| format!("unknown session {session_id}"))?
1061 .clone();
1062 ensure!(
1063 session.state.is_active(),
1064 "session {session_id} is already inactive"
1065 );
1066 let checkpoint = session
1067 .checkpoint
1068 .as_ref()
1069 .context("force stop requires an existing recovery archive")?;
1070 {
1073 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1074 verify_installed_checkpoint_gate(session_id, checkpoint)
1075 .context("verify the recovery archive before force stopping")?;
1076 }
1077 self.stop_target_and_settle(session_id, &session, executor, &persist)
1078 }
1079
1080 fn stop_target_and_settle(
1085 &mut self,
1086 session_id: &str,
1087 session: &SessionRecord,
1088 executor: &impl CommandExecutor,
1089 persist: &impl Fn(&SessionRecord) -> Result<()>,
1090 ) -> Result<bool> {
1091 crate::worker_lifecycle::run_blocking(
1092 session_id,
1093 "stop target and settle",
1094 executor,
1095 || {
1096 let mut deferred = false;
1097 if let Some(locator) = &session.target {
1098 let backend = backend_locator(locator, session, &self.config)?;
1099 if self.state.subagents.contains_key(session_id) {
1103 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1104 .execute(executor)?;
1105 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1106 plan.execute(executor)?;
1107 deferred = true;
1108 } else {
1109 execute_target_cleanup(&backend, session, &self.config, executor)?;
1110 }
1111 }
1112 if let Some(worktree) = &session.managed_worktree {
1113 retire_managed_worktree(executor, worktree)
1114 .context("retire managed raw-session worktree after stopping the target")?;
1115 }
1116 let record = self.state.sessions.get_mut(session_id).unwrap();
1117 record.state = SessionState::Stopped;
1118 if !deferred {
1119 record.target = None;
1120 }
1121 record.updated_at = now();
1122 record.last_error = None;
1123 record.last_checkpoint_error = None;
1124 persist_session_record_transition_or_restore(
1125 &mut self.state,
1126 session_id,
1127 session,
1128 "persist stopped state after tearing down the current target",
1129 persist,
1130 )?;
1131 Ok(deferred)
1132 },
1133 )
1134 }
1135
1136 pub fn destroy_session_controlled(
1140 &mut self,
1141 session_id: &str,
1142 executor: &impl CommandExecutor,
1143 ) -> Result<()> {
1144 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1145 }
1146
1147 pub fn destroy_session_controlled_with(
1155 &mut self,
1156 session_id: &str,
1157 executor: &impl CommandExecutor,
1158 branch: BranchDisposition,
1159 ) -> Result<()> {
1160 self.destroy_session_controlled_with_checkout(
1161 session_id,
1162 executor,
1163 branch,
1164 CheckoutDisposition::Remove,
1165 )
1166 .map(|_| ())
1167 }
1168
1169 pub(crate) fn destroy_session_controlled_with_checkout(
1175 &mut self,
1176 session_id: &str,
1177 executor: &impl CommandExecutor,
1178 branch: BranchDisposition,
1179 checkout: CheckoutDisposition,
1180 ) -> Result<Option<PathBuf>> {
1181 let session = self
1182 .state
1183 .sessions
1184 .get(session_id)
1185 .with_context(|| format!("unknown session {session_id}"))?
1186 .clone();
1187 if session.state.is_active() {
1188 bail!("refusing to destroy active session {session_id}");
1189 }
1190 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1191 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1192 }
1193 let mut retained_checkout = None;
1194 if let Some(worktree) = &session.managed_worktree {
1195 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1196 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1197 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1198 "check the managed raw-session worktree for uncommitted changes",
1199 )?);
1200 if keep {
1201 retained_checkout = Some(worktree.worktree_root.clone());
1205 } else {
1206 cleanup_managed_worktree(executor, worktree, branch)
1207 .context("remove managed raw-session worktree")?;
1208 }
1209 }
1210 if let Some(checkpoint) = &session.checkpoint
1211 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1212 && error.kind() != std::io::ErrorKind::NotFound
1213 {
1214 return Err(error).with_context(|| {
1215 format!(
1216 "remove session recovery archive {}",
1217 checkpoint.archive_path.display()
1218 )
1219 });
1220 }
1221 mj_core::attachment::AttachmentStore::controller(session_id)?
1222 .remove_session_data()
1223 .context("remove session image attachments")?;
1224 crate::database::delete_session(session_id)
1225 .context("destroy stopped session in database")?;
1226 self.state.subagents.remove(session_id);
1227 self.state.destroy_stopped_session(session_id)?;
1228 Ok(retained_checkout)
1229 }
1230
1231 pub fn force_destroy_session(
1243 &mut self,
1244 session_id: &str,
1245 executor: &impl CommandExecutor,
1246 branch: BranchDisposition,
1247 ) -> Result<()> {
1248 self.force_destroy_session_with(
1249 session_id,
1250 executor,
1251 branch,
1252 crate::database::delete_session,
1253 )
1254 }
1255
1256 fn force_destroy_session_with(
1257 &mut self,
1258 session_id: &str,
1259 executor: &impl CommandExecutor,
1260 branch: BranchDisposition,
1261 delete: impl Fn(&str) -> Result<()>,
1262 ) -> Result<()> {
1263 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1264 if result.is_err() {
1265 self.settle_failed_destroy(session_id);
1266 }
1267 result
1268 }
1269
1270 fn settle_failed_destroy(&mut self, session_id: &str) {
1276 let Some(record) = self.state.sessions.get(session_id) else {
1277 return;
1278 };
1279 if !record.state.is_active() || record.state == SessionState::Error {
1280 return;
1281 }
1282 let previous = record.clone();
1283 let record = self.state.sessions.get_mut(session_id).unwrap();
1284 record.state = SessionState::Error;
1285 record.updated_at = now();
1286 if let Err(error) = persist_session_record_transition_or_restore(
1287 &mut self.state,
1288 session_id,
1289 &previous,
1290 "persist the failed destroy",
1291 &crate::database::save_lifecycle_session,
1292 ) {
1293 tracing::warn!(
1294 session_id,
1295 error = format!("{error:#}"),
1296 "could not record a failed destroy"
1297 );
1298 }
1299 }
1300
1301 fn force_destroy_session_steps(
1302 &mut self,
1303 session_id: &str,
1304 executor: &impl CommandExecutor,
1305 branch: BranchDisposition,
1306 delete: impl Fn(&str) -> Result<()>,
1307 ) -> Result<()> {
1308 crate::worker_lifecycle::run_blocking(
1309 session_id,
1310 "force destroy session steps",
1311 executor,
1312 || {
1313 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1314 let session = self
1315 .state
1316 .sessions
1317 .get(session_id)
1318 .with_context(|| format!("unknown session {session_id}"))?
1319 .clone();
1320 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1321 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1322 }
1323 if let Some(locator) = &session.target {
1327 let backend = backend_locator(locator, &session, &self.config)?;
1328 if self.state.subagents.contains_key(session_id) {
1329 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1330 .execute(executor)?;
1331 } else {
1332 execute_target_cleanup(&backend, &session, &self.config, executor)?;
1333 }
1334 }
1335 if let Some(worktree) = &session.managed_worktree {
1336 cleanup_managed_worktree(executor, worktree, branch)
1337 .context("remove managed raw-session worktree")?;
1338 }
1339 if let Some(checkpoint) = &session.checkpoint
1340 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1341 && error.kind() != std::io::ErrorKind::NotFound
1342 {
1343 return Err(error).with_context(|| {
1344 format!(
1345 "remove session recovery archive {}",
1346 checkpoint.archive_path.display()
1347 )
1348 });
1349 }
1350 mj_core::attachment::AttachmentStore::controller(session_id)?
1351 .remove_session_data()
1352 .context("remove session image attachments")?;
1353 delete(session_id).context("force destroy session in database")?;
1354 self.state.subagents.remove(session_id);
1355 self.state.destroy_session_force(session_id)?;
1356 Ok(())
1357 },
1358 )
1359 }
1360}
1361
1362pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1377 match session.state {
1378 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1379 SessionState::Error if subagent => true,
1380 SessionState::Closing
1381 | SessionState::Destroying
1382 | SessionState::Error
1383 | SessionState::Lost => session.target.is_none(),
1384 SessionState::Running
1385 | SessionState::Disconnected
1386 | SessionState::Checkpointing
1387 | SessionState::Stopped
1388 | SessionState::DestroyedWithDataLoss => false,
1389 }
1390}
1391
1392fn execute_target_cleanup(
1396 backend: &targets::TargetLocator,
1397 session: &SessionRecord,
1398 config: &mj_core::config::Config,
1399 executor: &impl CommandExecutor,
1400) -> Result<()> {
1401 let session_id = session.id.as_str();
1402 let _owner = crate::worker_lifecycle::require(session_id)?;
1403 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1404 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1405 Ok(true) => {
1406 tracing::warn!(
1407 session_id,
1408 error = format!("{cleanup_error:#}"),
1409 "target cleanup command failed, but the target was confirmed absent"
1410 );
1411 }
1412 Ok(false) => {
1413 tracing::error!(
1414 session_id,
1415 error = format!("{cleanup_error:#}"),
1416 "target cleanup failed and the target is still present"
1417 );
1418 return Err(cleanup_error);
1419 }
1420 Err(probe_error) => {
1421 tracing::error!(
1422 session_id,
1423 cleanup_error = format!("{cleanup_error:#}"),
1424 probe_error = format!("{probe_error:#}"),
1425 "target cleanup failed and exact absence could not be confirmed"
1426 );
1427 return Err(cleanup_error.context(format!(
1428 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1429 )));
1430 }
1431 }
1432 }
1433 if let Some(release) = BuildStateRelease::for_target(session, backend, config) {
1434 release.run(executor);
1435 }
1436 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
1439 crate::database::finish_worker_restart(session_id, &intent.operation_id)?;
1440 }
1441 Ok(())
1442}
1443
1444fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1445 record.state = SessionState::Closing;
1446 record.updated_at = updated_at;
1447 record.last_checkpoint_error = None;
1448}
1449
1450fn apply_close_checkpoint_failure(
1457 record: &mut SessionRecord,
1458 previous: &SessionRecord,
1459 error: &anyhow::Error,
1460 updated_at: String,
1461) {
1462 if WorkerRestartLeftNoWorker::marks(error) {
1463 record.state = SessionState::Error;
1464 record.last_error = Some(format!(
1465 "close failed and left the session without a live worker; retry the close, \
1466 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1467 ));
1468 } else {
1469 record.state = state_after_unsealed_close(previous);
1470 }
1471 record.last_checkpoint_error = Some(format!("{error:#}"));
1472 record.updated_at = updated_at;
1473}
1474
1475fn close_was_refused(error: &anyhow::Error) -> bool {
1483 error.chain().any(|cause| {
1484 cause
1485 .downcast_ref::<crate::worker_client::RelayRejected>()
1486 .is_some_and(|rejected| !rejected.is_retryable())
1487 })
1488}
1489
1490fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1491 if previous.state == SessionState::Closing {
1492 SessionState::Running
1493 } else {
1494 previous.state
1495 }
1496}
1497
1498fn apply_interrupted_close_error(
1499 record: &mut SessionRecord,
1500 error: &anyhow::Error,
1501 updated_at: &str,
1502) {
1503 let destroying = record.state == SessionState::Destroying;
1504 if !destroying {
1505 record.state = SessionState::Closing;
1506 }
1507 record.updated_at = updated_at.to_owned();
1508 record.last_error = Some(if destroying {
1509 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1510 } else {
1511 format!("close is safely resumable from its verified checkpoint: {error:#}")
1512 });
1513}
1514
1515#[cfg(test)]
1516mod tests;