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