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        let mut controller_owner = self.owner();
6        let durable = durable_session_state(controller_owner.controller(), session_id);
7        let active = controller_owner
8            .lifecycle
9            .get_mut(session_id)
10            .with_context(|| {
11                format!("no lifecycle operation is running for session {session_id}")
12            })?;
13        ensure!(
14            lifecycle_cancellable(active.kind, durable),
15            "stop of {session_id} has passed its verified checkpoint and is removing the target; \
16             it cannot be cancelled"
17        );
18        ensure!(
19            active.request_cancel(),
20            "lifecycle operation is no longer cancellable"
21        );
22        drop(controller_owner);
23        self.publish_revision();
24        Ok(())
25    }
26
27    /// Let storage cleanup drain briefly, then cancel and join every lifecycle
28    /// owner before the daemon closes its session manager and database writer.
29    /// One shared deadline bounds all cleanup tasks rather than granting eight
30    /// seconds to each session serially.
31    pub(super) async fn cancel_and_wait_lifecycles(&self) -> Result<()> {
32        let mut pending = {
33            let mut lifecycle_owner = self.owner();
34            let lifecycle = &mut lifecycle_owner.lifecycle;
35            lifecycle
36                .iter_mut()
37                .filter(|(_, active)| active.is_running())
38                .map(|(session_id, active)| {
39                    if active.kind != LifecycleKind::Cleanup {
40                        active.request_cancel();
41                    }
42                    let stage = active
43                        .active_stages
44                        .keys()
45                        .next_back()
46                        .map(|stage| stage.label())
47                        .unwrap_or_else(|| "container cleanup".to_owned());
48                    (
49                        session_id.clone(),
50                        active.operation_id.clone(),
51                        active.kind,
52                        stage,
53                        active.result.clone(),
54                    )
55                })
56                .collect::<Vec<_>>()
57        };
58        let cleanup_deadline = tokio::time::Instant::now() + Duration::from_secs(8);
59        for (session_id, operation_id, kind, stage, result) in &mut pending {
60            if *kind != LifecycleKind::Cleanup || result.borrow().is_some() {
61                continue;
62            }
63            tracing::info!(%session_id, %stage, "daemon shutdown is waiting for deferred cleanup");
64            self.set_lifecycle_notice(
65                session_id,
66                operation_id,
67                &format!("Daemon shutdown is waiting for {stage}"),
68            );
69            let finished = tokio::time::timeout_at(cleanup_deadline, async {
70                while result.borrow_and_update().is_none() {
71                    result.changed().await.with_context(|| {
72                        format!("cleanup owner stopped without a result for session {session_id}")
73                    })?;
74                }
75                Ok::<_, anyhow::Error>(())
76            })
77            .await;
78            match finished {
79                Ok(result) => result?,
80                Err(_) => {
81                    tracing::warn!(%session_id, %stage, "deferred cleanup exceeded the daemon shutdown drain deadline");
82                    self.cancel_operation(session_id, operation_id);
83                }
84            }
85        }
86        let join_started = tokio::time::Instant::now();
87        let join_deadline = join_started + Duration::from_secs(1);
88        // Launch cancellation is followed by independent, bounded teardown.
89        // Join that owner before closing the store it must settle.
90        let startup_cleanup_deadline = join_started
91            + crate::controller::FAILED_STARTUP_CLEANUP_TIMEOUT
92            + Duration::from_secs(1);
93        for (session_id, operation_id, kind, stage, mut result) in pending {
94            let join_deadline = if matches!(
95                kind,
96                LifecycleKind::Create
97                    | LifecycleKind::Resume
98                    | LifecycleKind::Unpark
99                    | LifecycleKind::StartupCleanup
100            ) {
101                startup_cleanup_deadline
102            } else {
103                join_deadline
104            };
105            self.cancel_operation(&session_id, &operation_id);
106            let joined = tokio::time::timeout_at(join_deadline, async {
107                while result.borrow_and_update().is_none() {
108                    result.changed().await.with_context(|| {
109                        format!("lifecycle owner stopped without a result for session {session_id}")
110                    })?;
111                }
112                Ok::<_, anyhow::Error>(())
113            })
114            .await;
115            if joined.is_err() {
116                bail!(
117                    "timed out cancelling lifecycle owner for session {session_id} while {stage}"
118                );
119            }
120            joined.expect("checked timeout")?;
121        }
122        Ok(())
123    }
124
125    /// Every lifecycle operation running now.
126    ///
127    /// The dashboard receives these through a watch channel built by its own
128    /// poller, which the phone server does not have; rather than plumb that
129    /// channel through the session-manager handle, the phone loop reads the
130    /// same state directly. The read is a mutex acquisition over a small map,
131    /// and it happens once per published snapshot, so it never blocks the
132    /// loop the way an await on the async snapshot path would.
133    pub fn active_lifecycles(&self) -> Vec<RuntimeLifecycleView> {
134        let controller_owner = self.owner();
135        Self::active_lifecycles_with(&controller_owner)
136    }
137
138    /// Project records and ownership from the same transition boundary.
139    pub(super) fn active_lifecycles_with(owner: &RuntimeStateOwner) -> Vec<RuntimeLifecycleView> {
140        let controller = owner.controller();
141        owner
142            .lifecycle
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        let owner = self.owner();
171        if owner.close_requested.contains(session_id) {
172            return Some(SessionState::Closing);
173        }
174        owner
175            .controller()
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.owner()
185            .controller()
186            .state
187            .sessions
188            .get(session_id)
189            .cloned()
190    }
191
192    pub async fn workspace_session_handle(
193        &self,
194        session_id: &str,
195    ) -> Result<crate::session_manager::ManagedSessionHandle> {
196        let record = self.session_record(session_id).context("unknown session")?;
197        ensure!(
198            record.target.is_some()
199                && record.state == SessionState::Running
200                && !self.close_is_requested(session_id),
201            "session must have a live running target for file injection"
202        );
203        self.session_manager.session(session_id.to_owned()).await
204    }
205
206    /// Checkpoint a session now and publish the result, the way the daemon's
207    /// own checkpoint action does.
208    ///
209    /// The API's bundle export needs a fresh archive for a running session. Only
210    /// that session's own lifecycle operation can conflict with its checkpoint,
211    /// so this refuses when the session itself is mid-operation and returns a
212    /// [`SessionLifecycleBusy`] the export path can fall back on. It must not
213    /// take the process-wide lifecycle guard: that rejected every export while
214    /// any unrelated session anywhere was mid-lifecycle (#1010).
215    pub async fn checkpoint_session_now(
216        &self,
217        session_id: &str,
218    ) -> Result<mj_core::state::CheckpointMetadata> {
219        let _upgrade_work = crate::upgrade::activity("requested checkpoint")?;
220        if let Some(busy) = self.session_lifecycle_busy(session_id) {
221            return Err(anyhow::Error::new(busy));
222        }
223        let session_id = session_id.to_owned();
224        // Checkpoint capture and cleanup run synchronous filesystem, database
225        // and SSH operations between relay awaits. Keep the whole future,
226        // including its destructors, off the daemon's event loop.
227        let checkpoint = blocking(move || {
228            let mut controller = Controller::load()?;
229            mj_core::runtime::block_on(controller.checkpoint_session(&session_id))?
230        })
231        .await?;
232        refresh_runtime_controller(self).await;
233        Ok(checkpoint)
234    }
235
236    /// The lifecycle operation this specific session is running, if any, named
237    /// along with its age. A checkpoint conflicts only with its own session's
238    /// operations, never with another session's (#1010).
239    pub(super) fn session_lifecycle_busy(&self, session_id: &str) -> Option<SessionLifecycleBusy> {
240        let lifecycle_owner = self.owner();
241        let lifecycle = &lifecycle_owner.lifecycle;
242        let active = lifecycle.get(session_id)?;
243        active
244            .is_running()
245            .then(|| describe_lifecycle_busy(session_id, active))
246    }
247
248    /// Any session's running lifecycle operation, named the same way. The
249    /// config-rename guard reports this, so its refusal says which operation
250    /// stands in the way rather than only that one does (#1010).
251    pub(super) fn any_lifecycle_busy(&self) -> Option<SessionLifecycleBusy> {
252        self.owner()
253            .lifecycle
254            .iter()
255            .find(|(_, active)| active.is_running())
256            .map(|(session_id, active)| describe_lifecycle_busy(session_id, active))
257    }
258
259    /// In-memory records and ownership sampled with the same lock order as
260    /// completion. A web publish must not pair old records with a new absence
261    /// of ownership, even while its background database reload is in flight.
262    pub fn session_projection(
263        &self,
264    ) -> (
265        mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
266        Vec<RuntimeLifecycleView>,
267    ) {
268        let controller_owner = self.owner();
269        let operations = Self::active_lifecycles_with(&controller_owner);
270        (controller_owner.projected_records(), operations)
271    }
272
273    /// What the resume dialog lists: every inactive session that is not a
274    /// sub-agent, with the import de-duplication data from every record.
275    /// Only shared map handles cross the owner lock; the copies are made after.
276    pub(super) fn resume_candidates(&self) -> mj_client::daemon::ResumeCandidates {
277        let (records, subagents, moves, config) = {
278            let owner = self.owner();
279            (
280                owner.projected_records(),
281                owner.controller().state.subagents.clone(),
282                owner
283                    .committed()
284                    .map(|committed| committed.moves.clone())
285                    .unwrap_or_default(),
286                owner.controller().config.clone(),
287            )
288        };
289        let mut candidates = mj_client::daemon::ResumeCandidates::default();
290        for (id, record) in &records {
291            if let Some(native_session_id) = &record.native_session_id {
292                candidates
293                    .adopted_native_sessions
294                    .push((record.harness_kind, native_session_id.clone()));
295            }
296            if let Some(checkout) = &record.managed_worktree
297                && checkout.target == mj_core::state::ManagedWorktreeTarget::Local
298            {
299                candidates
300                    .local_checkout_roots
301                    .push(checkout.worktree_root.clone());
302            }
303            if record.state.is_active()
304                || subagents.contains_key(id)
305                || mj_core::native_agent::is_view_id(id)
306            {
307                continue;
308            }
309            if let Some(operation) = moves.get(id) {
310                candidates.moves.push(operation.clone());
311            }
312            candidates
313                .candidates
314                .push(mj_client::daemon::ResumeCandidate::of(record, &config));
315        }
316        candidates
317    }
318
319    /// The session `mj go` opens in `workspace_id`: `last_session_id` while it
320    /// is still eligible, otherwise the most recently updated eligible one.
321    pub(super) fn go_startup_session(
322        &self,
323        workspace_id: &str,
324        last_session_id: Option<&str>,
325    ) -> Option<SessionRecord> {
326        let (records, subagents) = {
327            let owner = self.owner();
328            (
329                owner.projected_records(),
330                owner.controller().state.subagents.clone(),
331            )
332        };
333        let eligible = |session: &&SessionRecord| {
334            session.workspace_id == workspace_id
335                && !session.archived
336                && !subagents.contains_key(&session.id)
337                && !mj_core::native_agent::is_view_id(&session.id)
338                && session.state != SessionState::DestroyedWithDataLoss
339        };
340        last_session_id
341            .and_then(|id| records.get(id))
342            .filter(eligible)
343            .or_else(|| {
344                records
345                    .values()
346                    .filter(eligible)
347                    .max_by_key(|session| &session.updated_at)
348            })
349            .cloned()
350    }
351
352    pub(crate) fn worker_controller_projection(&self) -> Controller {
353        self.owner().pollable_worker_inputs().controller()
354    }
355
356    pub(crate) fn controller_projection(&self) -> Controller {
357        let owner = self.owner();
358        let mut state = owner.controller().state.clone();
359        state.sessions = owner.projected_records();
360        Controller {
361            config: owner.controller().config.clone(),
362            state,
363        }
364    }
365
366    pub(crate) fn active_controller_projection(&self) -> Controller {
367        let owner = self.owner();
368        Controller {
369            config: owner.controller().config.clone(),
370            state: mj_core::state::State {
371                sessions: owner
372                    .indexes
373                    .active
374                    .keys()
375                    .filter_map(|id| {
376                        owner
377                            .controller()
378                            .state
379                            .sessions
380                            .get(id)
381                            .map(|record| (id.clone(), record.clone()))
382                    })
383                    .collect(),
384                ..Default::default()
385            },
386        }
387    }
388
389    pub fn cancel_lifecycle_if_active(&self, session_id: &str) {
390        if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
391            active.request_cancel();
392            self.publish_revision();
393        }
394    }
395
396    pub(super) fn set_lifecycle_resume_destination(
397        &self,
398        session_id: &str,
399        profile_id: String,
400        target_id: String,
401    ) {
402        if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
403            active.resume_destination = Some((profile_id, target_id));
404            self.publish_revision();
405        }
406    }
407
408    pub(super) fn change_lifecycle_stage(
409        &self,
410        session_id: &str,
411        operation_id: &str,
412        stage: ProvisionStage,
413        active: bool,
414    ) {
415        let changed = {
416            let mut lifecycle_owner = self.owner();
417            let lifecycle = &mut lifecycle_owner.lifecycle;
418            let Some(operation) = lifecycle.get_mut(session_id) else {
419                return;
420            };
421            if operation.operation_id != operation_id || !operation.is_running() {
422                return;
423            }
424            if active {
425                let entry = operation
426                    .active_stages
427                    .entry(stage)
428                    .or_insert_with(|| (0, epoch_seconds()));
429                entry.0 += 1;
430                entry.0 == 1
431            } else {
432                let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
433                    return;
434                };
435                *count -= 1;
436                if *count == 0 {
437                    operation.active_stages.remove(&stage);
438                    true
439                } else {
440                    false
441                }
442            }
443        };
444        if changed {
445            self.publish_revision();
446        }
447    }
448
449    /// Record something the daemon did on its own, for every attached surface
450    /// to report once.
451    pub(crate) fn push_notice(&self, session_id: &str, text: impl Into<String>) {
452        const RETAINED_NOTICES: usize = 32;
453
454        let notice = RuntimeNotice {
455            id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
456            session_id: session_id.to_owned(),
457            text: text.into(),
458        };
459        {
460            let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
461            notices.push_back(notice);
462            while notices.len() > RETAINED_NOTICES {
463                notices.pop_front();
464            }
465        }
466        self.publish_revision();
467    }
468
469    /// Connect the quota poller's wake-up channel. The poller is the daemon's
470    /// only prober; a surface that wants a fresh reading asks the daemon, and
471    /// this is how the ask reaches the poller.
472    pub(crate) fn attach_quota_refresh(&self, refresh: tokio::sync::mpsc::Sender<()>) {
473        self.quota
474            .lock()
475            .unwrap_or_else(PoisonError::into_inner)
476            .refresh = Some(refresh);
477    }
478
479    /// Ask the quota poller to probe now. Requests that arrive while one is
480    /// already waiting are one request.
481    pub(crate) fn request_quota_refresh(&self) -> Result<()> {
482        let quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
483        let refresh = quota
484            .refresh
485            .as_ref()
486            .context("the daemon's quota service is not running")?;
487        // A full channel means a refresh is already queued.
488        let _ = refresh.try_send(());
489        Ok(())
490    }
491
492    /// Replace what the daemon says about quota, and wake every attached
493    /// surface when it changed.
494    pub(crate) fn publish_quotas(&self, snapshot: mj_client::quota::QuotaSnapshot) {
495        {
496            let mut quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
497            if quota.snapshot == snapshot {
498                return;
499            }
500            quota.snapshot = snapshot;
501        }
502        self.publish_revision();
503    }
504
505    pub(super) fn reserve_move_destination(&self, session_id: &str, operation_id: &str) {
506        let mut owner = self.owner();
507        if let Some(active) = owner.lifecycle.get_mut(session_id)
508            && active.operation_id == operation_id
509            && active.kind == LifecycleKind::Move
510        {
511            active.phase = match active.phase {
512                LifecyclePhase::Executing => LifecyclePhase::MovingDestination,
513                LifecyclePhase::Cancelling => LifecyclePhase::CancellingMoveDestination,
514                _ => return,
515            };
516            self.publish_revision();
517        }
518    }
519
520    pub(super) fn set_lifecycle_notice(&self, session_id: &str, operation_id: &str, notice: &str) {
521        if let Some(active) = self.owner().lifecycle.get_mut(session_id)
522            && active.operation_id == operation_id
523            && active.is_running()
524        {
525            active.notice = Some(notice.to_owned());
526            self.publish_revision();
527        }
528    }
529
530    fn cancel_operation(&self, session_id: &str, operation_id: &str) {
531        if let Some(active) = self.owner().lifecycle.get_mut(session_id)
532            && active.operation_id == operation_id
533            && active.request_cancel()
534        {
535            self.publish_revision();
536        }
537    }
538}