Skip to main content

mj_controller/daemon/
snapshot.rs

1use super::*;
2
3impl RuntimeState {
4    pub async fn reload_controller(&self) -> Result<()> {
5        // Serialize installs so an earlier phone publication cannot overwrite
6        // a later completed lifecycle with the controller snapshot it loaded.
7        let _mutation = self.config_mutation.lock().await;
8        let controller_loader = self.controller_loader;
9        let controller = tokio::task::spawn_blocking(controller_loader)
10            .await
11            .context("daemon controller reload task panicked")??;
12        let session_count = controller.state.sessions.len();
13        *self
14            .controller
15            .lock()
16            .unwrap_or_else(PoisonError::into_inner) = controller;
17        let revision = self.publish_revision();
18        tracing::debug!(revision, session_count, "daemon controller state reloaded");
19        Ok(())
20    }
21
22    /// Definitive missing-target evidence belongs to the daemon, including
23    /// when no terminal is attached. Generic connection failures stay transient.
24    pub(super) fn missing_target_record(
25        &self,
26        session_id: &str,
27        view: &ManagedSessionView,
28    ) -> Option<(String, String)> {
29        let Some(ViewError::TargetMissing(detail)) = &view.error else {
30            return None;
31        };
32        if view.connected {
33            return None;
34        }
35        if self
36            .lifecycle
37            .lock()
38            .unwrap_or_else(PoisonError::into_inner)
39            .get(session_id)
40            .is_some_and(|active| active.result.borrow().is_none())
41        {
42            return None;
43        }
44        let controller = self
45            .controller
46            .lock()
47            .unwrap_or_else(PoisonError::into_inner);
48        let session = controller.state.sessions.get(session_id)?;
49        if !matches!(
50            session.state,
51            SessionState::Running | SessionState::Disconnected
52        ) {
53            return None;
54        }
55        Some((detail.clone(), session.updated_at.clone()))
56    }
57
58    pub(super) async fn persist_missing_target(
59        self: &Arc<Self>,
60        session_id: &str,
61        detail: String,
62        observed_updated_at: String,
63    ) -> Result<()> {
64        let changed = blocking({
65            let session_id = session_id.to_owned();
66            let detail = detail.clone();
67            move || {
68                crate::database::mark_session_target_missing_if_current(
69                    &session_id,
70                    &detail,
71                    &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
72                    &observed_updated_at,
73                )
74            }
75        })
76        .await?;
77        if changed.is_some() {
78            self.reload_controller().await?;
79            // A lost session has no checkpoint and no target, so its record
80            // offers nothing but Destroy. Report it once and discard it
81            // instead of leaving a tombstone behind.
82            if changed == Some(SessionState::Lost) {
83                self.discard_lost_session(session_id.to_owned()).await;
84            } else {
85                self.push_notice(session_id, detail);
86            }
87        }
88        Ok(())
89    }
90
91    /// One pass of the harness readiness wait: which live sessions have run
92    /// out of time to become usable.
93    ///
94    /// Locks are taken one at a time and nothing here touches the store, so
95    /// this belongs on the daemon's background tick. The write it leads to
96    /// does not; see [`Self::fail_unready_session`].
97    pub(super) fn sessions_without_a_usable_harness(&self) -> Vec<UnreadySession> {
98        let busy = {
99            let lifecycle = self
100                .lifecycle
101                .lock()
102                .unwrap_or_else(PoisonError::into_inner);
103            lifecycle
104                .iter()
105                .filter(|(_, active)| active.result.borrow().is_none())
106                .map(|(session_id, _)| session_id.clone())
107                .collect::<std::collections::BTreeSet<_>>()
108        };
109        let live = {
110            let controller = self
111                .controller
112                .lock()
113                .unwrap_or_else(PoisonError::into_inner);
114            controller
115                .state
116                .sessions
117                .values()
118                .filter(|record| record.state == SessionState::Running)
119                .filter(|record| !busy.contains(&record.id))
120                .map(|record| (record.id.clone(), record.updated_at.clone()))
121                .collect::<Vec<_>>()
122        };
123        let now = chrono::Utc::now();
124        let observations = {
125            let sessions = self.sessions.lock().unwrap_or_else(PoisonError::into_inner);
126            live.into_iter()
127                .filter(|(session_id, _)| !self.close_is_requested(session_id))
128                .map(|(session_id, updated_at)| {
129                    let view = sessions.get(&session_id);
130                    ReadinessObservation {
131                        harness_ready: view.is_some_and(|view| {
132                            view.operational.as_ref().is_some_and(
133                                mj_core::relay::RelayOperationalState::native_session_is_ready,
134                            )
135                        }),
136                        record_age: record_age(&updated_at, now),
137                        updated_at,
138                        detail: view
139                            .and_then(|view| view.error.as_ref())
140                            .map(|error| error.detail().to_owned()),
141                        session_id,
142                    }
143                })
144                .collect::<Vec<_>>()
145        };
146        self.harness_readiness
147            .lock()
148            .unwrap_or_else(PoisonError::into_inner)
149            .observe(std::time::Instant::now(), observations)
150    }
151
152    /// Fail one session whose harness never became usable, so it stops looking
153    /// like a session a driver can wait for.
154    pub(super) async fn fail_unready_session(&self, unready: UnreadySession) {
155        let cause = unready.cause();
156        let session_id = unready.session_id.clone();
157        tracing::warn!(%session_id, waited_seconds = unready.waited.as_secs(), "{cause}");
158        let applied = blocking({
159            let session_id = session_id.clone();
160            let cause = cause.clone();
161            move || {
162                let mut controller = Controller::load()?;
163                controller.fail_unready_session(&session_id, &cause, &unready.observed_updated_at)
164            }
165        })
166        .await;
167        match applied {
168            Ok(true) => {
169                if let Err(error) = self.reload_controller().await {
170                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after failing an unready session");
171                }
172                self.push_notice(&session_id, cause);
173                self.publish_revision();
174            }
175            Ok(false) => {}
176            Err(error) => tracing::warn!(
177                %session_id,
178                error = format!("{error:#}"),
179                "could not record that a session's harness never became usable"
180            ),
181        }
182    }
183
184    pub(super) async fn publish_session(
185        &self,
186        session_id: String,
187        view: ManagedSessionView,
188    ) -> Result<()> {
189        let connected = view.connected;
190        let has_snapshot = view.snapshot.is_some();
191        tracing::debug!(
192            %session_id,
193            connected,
194            has_snapshot,
195            "daemon received a session view"
196        );
197        if let Some(snapshot) = view.snapshot.as_ref() {
198            let controller = self
199                .controller
200                .lock()
201                .unwrap_or_else(PoisonError::into_inner);
202            if let Some(session) = controller.state.sessions.get(&session_id).cloned() {
203                // An upgrade only ever runs on a quiet session, and a session
204                // in a turn publishes a view every 150 ms. Skipping those here
205                // keeps the config clone off the streaming path; the
206                // coordinator still decides, from `quiet`, whether to act.
207                let quiet =
208                    view.connected && snapshot.operational.safe_to_replace(session.harness_kind);
209                if quiet {
210                    self.worker_upgrade_observer
211                        .observe(WorkerUpgradeObservation {
212                            session: session.clone(),
213                            config: controller.config.clone(),
214                            worker_build: snapshot.worker_build.clone(),
215                            quiet,
216                        });
217                }
218                self.recovery_observer.observe(RecoveryObservation {
219                    checkpoint_safe: snapshot
220                        .operational
221                        .safe_for_checkpoint(session.harness_kind),
222                    session,
223                    config: controller.config.clone(),
224                    latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
225                    execution: snapshot.materialized.execution,
226                });
227            }
228        }
229        self.sessions
230            .lock()
231            .unwrap_or_else(PoisonError::into_inner)
232            .insert(
233                session_id.clone(),
234                RuntimeSessionView::from_managed(session_id, view),
235            );
236        reach_test_hook("relay_projection_before_revision_publication").await?;
237        self.publish_revision();
238        Ok(())
239    }
240
241    pub(super) async fn runtime_snapshot(
242        &self,
243        workspace_id: &str,
244        after_revision: u64,
245        all_workspaces: bool,
246    ) -> Result<RuntimeSnapshot> {
247        let mut revisions = self.revisions.subscribe();
248        if *revisions.borrow_and_update() <= after_revision {
249            let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
250        }
251        let revision = self.revisions.current();
252        let moves = blocking(crate::database::load_move_operations).await?;
253        let workspace_names = blocking(crate::database::list_workspaces)
254            .await?
255            .into_iter()
256            .map(|workspace| (workspace.id, workspace.name))
257            .collect();
258        let session_ids = if all_workspaces {
259            self.controller
260                .lock()
261                .unwrap_or_else(PoisonError::into_inner)
262                .state
263                .sessions
264                .keys()
265                .cloned()
266                .collect::<BTreeSet<_>>()
267        } else {
268            let workspace_id = workspace_id.to_owned();
269            blocking(move || crate::database::session_ids_for_workspace(&workspace_id))
270                .await?
271                .into_iter()
272                .collect()
273        };
274        let native_owners = session_ids.clone();
275        let native_agents = blocking(move || {
276            let mut agents = Vec::new();
277            for owner in native_owners {
278                agents.extend(crate::database::load_native_agent_summaries(&owner)?);
279            }
280            Ok(agents)
281        })
282        .await?;
283        let sessions = self
284            .sessions
285            .lock()
286            .unwrap_or_else(PoisonError::into_inner)
287            .iter()
288            .filter(|(session_id, _)| session_ids.contains(*session_id))
289            .map(|(_, view)| view.clone())
290            .collect();
291        // Match the controller -> lifecycle lock order used by worker polling.
292        // Completion reloads records before publishing its result, so holding
293        // this guard prevents an absent operation paired with older records.
294        let controller = self
295            .controller
296            .lock()
297            .unwrap_or_else(PoisonError::into_inner);
298        let lifecycles = self
299            .lifecycle
300            .lock()
301            .unwrap_or_else(PoisonError::into_inner)
302            .iter()
303            .filter(|(session_id, active)| {
304                (all_workspaces
305                    || session_ids.contains(*session_id)
306                    || active.resume_workspace_id.as_deref() == Some(workspace_id))
307                    && active.is_visible()
308            })
309            .map(|(session_id, active)| RuntimeLifecycleView {
310                operation_id: active.operation_id.clone(),
311                cancellable: active.is_cancellable()
312                    && lifecycle_cancellable(
313                        active.kind,
314                        durable_session_state(&controller, session_id),
315                    ),
316                session_id: session_id.clone(),
317                kind: active.kind.into(),
318                started_at_epoch_seconds: active.started_at_epoch_seconds,
319                active_stages: active
320                    .active_stages
321                    .iter()
322                    .map(|(stage, (_, started_at))| (*stage, *started_at))
323                    .collect(),
324                resume_destination: active.resume_destination.clone(),
325                notice: active.notice.clone(),
326            })
327            .collect();
328        let reviews = self
329            .review_host
330            .views()
331            .into_iter()
332            .filter(|review| session_ids.contains(&review.session_id))
333            .collect();
334        let notices = self
335            .notices
336            .lock()
337            .unwrap_or_else(PoisonError::into_inner)
338            .iter()
339            .filter(|notice| notice_reaches_workspace(notice, &session_ids))
340            .cloned()
341            .collect();
342        let records = runtime_records_for_workspace(&controller, &session_ids);
343        let subagents = runtime_subagents_for_workspace(&controller, &records);
344        Ok(RuntimeSnapshot {
345            native_agents,
346            workspace_names,
347            moves: moves
348                .into_iter()
349                .filter(|operation| session_ids.contains(&operation.selection.session_id))
350                .collect(),
351            revision,
352            config: controller.config.clone(),
353            records,
354            sessions,
355            lifecycles,
356            reviews,
357            notices,
358            subagents,
359        })
360    }
361}
362
363/// Whether a notice belongs in one workspace's snapshot.
364///
365/// A notice usually belongs to a session, and only the workspace holding that
366/// session shows it. An empty session id marks a notice the daemon owns
367/// instead, such as a background container image download; there is no
368/// workspace to attach it to, and every workspace wants to see it.
369pub(super) fn notice_reaches_workspace(
370    notice: &RuntimeNotice,
371    session_ids: &BTreeSet<String>,
372) -> bool {
373    notice.session_id.is_empty() || session_ids.contains(&notice.session_id)
374}
375
376/// How long ago a durable record was written, when its timestamp can be read.
377///
378/// A record written in the future — a clock that moved backwards, a store
379/// written by another host — reports no age rather than a wrapped one, so the
380/// readiness wait leaves those sessions alone.
381fn record_age(updated_at: &str, now: chrono::DateTime<chrono::Utc>) -> Option<Duration> {
382    let written = chrono::DateTime::parse_from_rfc3339(updated_at).ok()?;
383    now.signed_duration_since(written.with_timezone(&chrono::Utc))
384        .to_std()
385        .ok()
386}