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 pub(crate) async fn persist_harness_preparation_failure(
685 &self,
686 session_id: &str,
687 cause: &str,
688 observed_updated_at: &str,
689 ) -> Result<()> {
690 let session_id = session_id.to_owned();
691 let cause = cause.to_owned();
692 let observed_updated_at = observed_updated_at.to_owned();
693 tokio::task::spawn_blocking(move || {
694 let mut controller = Controller::load()?;
695 controller.fail_unready_session(&session_id, &cause, &observed_updated_at)
696 })
697 .await
698 .context("record harness preparation failure task panicked")??;
699 Ok(())
700 }
701
702 fn fail_unready_session_with(
703 &mut self,
704 session_id: &str,
705 cause: &str,
706 observed_updated_at: &str,
707 persist: impl Fn(&SessionRecord) -> Result<()>,
708 ) -> Result<bool> {
709 let Some(record) = self.state.sessions.get_mut(session_id) else {
710 return Ok(false);
711 };
712 if record.state != SessionState::Running || record.updated_at != observed_updated_at {
713 return Ok(false);
714 }
715 let previous = record.clone();
716 record.state = SessionState::Error;
717 record.updated_at = now();
718 record.last_error = Some(cause.to_owned());
719 persist_session_record_transition_or_restore(
720 &mut self.state,
721 session_id,
722 &previous,
723 "persist a session whose harness never became usable",
724 &persist,
725 )?;
726 Ok(true)
727 }
728
729 pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
733 let Some(record) = self.state.sessions.get(session_id) else {
734 return Ok(false);
735 };
736 let previous = record.clone();
739 let record = self.state.sessions.get_mut(session_id).unwrap();
740 record.last_error = Some(cause.to_owned());
741 record.updated_at = now();
742 persist_session_record_transition_or_restore(
743 &mut self.state,
744 session_id,
745 &previous,
746 "persist the reason a close did not finish",
747 &crate::database::save_lifecycle_session,
748 )?;
749 Ok(true)
750 }
751
752 pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
757 let Some(record) = self.state.sessions.get(session_id) else {
758 return Ok(false);
759 };
760 if record.public_error().is_none() {
761 return Ok(false);
762 }
763 let previous = record.clone();
764 let record = self.state.sessions.get_mut(session_id).unwrap();
765 record.last_error = None;
766 record.updated_at = now();
767 persist_session_record_transition_or_restore(
768 &mut self.state,
769 session_id,
770 &previous,
771 "clear the reason a close did not finish",
772 &crate::database::save_lifecycle_session,
773 )?;
774 Ok(true)
775 }
776
777 pub fn suspend_session_without_checkpoint(
790 &mut self,
791 session_id: &str,
792 executor: &impl CommandExecutor,
793 ) -> Result<bool> {
794 self.suspend_session_without_checkpoint_with(
795 session_id,
796 executor,
797 crate::database::save_lifecycle_session,
798 )
799 }
800
801 fn suspend_session_without_checkpoint_with(
802 &mut self,
803 session_id: &str,
804 executor: &impl CommandExecutor,
805 persist: impl Fn(&SessionRecord) -> Result<()>,
806 ) -> Result<bool> {
807 let session = self
808 .state
809 .sessions
810 .get(session_id)
811 .with_context(|| format!("unknown session {session_id}"))?
812 .clone();
813 ensure!(
814 has_nothing_to_checkpoint(&session, self.state.subagents.contains_key(session_id)),
815 "session {session_id} has a workspace to checkpoint; suspend it instead"
816 );
817 self.stop_target_and_settle(session_id, &session, executor, &persist)
818 }
819
820 fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
821 let record = self.state.sessions.get_mut(session_id).unwrap();
822 apply_interrupted_close_error(record, error, &now());
823 self.persist_session_state(session_id)
824 }
825
826 fn destroy_after_verified_checkpoint(
829 &mut self,
830 session_id: &str,
831 verified: &CheckpointMetadata,
832 executor: &impl CommandExecutor,
833 ) -> Result<bool> {
834 self.destroy_after_verified_checkpoint_with(
835 session_id,
836 verified,
837 executor,
838 crate::database::save_lifecycle_session,
839 )
840 }
841
842 fn destroy_after_verified_checkpoint_with(
843 &mut self,
844 session_id: &str,
845 verified: &CheckpointMetadata,
846 executor: &impl CommandExecutor,
847 persist: impl Fn(&SessionRecord) -> Result<()>,
848 ) -> Result<bool> {
849 crate::worker_lifecycle::run_blocking(
850 session_id,
851 "destroy after verified checkpoint with",
852 executor,
853 || {
854 let session = self
855 .state
856 .sessions
857 .get(session_id)
858 .with_context(|| format!("unknown session {session_id}"))?
859 .clone();
860 ensure!(
861 matches!(
862 session.state,
863 SessionState::Closing | SessionState::Destroying
864 ),
865 "refusing to destroy session {session_id}: it is not closing or destroying"
866 );
867 ensure!(
868 session.checkpoint.as_ref() == Some(verified),
869 "refusing to destroy session {session_id}: verified checkpoint gate is stale"
870 );
871 if session.state == SessionState::Closing {
872 let record = self.state.sessions.get_mut(session_id).unwrap();
873 record.state = SessionState::Destroying;
874 record.updated_at = now();
875 record.last_error = None;
876 persist_session_record_transition_or_restore(
877 &mut self.state,
878 session_id,
879 &session,
880 "persist destroying state before target cleanup",
881 &persist,
882 )?;
883 }
884
885 let destroying = self
886 .state
887 .sessions
888 .get(session_id)
889 .expect("destroying session disappeared")
890 .clone();
891 {
892 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
893 verify_installed_checkpoint_gate(session_id, verified)?;
894 }
895 if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
900 tracing::warn!(
901 session_id,
902 error = format!("{error:#}"),
903 "could not record that the second-opinion conversation ends with this target"
904 );
905 }
906 let locator = destroying
907 .target
908 .as_ref()
909 .context("session has no target")?;
910 let backend = backend_locator(locator, &destroying, &self.config)?;
911 let deferred = if self.state.subagents.contains_key(session_id) {
912 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
913 .execute(executor)?;
914 false
915 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
916 plan.execute(executor)?;
917 true
918 } else {
919 execute_target_cleanup(&backend, &destroying, &self.config, executor)?;
920 false
921 };
922 if let Some(worktree) = &destroying.managed_worktree {
923 retire_managed_worktree(executor, worktree)
924 .context("retire managed raw-session worktree after verified close")?;
925 }
926 let record = self.state.sessions.get_mut(session_id).unwrap();
927 record.state = SessionState::Stopped;
928 if !deferred {
929 record.target = None;
930 }
931 record.updated_at = now();
932 record.last_error = None;
933 persist_session_record_transition_or_restore(
934 &mut self.state,
935 session_id,
936 &destroying,
937 "persist stopped state after target cleanup",
938 &persist,
939 )?;
940 Ok(deferred)
941 },
942 )
943 }
944
945 pub fn cleanup_stopped_target(
949 &mut self,
950 session_id: &str,
951 executor: &impl CommandExecutor,
952 ) -> Result<()> {
953 self.cleanup_stopped_target_with(
954 session_id,
955 executor,
956 crate::database::save_lifecycle_session,
957 )
958 }
959
960 fn cleanup_stopped_target_with(
961 &mut self,
962 session_id: &str,
963 executor: &impl CommandExecutor,
964 persist: impl Fn(&SessionRecord) -> Result<()>,
965 ) -> Result<()> {
966 crate::worker_lifecycle::run_blocking(
967 session_id,
968 "cleanup stopped target with",
969 executor,
970 || {
971 let previous = self
972 .state
973 .sessions
974 .get(session_id)
975 .with_context(|| format!("unknown session {session_id}"))?
976 .clone();
977 ensure!(
978 previous.state == SessionState::Stopped,
979 "refusing deferred cleanup for active session {session_id}"
980 );
981 let Some(locator) = previous.target.as_ref() else {
982 return Ok(());
983 };
984 let backend = backend_locator(locator, &previous, &self.config)?;
985 ensure!(
986 targets::quiesce_plan(&backend, session_id)?.is_some(),
987 "session {session_id} retained a non-Podman target after stopping"
988 );
989 if let Err(error) =
990 execute_target_cleanup(&backend, &previous, &self.config, executor)
991 {
992 let record = self.state.sessions.get_mut(session_id).unwrap();
993 record.updated_at = now();
994 record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
995 let persisted = persist_session_record_transition_or_restore(
996 &mut self.state,
997 session_id,
998 &previous,
999 "persist deferred target cleanup failure",
1000 &persist,
1001 );
1002 return match persisted {
1003 Ok(()) => Err(error),
1004 Err(persist_error) => Err(error.context(format!(
1005 "also failed to persist deferred target cleanup failure: {persist_error:#}"
1006 ))),
1007 };
1008 }
1009 let record = self.state.sessions.get_mut(session_id).unwrap();
1010 record.target = None;
1011 record.updated_at = now();
1012 record.last_error = None;
1013 persist_session_record_transition_or_restore(
1014 &mut self.state,
1015 session_id,
1016 &previous,
1017 "persist completion of deferred Podman target cleanup",
1018 &persist,
1019 )
1020 },
1021 )
1022 }
1023
1024 pub fn force_stop(
1027 &mut self,
1028 session_id: &str,
1029 executor: &impl CommandExecutor,
1030 ) -> Result<bool> {
1031 self.force_stop_with(
1032 session_id,
1033 executor,
1034 crate::database::save_lifecycle_session,
1035 )
1036 }
1037
1038 fn force_stop_with(
1039 &mut self,
1040 session_id: &str,
1041 executor: &impl CommandExecutor,
1042 persist: impl Fn(&SessionRecord) -> Result<()>,
1043 ) -> Result<bool> {
1044 let session = self
1045 .state
1046 .sessions
1047 .get(session_id)
1048 .with_context(|| format!("unknown session {session_id}"))?
1049 .clone();
1050 ensure!(
1051 session.state.is_active(),
1052 "session {session_id} is already inactive"
1053 );
1054 let checkpoint = session
1055 .checkpoint
1056 .as_ref()
1057 .context("force stop requires an existing recovery archive")?;
1058 {
1061 let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
1062 verify_installed_checkpoint_gate(session_id, checkpoint)
1063 .context("verify the recovery archive before force stopping")?;
1064 }
1065 self.stop_target_and_settle(session_id, &session, executor, &persist)
1066 }
1067
1068 fn stop_target_and_settle(
1073 &mut self,
1074 session_id: &str,
1075 session: &SessionRecord,
1076 executor: &impl CommandExecutor,
1077 persist: &impl Fn(&SessionRecord) -> Result<()>,
1078 ) -> Result<bool> {
1079 crate::worker_lifecycle::run_blocking(
1080 session_id,
1081 "stop target and settle",
1082 executor,
1083 || {
1084 let mut deferred = false;
1085 if let Some(locator) = &session.target {
1086 let backend = backend_locator(locator, session, &self.config)?;
1087 if self.state.subagents.contains_key(session_id) {
1091 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1092 .execute(executor)?;
1093 } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
1094 plan.execute(executor)?;
1095 deferred = true;
1096 } else {
1097 execute_target_cleanup(&backend, session, &self.config, executor)?;
1098 }
1099 }
1100 if let Some(worktree) = &session.managed_worktree {
1101 retire_managed_worktree(executor, worktree)
1102 .context("retire managed raw-session worktree after stopping the target")?;
1103 }
1104 let record = self.state.sessions.get_mut(session_id).unwrap();
1105 record.state = SessionState::Stopped;
1106 if !deferred {
1107 record.target = None;
1108 }
1109 record.updated_at = now();
1110 record.last_error = None;
1111 record.last_checkpoint_error = None;
1112 persist_session_record_transition_or_restore(
1113 &mut self.state,
1114 session_id,
1115 session,
1116 "persist stopped state after tearing down the current target",
1117 persist,
1118 )?;
1119 Ok(deferred)
1120 },
1121 )
1122 }
1123
1124 pub fn destroy_session_controlled(
1128 &mut self,
1129 session_id: &str,
1130 executor: &impl CommandExecutor,
1131 ) -> Result<()> {
1132 self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
1133 }
1134
1135 pub fn destroy_session_controlled_with(
1143 &mut self,
1144 session_id: &str,
1145 executor: &impl CommandExecutor,
1146 branch: BranchDisposition,
1147 ) -> Result<()> {
1148 self.destroy_session_controlled_with_checkout(
1149 session_id,
1150 executor,
1151 branch,
1152 CheckoutDisposition::Remove,
1153 )
1154 .map(|_| ())
1155 }
1156
1157 pub(crate) fn destroy_session_controlled_with_checkout(
1163 &mut self,
1164 session_id: &str,
1165 executor: &impl CommandExecutor,
1166 branch: BranchDisposition,
1167 checkout: CheckoutDisposition,
1168 ) -> Result<Option<PathBuf>> {
1169 let session = self
1170 .state
1171 .sessions
1172 .get(session_id)
1173 .with_context(|| format!("unknown session {session_id}"))?
1174 .clone();
1175 if session.state.is_active() {
1176 bail!("refusing to destroy active session {session_id}");
1177 }
1178 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1179 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1180 }
1181 let mut retained_checkout = None;
1182 if let Some(worktree) = &session.managed_worktree {
1183 let keep = checkout == CheckoutDisposition::KeepWhenDirty
1184 && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
1185 || managed_worktree_checkout_is_dirty(executor, worktree).context(
1186 "check the managed raw-session worktree for uncommitted changes",
1187 )?);
1188 if keep {
1189 retained_checkout = Some(worktree.worktree_root.clone());
1193 } else {
1194 cleanup_managed_worktree(executor, worktree, branch)
1195 .context("remove managed raw-session worktree")?;
1196 }
1197 }
1198 if let Some(checkpoint) = &session.checkpoint
1199 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1200 && error.kind() != std::io::ErrorKind::NotFound
1201 {
1202 return Err(error).with_context(|| {
1203 format!(
1204 "remove session recovery archive {}",
1205 checkpoint.archive_path.display()
1206 )
1207 });
1208 }
1209 mj_core::attachment::AttachmentStore::controller(session_id)?
1210 .remove_session_data()
1211 .context("remove session image attachments")?;
1212 crate::database::delete_session(session_id)
1213 .context("destroy stopped session in database")?;
1214 self.state.subagents.remove(session_id);
1215 self.state.destroy_stopped_session(session_id)?;
1216 Ok(retained_checkout)
1217 }
1218
1219 pub fn force_destroy_session(
1231 &mut self,
1232 session_id: &str,
1233 executor: &impl CommandExecutor,
1234 branch: BranchDisposition,
1235 ) -> Result<()> {
1236 self.force_destroy_session_with(
1237 session_id,
1238 executor,
1239 branch,
1240 crate::database::delete_session,
1241 )
1242 }
1243
1244 fn force_destroy_session_with(
1245 &mut self,
1246 session_id: &str,
1247 executor: &impl CommandExecutor,
1248 branch: BranchDisposition,
1249 delete: impl Fn(&str) -> Result<()>,
1250 ) -> Result<()> {
1251 let result = self.force_destroy_session_steps(session_id, executor, branch, delete);
1252 if result.is_err() {
1253 self.settle_failed_destroy(session_id);
1254 }
1255 result
1256 }
1257
1258 fn settle_failed_destroy(&mut self, session_id: &str) {
1264 let Some(record) = self.state.sessions.get(session_id) else {
1265 return;
1266 };
1267 if !record.state.is_active() || record.state == SessionState::Error {
1268 return;
1269 }
1270 let previous = record.clone();
1271 let record = self.state.sessions.get_mut(session_id).unwrap();
1272 record.state = SessionState::Error;
1273 record.updated_at = now();
1274 if let Err(error) = persist_session_record_transition_or_restore(
1275 &mut self.state,
1276 session_id,
1277 &previous,
1278 "persist the failed destroy",
1279 &crate::database::save_lifecycle_session,
1280 ) {
1281 tracing::warn!(
1282 session_id,
1283 error = format!("{error:#}"),
1284 "could not record a failed destroy"
1285 );
1286 }
1287 }
1288
1289 fn force_destroy_session_steps(
1290 &mut self,
1291 session_id: &str,
1292 executor: &impl CommandExecutor,
1293 branch: BranchDisposition,
1294 delete: impl Fn(&str) -> Result<()>,
1295 ) -> Result<()> {
1296 crate::worker_lifecycle::run_blocking(
1297 session_id,
1298 "force destroy session steps",
1299 executor,
1300 || {
1301 crate::worker_lifecycle::require(session_id)?.verify_cached_target(&self.state)?;
1302 let session = self
1303 .state
1304 .sessions
1305 .get(session_id)
1306 .with_context(|| format!("unknown session {session_id}"))?
1307 .clone();
1308 if let Some(mut operation) = crate::database::load_move_operation(session_id)? {
1309 self.cleanup_prepared_move_destination(&mut operation, executor)?;
1310 }
1311 if let Some(locator) = &session.target {
1315 let backend = backend_locator(locator, &session, &self.config)?;
1316 if self.state.subagents.contains_key(session_id) {
1317 targets::borrowed_worker_cleanup_plan(&backend, session_id)?
1318 .execute(executor)?;
1319 } else {
1320 execute_target_cleanup(&backend, &session, &self.config, executor)?;
1321 }
1322 }
1323 if let Some(worktree) = &session.managed_worktree {
1324 cleanup_managed_worktree(executor, worktree, branch)
1325 .context("remove managed raw-session worktree")?;
1326 }
1327 if let Some(checkpoint) = &session.checkpoint
1328 && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
1329 && error.kind() != std::io::ErrorKind::NotFound
1330 {
1331 return Err(error).with_context(|| {
1332 format!(
1333 "remove session recovery archive {}",
1334 checkpoint.archive_path.display()
1335 )
1336 });
1337 }
1338 mj_core::attachment::AttachmentStore::controller(session_id)?
1339 .remove_session_data()
1340 .context("remove session image attachments")?;
1341 delete(session_id).context("force destroy session in database")?;
1342 self.state.subagents.remove(session_id);
1343 self.state.destroy_session_force(session_id)?;
1344 Ok(())
1345 },
1346 )
1347 }
1348}
1349
1350pub fn has_nothing_to_checkpoint(session: &SessionRecord, subagent: bool) -> bool {
1365 match session.state {
1366 SessionState::Provisioning | SessionState::Parked | SessionState::StartupCleanup => true,
1367 SessionState::Error if subagent => true,
1368 SessionState::Closing
1369 | SessionState::Destroying
1370 | SessionState::Error
1371 | SessionState::Lost => session.target.is_none(),
1372 SessionState::Running
1373 | SessionState::Disconnected
1374 | SessionState::Checkpointing
1375 | SessionState::Stopped
1376 | SessionState::DestroyedWithDataLoss => false,
1377 }
1378}
1379
1380fn execute_target_cleanup(
1384 backend: &targets::TargetLocator,
1385 session: &SessionRecord,
1386 config: &mj_core::config::Config,
1387 executor: &impl CommandExecutor,
1388) -> Result<()> {
1389 let session_id = session.id.as_str();
1390 let _owner = crate::worker_lifecycle::require(session_id)?;
1391 if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1392 match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1393 Ok(true) => {
1394 tracing::warn!(
1395 session_id,
1396 error = format!("{cleanup_error:#}"),
1397 "target cleanup command failed, but the target was confirmed absent"
1398 );
1399 }
1400 Ok(false) => {
1401 tracing::error!(
1402 session_id,
1403 error = format!("{cleanup_error:#}"),
1404 "target cleanup failed and the target is still present"
1405 );
1406 return Err(cleanup_error);
1407 }
1408 Err(probe_error) => {
1409 tracing::error!(
1410 session_id,
1411 cleanup_error = format!("{cleanup_error:#}"),
1412 probe_error = format!("{probe_error:#}"),
1413 "target cleanup failed and exact absence could not be confirmed"
1414 );
1415 return Err(cleanup_error.context(format!(
1416 "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1417 )));
1418 }
1419 }
1420 }
1421 if let Some(release) = BuildStateRelease::for_target(session, backend, config) {
1422 release.run(executor);
1423 }
1424 if let Some(intent) = crate::database::load_worker_restart(session_id)? {
1427 crate::database::finish_worker_restart(session_id, &intent.operation_id)?;
1428 }
1429 Ok(())
1430}
1431
1432fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1433 record.state = SessionState::Closing;
1434 record.updated_at = updated_at;
1435 record.last_checkpoint_error = None;
1436}
1437
1438fn apply_close_checkpoint_failure(
1445 record: &mut SessionRecord,
1446 previous: &SessionRecord,
1447 error: &anyhow::Error,
1448 updated_at: String,
1449) {
1450 if WorkerRestartLeftNoWorker::marks(error) {
1451 record.state = SessionState::Error;
1452 record.last_error = Some(format!(
1453 "close failed and left the session without a live worker; retry the close, \
1454 resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1455 ));
1456 } else {
1457 record.state = state_after_unsealed_close(previous);
1458 }
1459 record.last_checkpoint_error = Some(format!("{error:#}"));
1460 record.updated_at = updated_at;
1461}
1462
1463fn close_was_refused(error: &anyhow::Error) -> bool {
1471 error.chain().any(|cause| {
1472 cause
1473 .downcast_ref::<crate::worker_client::RelayRejected>()
1474 .is_some_and(|rejected| !rejected.is_retryable())
1475 })
1476}
1477
1478fn state_after_unsealed_close(previous: &SessionRecord) -> SessionState {
1479 if previous.state == SessionState::Closing {
1480 SessionState::Running
1481 } else {
1482 previous.state
1483 }
1484}
1485
1486fn apply_interrupted_close_error(
1487 record: &mut SessionRecord,
1488 error: &anyhow::Error,
1489 updated_at: &str,
1490) {
1491 let destroying = record.state == SessionState::Destroying;
1492 if !destroying {
1493 record.state = SessionState::Closing;
1494 }
1495 record.updated_at = updated_at.to_owned();
1496 record.last_error = Some(if destroying {
1497 format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1498 } else {
1499 format!("close is safely resumable from its verified checkpoint: {error:#}")
1500 });
1501}
1502
1503#[cfg(test)]
1504mod tests;