Skip to main content

mj_controller/pollers/
lifecycle.rs

1use super::*;
2
3pub enum LifecycleSuccess {
4    Created,
5    Resumed {
6        profile_id: String,
7        target_id: String,
8    },
9    Moved(mj_core::state::MoveOutcome),
10    Closed,
11    ForceStopped,
12    DestroyedStopped,
13    ForceDestroyed,
14}
15
16pub struct LifecycleUpdate {
17    pub session_id: String,
18    pub result: std::result::Result<LifecycleSuccess, String>,
19    pub deferred_cleanup: bool,
20}
21
22/// Whether a close stopped partway and left the record mid-close with its
23/// target still present. Such a record cannot be closed again from the start:
24/// its worker is gone, so only recovery can finish it.
25pub fn is_interrupted_close(session: &SessionRecord) -> bool {
26    matches!(
27        session.state,
28        SessionState::Closing | SessionState::Destroying
29    ) && session.target.is_some()
30}
31
32pub fn interrupted_suspend_session_ids(controller: &Controller) -> Vec<String> {
33    controller
34        .state
35        .sessions
36        .values()
37        .filter(|session| is_interrupted_close(session))
38        .map(|session| session.id.clone())
39        .collect()
40}
41
42/// Sessions whose destroy failed and stopped partway: `Error`, still holding
43/// their target, with the destruction failure recorded by
44/// `record_lifecycle_failure`. The person asked for them to be destroyed, so
45/// startup finishes the removal (the next `mj destroy` does the same).
46pub fn interrupted_destroy_session_ids(controller: &Controller) -> Vec<String> {
47    controller
48        .state
49        .sessions
50        .values()
51        .filter(|session| {
52            session.state == SessionState::Error
53                && session.target.is_some()
54                && session.last_error.as_deref().is_some_and(|error| {
55                    error.starts_with(mj_core::state::DESTRUCTION_FAILURE_PREFIX)
56                })
57        })
58        .map(|session| session.id.clone())
59        .collect()
60}
61
62/// Why a record left in an in-flight lifecycle state has nobody to finish it,
63/// in words the user reads in `mj sessions` and the TUI.
64///
65/// Every in-flight state needs an owner that will complete it. A durable move
66/// intent owns its session, [`is_interrupted_close`] owns a close or teardown
67/// that still holds its target, and
68/// `database::recover_interrupted_checkpointing_sessions` returns an
69/// interrupted `Checkpointing` record to `Running` before the controller
70/// loads. What is left is a record whose operation died with the process, and
71/// it has to say so instead of waiting forever.
72///
73/// `None` means the state needs no reconciliation; callers exclude the owned
74/// sessions before asking.
75pub fn interrupted_lifecycle_cause(session: &SessionRecord) -> Option<String> {
76    match session.state {
77        // Provisioning has no durable operation behind it. Whatever the dead
78        // provision created is not named by this record, so the resource is
79        // recovered through `mj recover scan`, which can see it again once the
80        // record is no longer in flight.
81        SessionState::Provisioning => Some(
82            "the daemon stopped while this session was provisioning; anything it created \
83             is offered by `mj recover scan`"
84                .to_owned(),
85        ),
86        // An interrupted close or teardown that still holds its target is
87        // resumed rather than failed, so only the target-less residue reaches
88        // here: there is nothing left to tear down, and no relay through which
89        // to finish the close the record claims.
90        SessionState::Closing => Some(
91            "the daemon stopped while this session was closing, and it has no target left \
92             to close"
93                .to_owned(),
94        ),
95        SessionState::Destroying => Some(
96            "the daemon stopped while this session was being torn down, and it has no \
97             target left to remove"
98                .to_owned(),
99        ),
100        // A parked sub-agent is settled: its worker was stopped on purpose and
101        // it waits for its parent's next `send_input`.
102        SessionState::StartupCleanup
103        | SessionState::Checkpointing
104        | SessionState::Running
105        | SessionState::Disconnected
106        | SessionState::Parked
107        | SessionState::Stopped
108        | SessionState::Lost
109        | SessionState::Error
110        | SessionState::DestroyedWithDataLoss => None,
111    }
112}
113
114/// Every session whose in-flight lifecycle state has no owner, with the cause
115/// to record against it. `owned` names the sessions a durable move intent or
116/// another startup recovery has already claimed.
117pub fn unowned_interrupted_lifecycles(
118    controller: &Controller,
119    owned: &std::collections::BTreeSet<String>,
120) -> Vec<(String, String)> {
121    controller
122        .state
123        .sessions
124        .values()
125        .filter(|session| !owned.contains(&session.id) && !is_interrupted_close(session))
126        .filter_map(|session| {
127            interrupted_lifecycle_cause(session).map(|cause| (session.id.clone(), cause))
128        })
129        .collect()
130}
131
132pub fn reserve_recovery_or_cancel(
133    observer: &crate::recovery_gate::RecoveryObserver,
134    session_id: &str,
135    cancelled: &AtomicBool,
136) -> Result<crate::recovery_gate::RecoveryReservation> {
137    let reservation = observer.reserve(session_id);
138    // The reservation stops the next copy; cancelling preempts the one already
139    // running so a lifecycle operation never queues behind a long or wedged
140    // copy.
141    observer.cancel_busy(session_id);
142    while observer.is_busy(session_id) {
143        if cancelled.load(Ordering::Acquire) {
144            bail!("operation cancelled while waiting for recovery copy");
145        }
146        std::thread::sleep(Duration::from_millis(25));
147    }
148    Ok(reservation)
149}
150
151pub fn project_worker_title(
152    controller: &mut Controller,
153    update: &WorkerPollUpdate,
154) -> Option<Option<String>> {
155    let snapshot = update.view.snapshot.as_ref()?;
156    let session = controller.state.sessions.get_mut(&update.session_id)?;
157    let title = snapshot.resolved_title();
158    if session.acp_session_title == title {
159        return None;
160    }
161    session.acp_session_title = title.clone();
162    Some(title)
163}
164
165pub fn apply_worker_record_update(controller: &mut Controller, update: &WorkerPollUpdate) {
166    let Some(title) = project_worker_title(controller, update) else {
167        return;
168    };
169    let session_id = update.session_id.clone();
170    tokio::spawn(async move {
171        let result = tokio::task::spawn_blocking(move || {
172            crate::database::set_session_acp_title(&session_id, title.as_deref())
173        })
174        .await;
175        match result {
176            Ok(Ok(())) => {}
177            Ok(Err(error)) => tracing::warn!(%error, "could not persist relay title"),
178            Err(error) => tracing::warn!(%error, "relay title persistence task failed"),
179        }
180    });
181}