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