Skip to main content

mj_controller/controller/
lifecycle.rs

1//! Session close, force-stop, and permanent-destruction transitions.
2
3use std::path::PathBuf;
4use std::time::Duration;
5
6use anyhow::{Context, Result, bail, ensure};
7
8use crate::session_manager::{SessionManagerControl, new_command_id};
9use mj_core::state::{CheckpointMetadata, ManagedCheckoutKind, SessionRecord, SessionState};
10
11use crate::targets::{self, CommandExecutor, ProcessExecutor, ProvisionStage, ProvisionStageGuard};
12use mj_core::relay::{RelayCommand, RelayExecutionState};
13
14use super::backend::backend_locator;
15use super::checkpoint::{
16    CheckpointExportPolicy, LatchExclusivity, prune_replaced_checkpoint,
17    release_projection_behind_checkpoint, verify_installed_checkpoint_gate, wait_for_relay_closed,
18};
19use super::worker_restart::WorkerRestartLeftNoWorker;
20use super::worktree::{
21    cleanup_managed_worktree, managed_worktree_checkout_is_dirty, retire_managed_worktree,
22};
23use super::{Controller, now, persist_session_record_transition_or_restore};
24
25/// What destroying a session does with its managed worktree's git branch.
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub enum BranchDisposition {
28    /// Delete the branch with the rest of the session. Only for a branch
29    /// nobody has worked on, or when the user asks for it by name.
30    Delete,
31    /// Leave the branch in the repository. The default for every destroy,
32    /// because the branch may hold work the user still wants.
33    Keep,
34    /// Delete the branch only when every commit on it is reachable from some
35    /// other branch, local or remote-tracking, that is not a session branch.
36    /// Anything else keeps the branch, exactly as [`BranchDisposition::Keep`]
37    /// would. The archive job uses this so a branch whose work has landed
38    /// elsewhere does not pile up forever.
39    DeleteIfMerged,
40}
41
42/// What destroying a session does with its managed worktree's checkout.
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum CheckoutDisposition {
45    /// Remove the checkout whatever it holds. Every destroy a person asks for
46    /// takes this path, because the confirmation already covered the loss.
47    Remove,
48    /// Leave the checkout, and its branch, in place when it has uncommitted
49    /// changes. Destruction Mjolnir decides on its own never discards work the
50    /// user has not seen.
51    KeepWhenDirty,
52}
53
54/// What a verified close does with the target the session was running in.
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub(super) enum SourceTargetDisposition {
57    /// Verified checkpoint, sealed relay, then destroy the exact target.
58    Destroy,
59    /// Verified checkpoint, sealed relay; keep the target and its worker
60    /// daemon alive for an in-place harness replacement.
61    RetainForInPlaceSwap,
62}
63
64impl Controller {
65    /// Checkpoint, ask the harness to close, and only then tear down the exact
66    /// provisioned target. Checkpoint failure is deliberately non-destructive,
67    /// except when the checkpoint's worker restart left no live worker: that
68    /// records `Error` and keeps the target for a later resume or forced close.
69    pub async fn suspend_session(&mut self, session_id: &str) -> Result<()> {
70        self.suspend_session_controlled(session_id, &ProcessExecutor)
71            .await
72    }
73
74    pub async fn suspend_session_controlled(
75        &mut self,
76        session_id: &str,
77        executor: &(impl CommandExecutor + Sync),
78    ) -> Result<()> {
79        if self
80            .suspend_session_controlled_with_manager(
81                session_id,
82                executor,
83                None,
84                None,
85                SourceTargetDisposition::Destroy,
86                true,
87            )
88            .await?
89        {
90            self.cleanup_stopped_target(session_id, executor)?;
91        }
92        Ok(())
93    }
94
95    pub async fn suspend_session_managed_controlled(
96        &mut self,
97        session_id: &str,
98        executor: &(impl CommandExecutor + Sync),
99        manager: &SessionManagerControl,
100        acknowledge_unpublished_work: bool,
101    ) -> Result<bool> {
102        self.suspend_session_controlled_with_manager(
103            session_id,
104            executor,
105            Some(manager),
106            None,
107            SourceTargetDisposition::Destroy,
108            acknowledge_unpublished_work,
109        )
110        .await
111    }
112
113    pub(super) async fn suspend_session_for_move(
114        &mut self,
115        session_id: &str,
116        executor: &(impl CommandExecutor + Sync),
117        manager: &SessionManagerControl,
118        operation: &mut mj_core::state::MoveOperation,
119        preparation: Option<&mj_core::state::MovePreparation>,
120        disposition: SourceTargetDisposition,
121    ) -> Result<bool> {
122        self.prepare_move_source_checkpoint(session_id, executor, manager, operation)
123            .await?;
124        self.suspend_session_controlled_with_manager(
125            session_id,
126            executor,
127            Some(manager),
128            Some((operation, preparation)),
129            disposition,
130            true,
131        )
132        .await
133    }
134
135    async fn suspend_session_controlled_with_manager(
136        &mut self,
137        session_id: &str,
138        executor: &(impl CommandExecutor + Sync),
139        manager: Option<&SessionManagerControl>,
140        move_intent: Option<(
141            &mut mj_core::state::MoveOperation,
142            Option<&mj_core::state::MovePreparation>,
143        )>,
144        disposition: SourceTargetDisposition,
145        acknowledge_unpublished_work: bool,
146    ) -> Result<bool> {
147        let previous = self
148            .state
149            .sessions
150            .get(session_id)
151            .with_context(|| format!("unknown session {session_id}"))?
152            .clone();
153        let record = self.state.sessions.get_mut(session_id).unwrap();
154        // Persist the close intent before beginning its checkpoint. A process
155        // exit anywhere below must leave enough state for the next controller
156        // to retry the close, even when no checkpoint has been installed yet.
157        apply_close_checkpoint_started(record, now());
158        self.persist_session_transition_or_restore(
159            session_id,
160            &previous,
161            "persist closing state before checkpointing the session",
162        )?;
163
164        // Close seals the relay at the exact latched cursor, so this checkpoint
165        // keeps its exclusive connection until the relay reports Closed.
166        let mut latched = match self
167            .checkpoint_session_latched(
168                session_id,
169                executor,
170                manager,
171                LatchExclusivity::HoldThroughClose,
172                CheckpointExportPolicy::ReuseUnchangedArchive,
173            )
174            .await
175        {
176            Ok(latched) => latched,
177            Err(error) => {
178                let record = self.state.sessions.get_mut(session_id).unwrap();
179                // The target is kept even when the restart left no worker: a
180                // forced destroy and a resume's pre-clean both use it to tear
181                // down the dead container.
182                apply_close_checkpoint_failure(record, &previous, &error, now());
183                return Err(
184                    self.persist_failed_checkpoint_state_or_restore(session_id, &previous, error)
185                );
186            }
187        };
188
189        let artifact = latched.artifact.clone();
190        let publication = match previous.managed_worktree.as_ref().map(|owned| owned.kind) {
191            Some(ManagedCheckoutKind::Clone) => Some(super::publication::assess_clone_checkpoint(
192                &previous,
193                &artifact.metadata,
194            )),
195            Some(ManagedCheckoutKind::Worktree) => None,
196            None if previous.project_directory.is_none() => {
197                Some(super::publication::assess_network_checkpoint(
198                    &previous,
199                    &artifact.metadata,
200                    &self.config,
201                ))
202            }
203            None => None,
204        };
205        if !acknowledge_unpublished_work
206            && publication
207                .as_ref()
208                .is_some_and(|result| result.state != mj_core::state::PublicationState::Published)
209        {
210            latched.relay.cancel_abandoned_barrier().await?;
211            self.state
212                .sessions
213                .insert(session_id.to_owned(), previous.clone());
214            self.persist_session_transition_or_restore(
215                session_id,
216                &previous,
217                "restore session after unpublished-work preflight",
218            )?;
219            let _ = std::fs::remove_file(&artifact.metadata.archive_path);
220            return Err(mj_core::refusal::Refusal::precondition(
221                "the checkout has unpublished or unverified work; confirm suspension with acknowledge_unpublished_work=true",
222            ).into());
223        }
224        let record = self.state.sessions.get_mut(session_id).unwrap();
225        record.state = SessionState::Closing;
226        record.native_session_id = Some(artifact.native_session_id.clone());
227        record.checkpoint = Some(artifact.metadata.clone());
228        record.publication = publication;
229        record.updated_at = now();
230        record.last_error = None;
231        record.last_checkpoint_error = None;
232        self.persist_checkpoint_transition_or_restore(
233            session_id,
234            &previous,
235            "persist verified checkpoint and closing state before sealing the relay",
236        )?;
237        if let Some((operation, preparation)) = move_intent {
238            // The source is still behind an unsealed barrier. A destination
239            // preflight error must release it and leave its processes alive.
240            if let Err(error) = self.validate_move_checkpoint(operation, preparation, executor) {
241                let record = self.state.sessions.get_mut(session_id).unwrap();
242                record.state = previous.state;
243                record.last_error = Some(format!("{error:#}"));
244                self.persist_session_transition_or_restore(
245                    session_id,
246                    &previous,
247                    "restore source after move preflight failure",
248                )?;
249                return Err(error);
250            }
251            operation.checkpoint = Some(artifact.metadata.clone());
252            operation.updated_at = now();
253            crate::database::save_move_operation(operation)?;
254        }
255        prune_replaced_checkpoint(previous.checkpoint.as_ref(), &artifact.metadata);
256        // A stopping session will not checkpoint again, so this is its last
257        // chance to release what its checkpoint now covers.
258        release_projection_behind_checkpoint(session_id, &artifact.metadata);
259
260        let close_command_id = new_command_id("close")?;
261        let barrier_command_id = latched.barrier_command_id.clone();
262        let close_result = {
263            let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
264            latched
265                .relay
266                .connection_mut()
267                .submit(
268                    close_command_id,
269                    RelayCommand::Close {
270                        barrier_command_id: barrier_command_id.clone(),
271                        expected: latched.cursor.clone(),
272                    },
273                )
274                .await
275        };
276        if let Err(error) = close_result {
277            self.record_interrupted_close(session_id, &error)?;
278            return Err(error.context("seal verified checkpoint for close"));
279        }
280        let close_result = {
281            let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
282            latched
283                .relay
284                .connection_mut()
285                .submit(
286                    new_command_id("checkpoint-complete")?,
287                    RelayCommand::CompleteCheckpoint { barrier_command_id },
288                )
289                .await
290        };
291        if let Err(error) = close_result {
292            self.record_interrupted_close(session_id, &error)?;
293            return Err(error.context("release verified close checkpoint"));
294        }
295        let close_result = {
296            let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
297            wait_for_relay_closed(latched.relay.connection_mut()).await
298        };
299        if let Err(error) = close_result {
300            self.record_interrupted_close(session_id, &error)?;
301            return Err(error);
302        }
303        latched.relay.release();
304
305        if disposition == SourceTargetDisposition::RetainForInPlaceSwap {
306            // The record stays `Closing` with its verified checkpoint and its
307            // target. The worker daemon is deliberately left running: a crash
308            // between here and the in-place restore recovers through
309            // `recover_interrupted_close_managed`, which needs the daemon to
310            // answer `Closed`.
311            return Ok(false);
312        }
313        match self.destroy_after_verified_checkpoint(session_id, &artifact.metadata, executor) {
314            Ok(deferred) => Ok(deferred),
315            Err(error) => {
316                self.record_interrupted_close(session_id, &error)?;
317                Err(error)
318            }
319        }
320    }
321
322    /// Resume the durable closing state after a controller restart. If the
323    /// relay had accepted Close, wait for it and destroy through the exact
324    /// installed checkpoint gate. If it had not, take a fresh checkpoint;
325    /// the previously installed archive may have become stale after EOF
326    /// released its barrier.
327    pub async fn recover_interrupted_close_managed(
328        &mut self,
329        session_id: &str,
330        executor: &(impl CommandExecutor + Sync),
331        manager: &SessionManagerControl,
332    ) -> Result<bool> {
333        let (state, verified) = {
334            let session = self
335                .state
336                .sessions
337                .get(session_id)
338                .with_context(|| format!("unknown session {session_id}"))?;
339            ensure!(
340                matches!(
341                    session.state,
342                    SessionState::Closing | SessionState::Destroying
343                ),
344                "session {session_id} has no interrupted close to recover"
345            );
346            (session.state, session.checkpoint.clone())
347        };
348        if state == SessionState::Destroying {
349            let verified = verified.context("destroying session has no verified checkpoint")?;
350            return self.destroy_after_verified_checkpoint(session_id, &verified, executor);
351        }
352        ensure!(
353            state == SessionState::Closing,
354            "session {session_id} has no relay close to recover"
355        );
356        let handle = manager
357            .wait_for_session(session_id, Duration::from_secs(5))
358            .await?;
359        let mut lease = handle.lease_connection().await?;
360        let execution = lease.connection_mut().sync().await?.operational.execution;
361        match execution {
362            RelayExecutionState::Closed => {}
363            RelayExecutionState::Closing => {
364                let _closing = ProvisionStageGuard::new(executor, ProvisionStage::Closing);
365                wait_for_relay_closed(lease.connection_mut()).await?;
366            }
367            RelayExecutionState::Idle | RelayExecutionState::Running => {
368                lease.release();
369                return self
370                    .suspend_session_controlled_with_manager(
371                        session_id,
372                        executor,
373                        Some(manager),
374                        None,
375                        SourceTargetDisposition::Destroy,
376                        true,
377                    )
378                    .await;
379            }
380        }
381        lease.release();
382        let verified = verified.context("closed relay has no verified checkpoint")?;
383        self.destroy_after_verified_checkpoint(session_id, &verified, executor)
384    }
385
386    /// Record that an in-flight lifecycle state has no operation left to
387    /// finish it, so the session stops waiting for one.
388    ///
389    /// The state is re-checked against the freshly loaded record, because the
390    /// caller decided what to reconcile from a startup snapshot. Returns
391    /// whether anything changed.
392    pub fn fail_interrupted_lifecycle(&mut self, session_id: &str, cause: &str) -> Result<bool> {
393        self.fail_interrupted_lifecycle_with(
394            session_id,
395            cause,
396            crate::database::save_lifecycle_session,
397        )
398    }
399
400    fn fail_interrupted_lifecycle_with(
401        &mut self,
402        session_id: &str,
403        cause: &str,
404        persist: impl Fn(&SessionRecord) -> Result<()>,
405    ) -> Result<bool> {
406        let Some(record) = self.state.sessions.get_mut(session_id) else {
407            return Ok(false);
408        };
409        if crate::pollers::interrupted_lifecycle_cause(record).is_none() {
410            return Ok(false);
411        }
412        let previous = record.clone();
413        record.state = SessionState::Error;
414        record.updated_at = now();
415        record.last_error = Some(cause.to_owned());
416        persist_session_record_transition_or_restore(
417            &mut self.state,
418            session_id,
419            &previous,
420            "persist the failure of an interrupted lifecycle state",
421            &persist,
422        )?;
423        Ok(true)
424    }
425
426    /// Record that a live session's harness never became usable, so a driver
427    /// can close it and provision a replacement instead of waiting on a
428    /// session that will never take a prompt (#1090).
429    ///
430    /// Only a record that still looks exactly as the daemon observed it is
431    /// failed: a lifecycle operation that started after the observation owns
432    /// the session, and its own outcome must not be overwritten by this one.
433    /// Returns whether anything changed.
434    pub fn fail_unready_session(
435        &mut self,
436        session_id: &str,
437        cause: &str,
438        observed_updated_at: &str,
439    ) -> Result<bool> {
440        self.fail_unready_session_with(
441            session_id,
442            cause,
443            observed_updated_at,
444            crate::database::save_lifecycle_session,
445        )
446    }
447
448    fn fail_unready_session_with(
449        &mut self,
450        session_id: &str,
451        cause: &str,
452        observed_updated_at: &str,
453        persist: impl Fn(&SessionRecord) -> Result<()>,
454    ) -> Result<bool> {
455        let Some(record) = self.state.sessions.get_mut(session_id) else {
456            return Ok(false);
457        };
458        if record.state != SessionState::Running || record.updated_at != observed_updated_at {
459            return Ok(false);
460        }
461        let previous = record.clone();
462        record.state = SessionState::Error;
463        record.updated_at = now();
464        record.last_error = Some(cause.to_owned());
465        persist_session_record_transition_or_restore(
466            &mut self.state,
467            session_id,
468            &previous,
469            "persist a session whose harness never became usable",
470            &persist,
471        )?;
472        Ok(true)
473    }
474
475    /// Publish a safe lifecycle failure even if a previous internal error exists.
476    /// Detailed diagnostics belong in logs; the stored outcome must be visible
477    /// on all control surfaces. Reports whether the session still exists.
478    pub fn record_failed_close(&mut self, session_id: &str, cause: &str) -> Result<bool> {
479        let Some(record) = self.state.sessions.get(session_id) else {
480            return Ok(false);
481        };
482        // The supervisor logs the detailed failure. Always publish this safe
483        // outcome, including when an older raw error was already recorded.
484        let previous = record.clone();
485        let record = self.state.sessions.get_mut(session_id).unwrap();
486        record.last_error = Some(cause.to_owned());
487        record.updated_at = now();
488        persist_session_record_transition_or_restore(
489            &mut self.state,
490            session_id,
491            &previous,
492            "persist the reason a close did not finish",
493            &crate::database::save_lifecycle_session,
494        )?;
495        Ok(true)
496    }
497
498    /// Forget a recorded close failure, because something for this session has
499    /// since succeeded. Only the sentence a failed close wrote is cleared; a
500    /// raw error from any other operation is left alone. Reports whether
501    /// anything changed.
502    pub fn clear_recorded_close_failure(&mut self, session_id: &str) -> Result<bool> {
503        let Some(record) = self.state.sessions.get(session_id) else {
504            return Ok(false);
505        };
506        if record.public_error().is_none() {
507            return Ok(false);
508        }
509        let previous = record.clone();
510        let record = self.state.sessions.get_mut(session_id).unwrap();
511        record.last_error = None;
512        record.updated_at = now();
513        persist_session_record_transition_or_restore(
514            &mut self.state,
515            session_id,
516            &previous,
517            "clear the reason a close did not finish",
518            &crate::database::save_lifecycle_session,
519        )?;
520        Ok(true)
521    }
522
523    /// Close a session that has nothing to checkpoint.
524    ///
525    /// A record still provisioning never reached a running worker, and a
526    /// record with no target locator has no target to read a workspace from,
527    /// so in both cases there is no relay to latch and no harness state to
528    /// archive. Waiting for a relay that does not exist is what left a stuck
529    /// provisioning session unclosable. Any target the session did leave
530    /// behind is still torn down, and the checkpoint it already had is kept,
531    /// so this is a close, not a forced destroy.
532    ///
533    /// Returns whether target storage cleanup was deferred, like the graceful
534    /// close does.
535    pub fn suspend_session_without_checkpoint(
536        &mut self,
537        session_id: &str,
538        executor: &impl CommandExecutor,
539    ) -> Result<bool> {
540        self.suspend_session_without_checkpoint_with(
541            session_id,
542            executor,
543            crate::database::save_lifecycle_session,
544        )
545    }
546
547    fn suspend_session_without_checkpoint_with(
548        &mut self,
549        session_id: &str,
550        executor: &impl CommandExecutor,
551        persist: impl Fn(&SessionRecord) -> Result<()>,
552    ) -> Result<bool> {
553        let session = self
554            .state
555            .sessions
556            .get(session_id)
557            .with_context(|| format!("unknown session {session_id}"))?
558            .clone();
559        ensure!(
560            has_nothing_to_checkpoint(&session),
561            "session {session_id} has a workspace to checkpoint; suspend it instead"
562        );
563        self.stop_target_and_settle(session_id, &session, executor, &persist)
564    }
565
566    fn record_interrupted_close(&mut self, session_id: &str, error: &anyhow::Error) -> Result<()> {
567        let record = self.state.sessions.get_mut(session_id).unwrap();
568        apply_interrupted_close_error(record, error, &now());
569        self.persist_session_state(session_id)
570    }
571
572    /// Execute cleanup only after the close state machine has installed a
573    /// verified checkpoint on the record.
574    fn destroy_after_verified_checkpoint(
575        &mut self,
576        session_id: &str,
577        verified: &CheckpointMetadata,
578        executor: &impl CommandExecutor,
579    ) -> Result<bool> {
580        self.destroy_after_verified_checkpoint_with(
581            session_id,
582            verified,
583            executor,
584            crate::database::save_lifecycle_session,
585        )
586    }
587
588    fn destroy_after_verified_checkpoint_with(
589        &mut self,
590        session_id: &str,
591        verified: &CheckpointMetadata,
592        executor: &impl CommandExecutor,
593        persist: impl Fn(&SessionRecord) -> Result<()>,
594    ) -> Result<bool> {
595        let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
596        let _target_guard = target_mutex.lock().map_err(|_| {
597            anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
598        })?;
599        let session = self
600            .state
601            .sessions
602            .get(session_id)
603            .with_context(|| format!("unknown session {session_id}"))?
604            .clone();
605        ensure!(
606            matches!(
607                session.state,
608                SessionState::Closing | SessionState::Destroying
609            ),
610            "refusing to destroy session {session_id}: it is not closing or destroying"
611        );
612        ensure!(
613            session.checkpoint.as_ref() == Some(verified),
614            "refusing to destroy session {session_id}: verified checkpoint gate is stale"
615        );
616        if session.state == SessionState::Closing {
617            let record = self.state.sessions.get_mut(session_id).unwrap();
618            record.state = SessionState::Destroying;
619            record.updated_at = now();
620            record.last_error = None;
621            persist_session_record_transition_or_restore(
622                &mut self.state,
623                session_id,
624                &session,
625                "persist destroying state before target cleanup",
626                &persist,
627            )?;
628        }
629
630        let destroying = self
631            .state
632            .sessions
633            .get(session_id)
634            .expect("destroying session disappeared")
635            .clone();
636        {
637            let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
638            verify_installed_checkpoint_gate(session_id, verified)?;
639        }
640        // The reviewer's native session lives on the target that is about to
641        // go. Recording that now, before the target is torn down, is what
642        // stops a resumed session from trying to reload a conversation that no
643        // longer exists; its transcript is kept for reference either way.
644        if let Err(error) = crate::database::lose_reviewer_continuity(session_id) {
645            tracing::warn!(
646                session_id,
647                error = format!("{error:#}"),
648                "could not record that the second-opinion conversation ends with this target"
649            );
650        }
651        let locator = destroying
652            .target
653            .as_ref()
654            .context("session has no target")?;
655        let backend = backend_locator(locator, &destroying, &self.config)?;
656        let deferred = if self.state.subagents.contains_key(session_id) {
657            targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
658            false
659        } else if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
660            plan.execute(executor)?;
661            true
662        } else {
663            execute_target_cleanup(&backend, session_id, executor)?;
664            false
665        };
666        if let Some(worktree) = &destroying.managed_worktree {
667            retire_managed_worktree(executor, worktree)
668                .context("retire managed raw-session worktree after verified close")?;
669        }
670        let record = self.state.sessions.get_mut(session_id).unwrap();
671        record.state = SessionState::Stopped;
672        if !deferred {
673            record.target = None;
674        }
675        record.updated_at = now();
676        record.last_error = None;
677        persist_session_record_transition_or_restore(
678            &mut self.state,
679            session_id,
680            &destroying,
681            "persist stopped state after target cleanup",
682            &persist,
683        )?;
684        Ok(deferred)
685    }
686
687    /// Finish storage cleanup for a stopped Podman target retained by the
688    /// quiescence transition. The locator stays durable until every command
689    /// succeeds, making daemon restart and explicit retry idempotent.
690    pub fn cleanup_stopped_target(
691        &mut self,
692        session_id: &str,
693        executor: &impl CommandExecutor,
694    ) -> Result<()> {
695        self.cleanup_stopped_target_with(
696            session_id,
697            executor,
698            crate::database::save_lifecycle_session,
699        )
700    }
701
702    fn cleanup_stopped_target_with(
703        &mut self,
704        session_id: &str,
705        executor: &impl CommandExecutor,
706        persist: impl Fn(&SessionRecord) -> Result<()>,
707    ) -> Result<()> {
708        let target_mutex = crate::recovery_gate::worker_target_mutex(session_id);
709        let _target_guard = target_mutex.lock().map_err(|_| {
710            anyhow::anyhow!("worker target ownership lock poisoned for {session_id}")
711        })?;
712        let previous = self
713            .state
714            .sessions
715            .get(session_id)
716            .with_context(|| format!("unknown session {session_id}"))?
717            .clone();
718        ensure!(
719            previous.state == SessionState::Stopped,
720            "refusing deferred cleanup for active session {session_id}"
721        );
722        let Some(locator) = previous.target.as_ref() else {
723            return Ok(());
724        };
725        let backend = backend_locator(locator, &previous, &self.config)?;
726        ensure!(
727            targets::quiesce_plan(&backend, session_id)?.is_some(),
728            "session {session_id} retained a non-Podman target after stopping"
729        );
730        if let Err(error) = execute_target_cleanup(&backend, session_id, executor) {
731            let record = self.state.sessions.get_mut(session_id).unwrap();
732            record.updated_at = now();
733            record.last_error = Some(format!("deferred target cleanup failed: {error:#}"));
734            let persisted = persist_session_record_transition_or_restore(
735                &mut self.state,
736                session_id,
737                &previous,
738                "persist deferred target cleanup failure",
739                &persist,
740            );
741            return match persisted {
742                Ok(()) => Err(error),
743                Err(persist_error) => Err(error.context(format!(
744                    "also failed to persist deferred target cleanup failure: {persist_error:#}"
745                ))),
746            };
747        }
748        let record = self.state.sessions.get_mut(session_id).unwrap();
749        record.target = None;
750        record.updated_at = now();
751        record.last_error = None;
752        persist_session_record_transition_or_restore(
753            &mut self.state,
754            session_id,
755            &previous,
756            "persist completion of deferred Podman target cleanup",
757            &persist,
758        )
759    }
760
761    /// Tear down the current target without taking a fresh checkpoint, then
762    /// leave the logical session resumable from its latest verified archive.
763    pub fn force_stop(
764        &mut self,
765        session_id: &str,
766        executor: &impl CommandExecutor,
767    ) -> Result<bool> {
768        self.force_stop_with(
769            session_id,
770            executor,
771            crate::database::save_lifecycle_session,
772        )
773    }
774
775    fn force_stop_with(
776        &mut self,
777        session_id: &str,
778        executor: &impl CommandExecutor,
779        persist: impl Fn(&SessionRecord) -> Result<()>,
780    ) -> Result<bool> {
781        let session = self
782            .state
783            .sessions
784            .get(session_id)
785            .with_context(|| format!("unknown session {session_id}"))?
786            .clone();
787        ensure!(
788            session.state.is_active(),
789            "session {session_id} is already inactive"
790        );
791        let checkpoint = session
792            .checkpoint
793            .as_ref()
794            .context("force stop requires an existing recovery archive")?;
795        // Force stop skips a new checkpoint, never the checksum gate on the
796        // archive that makes the logical session resumable afterwards.
797        {
798            let _verifying = ProvisionStageGuard::new(executor, ProvisionStage::Verifying);
799            verify_installed_checkpoint_gate(session_id, checkpoint)
800                .context("verify the recovery archive before force stopping")?;
801        }
802        self.stop_target_and_settle(session_id, &session, executor, &persist)
803    }
804
805    /// Tear down whatever target the session holds and settle its record in
806    /// `Stopped`. Shared by force stop and by a close that has nothing to
807    /// checkpoint; neither takes a fresh archive, so neither may decide on its
808    /// own whether the session still has one.
809    fn stop_target_and_settle(
810        &mut self,
811        session_id: &str,
812        session: &SessionRecord,
813        executor: &impl CommandExecutor,
814        persist: &impl Fn(&SessionRecord) -> Result<()>,
815    ) -> Result<bool> {
816        let mut deferred = false;
817        if let Some(locator) = &session.target {
818            let backend = backend_locator(locator, session, &self.config)?;
819            if let Some(plan) = targets::quiesce_plan(&backend, session_id)? {
820                plan.execute(executor)?;
821                deferred = true;
822            } else {
823                execute_target_cleanup(&backend, session_id, executor)?;
824            }
825        }
826        if let Some(worktree) = &session.managed_worktree {
827            retire_managed_worktree(executor, worktree)
828                .context("retire managed raw-session worktree after stopping the target")?;
829        }
830        let record = self.state.sessions.get_mut(session_id).unwrap();
831        record.state = SessionState::Stopped;
832        if !deferred {
833            record.target = None;
834        }
835        record.updated_at = now();
836        record.last_error = None;
837        record.last_checkpoint_error = None;
838        persist_session_record_transition_or_restore(
839            &mut self.state,
840            session_id,
841            session,
842            "persist stopped state after tearing down the current target",
843            persist,
844        )?;
845        Ok(deferred)
846    }
847
848    /// Permanently destroy an inactive session and every artifact Hel owns for it.
849    /// External cleanup happens before the durable record is dropped so failures
850    /// remain visible and retryable.
851    pub fn destroy_session_controlled(
852        &mut self,
853        session_id: &str,
854        executor: &impl CommandExecutor,
855    ) -> Result<()> {
856        self.destroy_session_controlled_with(session_id, executor, BranchDisposition::Keep)
857    }
858
859    /// The same, with a say in what happens to the managed worktree's branch.
860    ///
861    /// A session that is already `Stopped` has had its checkout removed by
862    /// [`retire_managed_worktree`], so usually only the branch is left for
863    /// [`cleanup_managed_worktree`] to take. With [`BranchDisposition::Keep`]
864    /// the record, the checkpoint, and the attachments go and the branch
865    /// stays, which is what a destroy does unless the user asks otherwise.
866    pub fn destroy_session_controlled_with(
867        &mut self,
868        session_id: &str,
869        executor: &impl CommandExecutor,
870        branch: BranchDisposition,
871    ) -> Result<()> {
872        self.destroy_session_controlled_with_checkout(
873            session_id,
874            executor,
875            branch,
876            CheckoutDisposition::Remove,
877        )
878        .map(|_| ())
879    }
880
881    /// The same, with a say in whether a checkout holding uncommitted changes
882    /// survives the destruction.
883    ///
884    /// Answers with the checkout path that was kept, so the caller can name it
885    /// where the person will see it.
886    pub(crate) fn destroy_session_controlled_with_checkout(
887        &mut self,
888        session_id: &str,
889        executor: &impl CommandExecutor,
890        branch: BranchDisposition,
891        checkout: CheckoutDisposition,
892    ) -> Result<Option<PathBuf>> {
893        let session = self
894            .state
895            .sessions
896            .get(session_id)
897            .with_context(|| format!("unknown session {session_id}"))?
898            .clone();
899        if session.state.is_active() {
900            bail!("refusing to destroy active session {session_id}");
901        }
902        let mut retained_checkout = None;
903        if let Some(worktree) = &session.managed_worktree {
904            let keep = checkout == CheckoutDisposition::KeepWhenDirty
905                && (worktree.kind == mj_core::state::ManagedCheckoutKind::Clone
906                    || managed_worktree_checkout_is_dirty(executor, worktree).context(
907                        "check the managed raw-session worktree for uncommitted changes",
908                    )?);
909            if keep {
910                // Keeping the checkout keeps its branch with it: the commits
911                // the working tree is based on are the only way back to this
912                // work.
913                retained_checkout = Some(worktree.worktree_root.clone());
914            } else {
915                cleanup_managed_worktree(executor, worktree, branch)
916                    .context("remove managed raw-session worktree")?;
917            }
918        }
919        if let Some(checkpoint) = &session.checkpoint
920            && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
921            && error.kind() != std::io::ErrorKind::NotFound
922        {
923            return Err(error).with_context(|| {
924                format!(
925                    "remove session recovery archive {}",
926                    checkpoint.archive_path.display()
927                )
928            });
929        }
930        mj_core::attachment::AttachmentStore::controller(session_id)?
931            .remove_session_data()
932            .context("remove session image attachments")?;
933        crate::database::delete_session(session_id)
934            .context("destroy stopped session in database")?;
935        self.state.subagents.remove(session_id);
936        self.state.destroy_stopped_session(session_id)?;
937        Ok(retained_checkout)
938    }
939
940    /// Permanently destroy a session from any state, without checkpointing
941    /// and without requiring a recovery archive.
942    ///
943    /// Unlike [`Controller::destroy_session_controlled`], this accepts active
944    /// states: it tears the live target down with the same close plan a
945    /// verified close uses, so the owning process group dies before any files
946    /// go. External cleanup happens before the durable record is dropped so
947    /// failures stay visible and retryable; the recovery archive is removed,
948    /// which is what makes the destruction irreversible. The managed
949    /// worktree's checkout always goes; its branch goes only when `branch`
950    /// says so.
951    pub fn force_destroy_session(
952        &mut self,
953        session_id: &str,
954        executor: &impl CommandExecutor,
955        branch: BranchDisposition,
956    ) -> Result<()> {
957        self.force_destroy_session_with(
958            session_id,
959            executor,
960            branch,
961            crate::database::delete_session,
962        )
963    }
964
965    fn force_destroy_session_with(
966        &mut self,
967        session_id: &str,
968        executor: &impl CommandExecutor,
969        branch: BranchDisposition,
970        delete: impl Fn(&str) -> Result<()>,
971    ) -> Result<()> {
972        let session = self
973            .state
974            .sessions
975            .get(session_id)
976            .with_context(|| format!("unknown session {session_id}"))?
977            .clone();
978        // A session destroyed for good keeps nothing, including a broker an
979        // earlier failure left running; retiring it first also stops a live
980        // writer from recreating files under the teardown below.
981        if let Some(locator) = &session.target {
982            let backend = backend_locator(locator, &session, &self.config)?;
983            if self.state.subagents.contains_key(session_id) {
984                targets::borrowed_worker_cleanup_plan(&backend, session_id)?.execute(executor)?;
985            } else {
986                execute_target_cleanup(&backend, session_id, executor)?;
987            }
988        }
989        if let Some(worktree) = &session.managed_worktree {
990            cleanup_managed_worktree(executor, worktree, branch)
991                .context("remove managed raw-session worktree")?;
992        }
993        if let Some(checkpoint) = &session.checkpoint
994            && let Err(error) = std::fs::remove_file(&checkpoint.archive_path)
995            && error.kind() != std::io::ErrorKind::NotFound
996        {
997            return Err(error).with_context(|| {
998                format!(
999                    "remove session recovery archive {}",
1000                    checkpoint.archive_path.display()
1001                )
1002            });
1003        }
1004        mj_core::attachment::AttachmentStore::controller(session_id)?
1005            .remove_session_data()
1006            .context("remove session image attachments")?;
1007        delete(session_id).context("force destroy session in database")?;
1008        self.state.subagents.remove(session_id);
1009        self.state.destroy_session_force(session_id)?;
1010        Ok(())
1011    }
1012}
1013
1014/// Whether a close of this session has no workspace to archive.
1015///
1016/// Only a session that cannot be holding live work qualifies. A record still
1017/// `Provisioning` has never had a worker connected, so there is no relay to
1018/// latch and no harness state to capture. A record mid-close or already
1019/// failed, with no target locator left, has no target to read a workspace
1020/// from at all. Every other state may hold work and must take the graceful
1021/// close's checkpoint.
1022pub fn has_nothing_to_checkpoint(session: &SessionRecord) -> bool {
1023    match session.state {
1024        SessionState::Provisioning => true,
1025        SessionState::Closing
1026        | SessionState::Destroying
1027        | SessionState::Error
1028        | SessionState::Lost => session.target.is_none(),
1029        SessionState::Running
1030        | SessionState::Disconnected
1031        | SessionState::Checkpointing
1032        | SessionState::Stopped
1033        | SessionState::DestroyedWithDataLoss => false,
1034    }
1035}
1036
1037fn execute_target_cleanup(
1038    backend: &targets::TargetLocator,
1039    session_id: &str,
1040    executor: &impl CommandExecutor,
1041) -> Result<()> {
1042    if let Err(cleanup_error) = targets::close_plan(backend, session_id)?.execute(executor) {
1043        match targets::cleanup_target_is_confirmed_absent(backend, session_id, executor) {
1044            Ok(true) => {
1045                tracing::warn!(
1046                    session_id,
1047                    error = format!("{cleanup_error:#}"),
1048                    "target cleanup command failed, but the target was confirmed absent"
1049                );
1050            }
1051            Ok(false) => {
1052                tracing::error!(
1053                    session_id,
1054                    error = format!("{cleanup_error:#}"),
1055                    "target cleanup failed and the target is still present"
1056                );
1057                return Err(cleanup_error);
1058            }
1059            Err(probe_error) => {
1060                tracing::error!(
1061                    session_id,
1062                    cleanup_error = format!("{cleanup_error:#}"),
1063                    probe_error = format!("{probe_error:#}"),
1064                    "target cleanup failed and exact absence could not be confirmed"
1065                );
1066                return Err(cleanup_error.context(format!(
1067                    "target cleanup failed and exact absence could not be confirmed: {probe_error:#}"
1068                )));
1069            }
1070        }
1071    }
1072    Ok(())
1073}
1074
1075fn apply_close_checkpoint_started(record: &mut SessionRecord, updated_at: String) {
1076    record.state = SessionState::Closing;
1077    record.updated_at = updated_at;
1078    record.last_checkpoint_error = None;
1079}
1080
1081/// Record a close whose checkpoint failed.
1082///
1083/// An ordinary failure is non-destructive: the session returns to the state it
1084/// had. A restart that left no live worker cannot return to Running, because
1085/// nothing is listening there any more; it records `Error` so the session stops
1086/// being polled, and keeps its target for a later resume or forced close.
1087fn apply_close_checkpoint_failure(
1088    record: &mut SessionRecord,
1089    previous: &SessionRecord,
1090    error: &anyhow::Error,
1091    updated_at: String,
1092) {
1093    if WorkerRestartLeftNoWorker::marks(error) {
1094        record.state = SessionState::Error;
1095        record.last_error = Some(format!(
1096            "close failed and left the session without a live worker; retry the close, \
1097             resume from its checkpoint, or explicitly destroy it with mj destroy: {error:#}"
1098        ));
1099    } else {
1100        record.state = if previous.state == SessionState::Closing {
1101            SessionState::Running
1102        } else {
1103            previous.state
1104        };
1105    }
1106    record.last_checkpoint_error = Some(format!("{error:#}"));
1107    record.updated_at = updated_at;
1108}
1109
1110fn apply_interrupted_close_error(
1111    record: &mut SessionRecord,
1112    error: &anyhow::Error,
1113    updated_at: &str,
1114) {
1115    let destroying = record.state == SessionState::Destroying;
1116    if !destroying {
1117        record.state = SessionState::Closing;
1118    }
1119    record.updated_at = updated_at.to_owned();
1120    record.last_error = Some(if destroying {
1121        format!("target cleanup is safely retryable from its verified checkpoint: {error:#}")
1122    } else {
1123        format!("close is safely resumable from its verified checkpoint: {error:#}")
1124    });
1125}
1126
1127#[cfg(test)]
1128mod tests;