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, local_checkout_roots) = {
307            let owner = self.owner();
308            let controller = owner.controller();
309            let records = owner.projected_records();
310            let local_checkout_roots = records
311                .iter()
312                .filter_map(|(id, _)| match controller.state.checkout(id).ok()? {
313                    mj_core::state::Checkout::ManagedWorktree { worktree, .. }
314                        if worktree.target == mj_core::state::ManagedWorktreeTarget::Local =>
315                    {
316                        Some(worktree.worktree_root.clone())
317                    }
318                    _ => None,
319                })
320                .collect::<Vec<_>>();
321            (
322                records,
323                controller.state.subagents.clone(),
324                owner
325                    .committed()
326                    .map(|committed| committed.moves.clone())
327                    .unwrap_or_default(),
328                controller.config.clone(),
329                local_checkout_roots,
330            )
331        };
332        let mut candidates = mj_client::daemon::ResumeCandidates {
333            local_checkout_roots,
334            ..mj_client::daemon::ResumeCandidates::default()
335        };
336        for (id, record) in &records {
337            if let Some(native_session_id) = &record.native_session_id {
338                candidates
339                    .adopted_native_sessions
340                    .push((record.harness_kind, native_session_id.clone()));
341            }
342            if record.state.is_active()
343                || subagents.contains_key(id)
344                || mj_core::native_agent::is_view_id(id)
345            {
346                continue;
347            }
348            if let Some(operation) = moves.get(id) {
349                candidates.moves.push(operation.clone());
350            }
351            candidates
352                .candidates
353                .push(mj_client::daemon::ResumeCandidate::of(record, &config));
354        }
355        candidates
356    }
357
358    /// The session `mj go` opens in `workspace_id`: `last_session_id` while it
359    /// is still eligible, otherwise the most recently updated eligible one.
360    pub(super) fn go_startup_session(
361        &self,
362        workspace_id: &str,
363        last_session_id: Option<&str>,
364    ) -> Option<SessionRecord> {
365        let (records, subagents) = {
366            let owner = self.owner();
367            (
368                owner.projected_records(),
369                owner.controller().state.subagents.clone(),
370            )
371        };
372        let eligible = |session: &&SessionRecord| {
373            session.workspace_id == workspace_id
374                && !session.archived
375                && !subagents.contains_key(&session.id)
376                && !mj_core::native_agent::is_view_id(&session.id)
377                && session.state != SessionState::DestroyedWithDataLoss
378        };
379        last_session_id
380            .and_then(|id| records.get(id))
381            .filter(eligible)
382            .or_else(|| {
383                records
384                    .values()
385                    .filter(eligible)
386                    .max_by_key(|session| &session.updated_at)
387            })
388            .cloned()
389    }
390
391    pub(crate) fn worker_controller_projection(&self) -> Controller {
392        self.owner().pollable_worker_inputs().controller()
393    }
394
395    pub(crate) fn controller_projection(&self) -> Controller {
396        let owner = self.owner();
397        let mut state = owner.controller().state.clone();
398        state.sessions = owner.projected_records();
399        Controller {
400            config: owner.controller().config.clone(),
401            state,
402        }
403    }
404
405    pub(crate) fn active_controller_projection(&self) -> Controller {
406        let owner = self.owner();
407        Controller {
408            config: owner.controller().config.clone(),
409            state: mj_core::state::State {
410                sessions: owner
411                    .indexes
412                    .active
413                    .keys()
414                    .filter_map(|id| {
415                        owner
416                            .controller()
417                            .state
418                            .sessions
419                            .get(id)
420                            .map(|record| (id.clone(), record.clone()))
421                    })
422                    .collect(),
423                ..Default::default()
424            },
425        }
426    }
427
428    pub(super) fn set_lifecycle_resume_destination(
429        &self,
430        session_id: &str,
431        profile_id: String,
432        target_id: String,
433    ) {
434        if let Some(active) = self.owner().lifecycle.get_mut(session_id) {
435            active.resume_destination = Some((profile_id, target_id));
436            self.publish_revision();
437        }
438    }
439
440    pub(super) fn change_lifecycle_stage(
441        &self,
442        session_id: &str,
443        operation_id: &str,
444        stage: ProvisionStage,
445        active: bool,
446    ) {
447        let changed = {
448            let mut lifecycle_owner = self.owner();
449            let lifecycle = &mut lifecycle_owner.lifecycle;
450            let Some(operation) = lifecycle.get_mut(session_id) else {
451                return;
452            };
453            if operation.operation_id != operation_id || !operation.is_running() {
454                return;
455            }
456            if active {
457                let started_at = stage
458                    .started_at_epoch_seconds()
459                    .unwrap_or_else(epoch_seconds);
460                let entry = operation
461                    .active_stages
462                    .entry(stage)
463                    .or_insert_with(|| (0, started_at));
464                entry.0 += 1;
465                entry.0 == 1
466            } else {
467                let Some((count, _)) = operation.active_stages.get_mut(&stage) else {
468                    return;
469                };
470                *count -= 1;
471                if *count == 0 {
472                    operation.active_stages.remove(&stage);
473                    true
474                } else {
475                    false
476                }
477            }
478        };
479        if changed {
480            self.publish_revision();
481        }
482    }
483
484    /// Record something the daemon did on its own, for every attached surface
485    /// to report once.
486    pub(crate) fn push_notice(&self, session_id: &str, text: impl Into<String>) {
487        const RETAINED_NOTICES: usize = 32;
488
489        let notice = RuntimeNotice {
490            id: self.next_notice_id.fetch_add(1, Ordering::AcqRel),
491            session_id: session_id.to_owned(),
492            text: text.into(),
493        };
494        {
495            let mut notices = self.notices.lock().unwrap_or_else(PoisonError::into_inner);
496            notices.push_back(notice);
497            while notices.len() > RETAINED_NOTICES {
498                notices.pop_front();
499            }
500        }
501        self.publish_revision();
502    }
503
504    /// Connect the quota poller's wake-up channel. The poller is the daemon's
505    /// only prober; a surface that wants a fresh reading asks the daemon, and
506    /// this is how the ask reaches the poller.
507    pub(crate) fn attach_quota_refresh(&self, refresh: tokio::sync::mpsc::Sender<()>) {
508        self.quota
509            .lock()
510            .unwrap_or_else(PoisonError::into_inner)
511            .refresh = Some(refresh);
512    }
513
514    /// Ask the quota poller to probe now. Requests that arrive while one is
515    /// already waiting are one request.
516    pub(crate) fn request_quota_refresh(&self) -> Result<()> {
517        let quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
518        let refresh = quota
519            .refresh
520            .as_ref()
521            .context("the daemon's quota service is not running")?;
522        // A full channel means a refresh is already queued.
523        let _ = refresh.try_send(());
524        Ok(())
525    }
526
527    /// Replace what the daemon says about quota, and wake every attached
528    /// surface when it changed.
529    pub(crate) fn publish_quotas(&self, snapshot: mj_client::quota::QuotaSnapshot) {
530        {
531            let mut quota = self.quota.lock().unwrap_or_else(PoisonError::into_inner);
532            if quota.snapshot == snapshot {
533                return;
534            }
535            quota.snapshot = snapshot;
536        }
537        self.publish_revision();
538    }
539
540    pub(super) fn reserve_move_destination(&self, session_id: &str, operation_id: &str) {
541        let mut owner = self.owner();
542        if let Some(active) = owner.lifecycle.get_mut(session_id)
543            && active.operation_id == operation_id
544            && active.kind == LifecycleKind::Move
545        {
546            active.phase = match active.phase {
547                LifecyclePhase::Executing => LifecyclePhase::MovingDestination,
548                LifecyclePhase::Cancelling => LifecyclePhase::CancellingMoveDestination,
549                _ => return,
550            };
551            self.publish_revision();
552        }
553    }
554
555    pub(super) fn set_lifecycle_notice(&self, session_id: &str, operation_id: &str, notice: &str) {
556        if let Some(active) = self.owner().lifecycle.get_mut(session_id)
557            && active.operation_id == operation_id
558            && active.is_running()
559        {
560            active.notice = Some(notice.to_owned());
561            self.publish_revision();
562        }
563    }
564
565    fn cancel_operation(&self, session_id: &str, operation_id: &str) {
566        if let Some(active) = self.owner().lifecycle.get_mut(session_id)
567            && active.operation_id == operation_id
568            && active.request_cancel()
569        {
570            self.publish_revision();
571        }
572    }
573}