Skip to main content

mj_controller/daemon/
views.rs

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