Skip to main content

mj_controller/daemon/
views.rs

1use super::*;
2
3impl RuntimeState {
4    pub(super) fn cancel_lifecycle(&self, session_id: &str) -> Result<()> {
5        // Controller before lifecycle, the order worker polling takes.
6        let controller = self
7            .controller
8            .lock()
9            .unwrap_or_else(PoisonError::into_inner);
10        let lifecycle = self
11            .lifecycle
12            .lock()
13            .unwrap_or_else(PoisonError::into_inner);
14        let active = lifecycle.get(session_id).with_context(|| {
15            format!("no lifecycle operation is running for session {session_id}")
16        })?;
17        ensure!(
18            lifecycle_cancellable(active.kind, durable_session_state(&controller, session_id)),
19            "stop of {session_id} has passed its verified checkpoint and is removing the target; \
20             it cannot be cancelled"
21        );
22        ensure!(
23            active.request_cancel(),
24            "lifecycle operation is no longer cancellable"
25        );
26        drop(lifecycle);
27        drop(controller);
28        self.publish_revision();
29        Ok(())
30    }
31
32    /// Let storage cleanup drain briefly, then cancel and join every lifecycle
33    /// owner before the daemon closes its session manager and database writer.
34    /// One shared deadline bounds all cleanup tasks rather than granting eight
35    /// seconds to each session serially.
36    pub(super) async fn cancel_and_wait_lifecycles(&self) -> Result<()> {
37        let mut pending = {
38            let lifecycle = self
39                .lifecycle
40                .lock()
41                .unwrap_or_else(PoisonError::into_inner);
42            lifecycle
43                .iter()
44                .filter(|(_, active)| active.result.borrow().is_none())
45                .map(|(session_id, active)| {
46                    if active.kind != LifecycleKind::Cleanup {
47                        active.request_cancel();
48                    }
49                    let stage = active
50                        .active_stages
51                        .keys()
52                        .next_back()
53                        .map(|stage| stage.label())
54                        .unwrap_or_else(|| "container cleanup".to_owned());
55                    (
56                        session_id.clone(),
57                        active.kind,
58                        stage,
59                        active.cancelled.clone(),
60                        active.result.clone(),
61                    )
62                })
63                .collect::<Vec<_>>()
64        };
65        let cleanup_deadline = tokio::time::Instant::now() + Duration::from_secs(8);
66        for (session_id, kind, stage, cancelled, result) in &mut pending {
67            if *kind != LifecycleKind::Cleanup || result.borrow().is_some() {
68                continue;
69            }
70            tracing::info!(%session_id, %stage, "daemon shutdown is waiting for deferred cleanup");
71            self.set_lifecycle_notice(
72                session_id,
73                &format!("Daemon shutdown is waiting for {stage}"),
74            );
75            let finished = tokio::time::timeout_at(cleanup_deadline, async {
76                while result.borrow_and_update().is_none() {
77                    result.changed().await.with_context(|| {
78                        format!("cleanup owner stopped without a result for session {session_id}")
79                    })?;
80                }
81                Ok::<_, anyhow::Error>(())
82            })
83            .await;
84            match finished {
85                Ok(result) => result?,
86                Err(_) => {
87                    tracing::warn!(%session_id, %stage, "deferred cleanup exceeded the daemon shutdown drain deadline");
88                    cancelled.store(true, Ordering::Release);
89                }
90            }
91        }
92        let join_deadline = tokio::time::Instant::now() + Duration::from_secs(1);
93        for (session_id, _, stage, cancelled, mut result) in pending {
94            cancelled.store(true, Ordering::Release);
95            let joined = tokio::time::timeout_at(join_deadline, async {
96                while result.borrow_and_update().is_none() {
97                    result.changed().await.with_context(|| {
98                        format!("lifecycle owner stopped without a result for session {session_id}")
99                    })?;
100                }
101                Ok::<_, anyhow::Error>(())
102            })
103            .await;
104            if joined.is_err() {
105                bail!(
106                    "timed out cancelling lifecycle owner for session {session_id} while {stage}"
107                );
108            }
109            joined.expect("checked timeout")?;
110        }
111        Ok(())
112    }
113
114    /// Every lifecycle operation running now.
115    ///
116    /// The dashboard receives these through a watch channel built by its own
117    /// poller, which the phone server does not have; rather than plumb that
118    /// channel through the session-manager handle, the phone loop reads the
119    /// same state directly. The read is a mutex acquisition over a small map,
120    /// and it happens once per published snapshot, so it never blocks the
121    /// loop the way an await on the async snapshot path would.
122    pub fn active_lifecycles(&self) -> Vec<RuntimeLifecycleView> {
123        // Controller before lifecycle, the order worker polling takes. Both are
124        // plain mutex acquisitions over small maps, so a render loop calling
125        // this never awaits.
126        let controller = self
127            .controller
128            .lock()
129            .unwrap_or_else(PoisonError::into_inner);
130        self.active_lifecycles_with(&controller)
131    }
132
133    /// The same view for a caller that already holds the controller lock.
134    /// The lock is not reentrant, so taking it again here would deadlock the
135    /// daemon.
136    pub(super) fn active_lifecycles_with(
137        &self,
138        controller: &Controller,
139    ) -> Vec<RuntimeLifecycleView> {
140        self.lifecycle
141            .lock()
142            .unwrap_or_else(PoisonError::into_inner)
143            .iter()
144            .filter(|(_, active)| active.is_visible())
145            .map(|(session_id, active)| RuntimeLifecycleView {
146                operation_id: active.operation_id.clone(),
147                cancellable: active.is_cancellable()
148                    && lifecycle_cancellable(
149                        active.kind,
150                        durable_session_state(controller, session_id),
151                    ),
152                session_id: session_id.clone(),
153                kind: active.kind.into(),
154                started_at_epoch_seconds: active.started_at_epoch_seconds,
155                active_stages: active
156                    .active_stages
157                    .iter()
158                    .map(|(stage, (_, started_at))| (*stage, *started_at))
159                    .collect(),
160                resume_destination: active.resume_destination.clone(),
161                notice: active.notice.clone(),
162            })
163            .collect()
164    }
165
166    /// The lifecycle state of one in-memory record, or `None` when the daemon
167    /// holds no record for it. Reading one field costs one lock rather than a
168    /// clone of every record, which is what a poll wants.
169    pub fn session_state(&self, session_id: &str) -> Option<mj_core::state::SessionState> {
170        if self.close_is_requested(session_id) {
171            return Some(SessionState::Closing);
172        }
173        self.controller
174            .lock()
175            .unwrap_or_else(PoisonError::into_inner)
176            .state
177            .sessions
178            .get(session_id)
179            .map(|record| record.state)
180    }
181
182    /// One in-memory session record, or `None` when the daemon holds none.
183    pub fn session_record(&self, session_id: &str) -> Option<SessionRecord> {
184        self.controller
185            .lock()
186            .unwrap_or_else(PoisonError::into_inner)
187            .state
188            .sessions
189            .get(session_id)
190            .cloned()
191    }
192
193    pub async fn workspace_session_handle(
194        &self,
195        session_id: &str,
196    ) -> Result<crate::session_manager::ManagedSessionHandle> {
197        let record = self.session_record(session_id).context("unknown session")?;
198        ensure!(
199            record.target.is_some()
200                && record.state == SessionState::Running
201                && !self.close_is_requested(session_id),
202            "session must have a live running target for file injection"
203        );
204        self.session_manager.session(session_id.to_owned()).await
205    }
206
207    /// Checkpoint a session now and publish the result, the way the daemon's
208    /// own checkpoint action does.
209    ///
210    /// The API's bundle export needs a fresh archive for a running session. Only
211    /// that session's own lifecycle operation can conflict with its checkpoint,
212    /// so this refuses when the session itself is mid-operation and returns a
213    /// [`SessionLifecycleBusy`] the export path can fall back on. It must not
214    /// take the process-wide lifecycle guard: that rejected every export while
215    /// any unrelated session anywhere was mid-lifecycle (#1010).
216    pub async fn checkpoint_session_now(
217        &self,
218        session_id: &str,
219    ) -> Result<mj_core::state::CheckpointMetadata> {
220        if let Some(busy) = self.session_lifecycle_busy(session_id) {
221            return Err(anyhow::Error::new(busy));
222        }
223        let mut controller = blocking(Controller::load).await?;
224        let checkpoint = controller.checkpoint_session(session_id).await?;
225        refresh_runtime_controller(self).await;
226        Ok(checkpoint)
227    }
228
229    /// The lifecycle operation this specific session is running, if any, named
230    /// along with its age. A checkpoint conflicts only with its own session's
231    /// operations, never with another session's (#1010).
232    pub(super) fn session_lifecycle_busy(&self, session_id: &str) -> Option<SessionLifecycleBusy> {
233        let lifecycle = self
234            .lifecycle
235            .lock()
236            .unwrap_or_else(PoisonError::into_inner);
237        let active = lifecycle.get(session_id)?;
238        active
239            .result
240            .borrow()
241            .is_none()
242            .then(|| describe_lifecycle_busy(session_id, active))
243    }
244
245    /// Any session's running lifecycle operation, named the same way. The
246    /// config-rename guard reports this, so its refusal says which operation
247    /// stands in the way rather than only that one does (#1010).
248    pub(super) fn any_lifecycle_busy(&self) -> Option<SessionLifecycleBusy> {
249        self.lifecycle
250            .lock()
251            .unwrap_or_else(PoisonError::into_inner)
252            .iter()
253            .find(|(_, active)| active.result.borrow().is_none())
254            .map(|(session_id, active)| describe_lifecycle_busy(session_id, active))
255    }
256
257    /// In-memory records and ownership sampled with the same lock order as
258    /// completion. A web publish must not pair old records with a new absence
259    /// of ownership, even while its background database reload is in flight.
260    pub fn session_projection(
261        &self,
262    ) -> (BTreeMap<String, SessionRecord>, Vec<RuntimeLifecycleView>) {
263        let controller = self
264            .controller
265            .lock()
266            .unwrap_or_else(PoisonError::into_inner);
267        let operations = self.active_lifecycles_with(&controller);
268        let mut records = controller.state.sessions.clone();
269        for id in self
270            .close_requested
271            .lock()
272            .unwrap_or_else(PoisonError::into_inner)
273            .iter()
274        {
275            if let Some(record) = records.get_mut(id)
276                && record.state != SessionState::Stopped
277            {
278                record.state = SessionState::Closing;
279            }
280        }
281        (records, operations)
282    }
283
284    pub fn cancel_lifecycle_if_active(&self, session_id: &str) {
285        if let Some(active) = self
286            .lifecycle
287            .lock()
288            .unwrap_or_else(PoisonError::into_inner)
289            .get(session_id)
290        {
291            active.request_cancel();
292            self.publish_revision();
293        }
294    }
295
296    pub(super) fn set_lifecycle_resume_destination(
297        &self,
298        session_id: &str,
299        profile_id: String,
300        target_id: String,
301    ) {
302        if let Some(active) = self
303            .lifecycle
304            .lock()
305            .unwrap_or_else(PoisonError::into_inner)
306            .get_mut(session_id)
307        {
308            active.resume_destination = Some((profile_id, target_id));
309            self.publish_revision();
310        }
311    }
312
313    pub(super) fn change_lifecycle_stage(
314        &self,
315        session_id: &str,
316        stage: ProvisionStage,
317        active: bool,
318    ) {
319        let changed = {
320            let mut lifecycle = self
321                .lifecycle
322                .lock()
323                .unwrap_or_else(PoisonError::into_inner);
324            let Some(operation) = lifecycle.get_mut(session_id) else {
325                return;
326            };
327            if active {
328                let entry = operation
329                    .active_stages
330                    .entry(stage)
331                    .or_insert_with(|| (0, epoch_seconds()));
332                entry.0 += 1;
333                entry.0 == 1
334            } else {
335                let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
336                    return;
337                };
338                *count -= 1;
339                if *count == 0 {
340                    operation.active_stages.remove(&stage);
341                    true
342                } else {
343                    false
344                }
345            }
346        };
347        if changed {
348            self.publish_revision();
349        }
350    }
351
352    /// Record something the daemon did on its own, for every attached surface
353    /// to report once.
354    pub(super) fn push_notice(&self, session_id: &str, text: impl Into<String>) {
355        const RETAINED_NOTICES: usize = 32;
356
357        let notice = RuntimeNotice {
358            id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
359            session_id: session_id.to_owned(),
360            text: text.into(),
361        };
362        {
363            let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
364            notices.push_back(notice);
365            while notices.len() > RETAINED_NOTICES {
366                notices.pop_front();
367            }
368        }
369        self.publish_revision();
370    }
371
372    pub(super) fn set_lifecycle_notice(&self, session_id: &str, notice: &str) {
373        if let Some(active) = self
374            .lifecycle
375            .lock()
376            .unwrap_or_else(PoisonError::into_inner)
377            .get_mut(session_id)
378        {
379            if active.kind == LifecycleKind::Move && notice == "Preparing destination" {
380                active.move_source_closed = true;
381            }
382            active.notice = Some(notice.to_owned());
383            self.publish_revision();
384        }
385    }
386}