Skip to main content

mj_controller/daemon/
snapshot.rs

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