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_close_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
42pub fn spawn_interrupted_close_recovery(
43    session_id: String,
44    session_manager: SessionManagerControl,
45    recovery_observer: crate::recovery_gate::RecoveryObserver,
46    cancelled: Arc<AtomicBool>,
47    updates: tokio::sync::mpsc::UnboundedSender<LifecycleUpdate>,
48    tracker: Option<mj_client::operations::CriticalOperationTracker>,
49) -> tokio::task::JoinHandle<()> {
50    let guard = tracker.map(|tracker| {
51        tracker.begin_cancellable(
52            format!(
53                "recovering session {}",
54                mj_core::state::short_id(&session_id)
55            ),
56            cancelled.clone(),
57        )
58    });
59    let runtime = tokio::runtime::Handle::current();
60    tokio::spawn(async move {
61        let operation_session_id = session_id.clone();
62        let joined = tokio::task::spawn_blocking(move || {
63            (|| -> Result<bool> {
64                let _recovery_reservation = reserve_recovery_or_cancel(
65                    &recovery_observer,
66                    &operation_session_id,
67                    &cancelled,
68                )?;
69                let mut controller = Controller::load()?;
70                let executor = CancellableProcessExecutor::new(cancelled);
71                runtime.block_on(controller.recover_interrupted_close_managed(
72                    &operation_session_id,
73                    &executor,
74                    &session_manager,
75                ))
76            })()
77            .map_err(|error| format!("{error:#}"))
78        })
79        .await;
80        let (result, deferred_cleanup) = match joined {
81            Ok(Ok(deferred_cleanup)) => (Ok(LifecycleSuccess::Closed), deferred_cleanup),
82            Ok(Err(error)) => (Err(error), false),
83            Err(error) => (
84                Err(format!("interrupted close recovery task failed: {error}")),
85                false,
86            ),
87        };
88        if let Err(error) = updates.send(LifecycleUpdate {
89            session_id: session_id.clone(),
90            result,
91            deferred_cleanup,
92        }) {
93            tracing::debug!(%session_id, %error, "interrupted close result dropped after dashboard shutdown");
94        }
95        drop(guard);
96    })
97}
98
99pub fn reserve_recovery_or_cancel(
100    observer: &crate::recovery_gate::RecoveryObserver,
101    session_id: &str,
102    cancelled: &AtomicBool,
103) -> Result<crate::recovery_gate::RecoveryReservation> {
104    let reservation = observer.reserve(session_id);
105    // The reservation stops the next copy; cancelling preempts the one already
106    // running so a lifecycle operation never queues behind a long or wedged
107    // copy.
108    observer.cancel_busy(session_id);
109    while observer.is_busy(session_id) {
110        if cancelled.load(Ordering::Acquire) {
111            bail!("operation cancelled while waiting for recovery copy");
112        }
113        std::thread::sleep(Duration::from_millis(25));
114    }
115    Ok(reservation)
116}
117
118pub fn project_worker_title(
119    controller: &mut Controller,
120    update: &WorkerPollUpdate,
121) -> Option<Option<String>> {
122    let snapshot = update.view.snapshot.as_ref()?;
123    let session = controller.state.sessions.get_mut(&update.session_id)?;
124    let title = snapshot.resolved_title();
125    if session.acp_session_title == title {
126        return None;
127    }
128    session.acp_session_title = title.clone();
129    Some(title)
130}
131
132pub fn apply_worker_record_update(controller: &mut Controller, update: &WorkerPollUpdate) {
133    let Some(title) = project_worker_title(controller, update) else {
134        return;
135    };
136    let session_id = update.session_id.clone();
137    tokio::spawn(async move {
138        let result = tokio::task::spawn_blocking(move || {
139            crate::database::set_session_acp_title(&session_id, title.as_deref())
140        })
141        .await;
142        match result {
143            Ok(Ok(())) => {}
144            Ok(Err(error)) => tracing::warn!(%error, "could not persist relay title"),
145            Err(error) => tracing::warn!(%error, "relay title persistence task failed"),
146        }
147    });
148}