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,
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            self.push_notice(session_id, detail);
80        }
81        Ok(())
82    }
83
84    pub(super) async fn publish_session(
85        &self,
86        session_id: String,
87        view: ManagedSessionView,
88    ) -> Result<()> {
89        let connected = view.connected;
90        let has_snapshot = view.snapshot.is_some();
91        tracing::debug!(
92            %session_id,
93            connected,
94            has_snapshot,
95            "daemon received a session view"
96        );
97        if let Some(snapshot) = view.snapshot.as_ref() {
98            let controller = self
99                .controller
100                .lock()
101                .unwrap_or_else(PoisonError::into_inner);
102            if let Some(session) = controller.state.sessions.get(&session_id).cloned() {
103                // An upgrade only ever runs on a quiet session, and a session
104                // in a turn publishes a view every 150 ms. Skipping those here
105                // keeps the config clone off the streaming path; the
106                // coordinator still decides, from `quiet`, whether to act.
107                let quiet =
108                    view.connected && snapshot.operational.safe_to_replace(session.harness_kind);
109                if quiet {
110                    self.worker_upgrade_observer
111                        .observe(WorkerUpgradeObservation {
112                            session: session.clone(),
113                            config: controller.config.clone(),
114                            worker_build: snapshot.worker_build.clone(),
115                            quiet,
116                        });
117                }
118                self.recovery_observer.observe(RecoveryObservation {
119                    checkpoint_safe: snapshot
120                        .operational
121                        .safe_for_checkpoint(session.harness_kind),
122                    session,
123                    config: controller.config.clone(),
124                    latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
125                    execution: snapshot.materialized.execution,
126                });
127            }
128        }
129        self.sessions
130            .lock()
131            .unwrap_or_else(PoisonError::into_inner)
132            .insert(
133                session_id.clone(),
134                RuntimeSessionView::from_managed(session_id, view),
135            );
136        reach_test_hook("relay_projection_before_revision_publication").await?;
137        self.publish_revision();
138        Ok(())
139    }
140
141    pub(super) async fn runtime_snapshot(
142        &self,
143        workspace_id: &str,
144        after_revision: u64,
145        all_workspaces: bool,
146    ) -> Result<RuntimeSnapshot> {
147        let mut revisions = self.revisions.subscribe();
148        if *revisions.borrow_and_update() <= after_revision {
149            let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
150        }
151        let revision = self.revisions.current();
152        let moves = blocking(crate::database::load_move_operations).await?;
153        let workspace_names = blocking(crate::database::list_workspaces)
154            .await?
155            .into_iter()
156            .map(|workspace| (workspace.id, workspace.name))
157            .collect();
158        let session_ids = if all_workspaces {
159            self.controller
160                .lock()
161                .unwrap_or_else(PoisonError::into_inner)
162                .state
163                .sessions
164                .keys()
165                .cloned()
166                .collect::<BTreeSet<_>>()
167        } else {
168            let workspace_id = workspace_id.to_owned();
169            blocking(move || crate::database::session_ids_for_workspace(&workspace_id))
170                .await?
171                .into_iter()
172                .collect()
173        };
174        let sessions = self
175            .sessions
176            .lock()
177            .unwrap_or_else(PoisonError::into_inner)
178            .iter()
179            .filter(|(session_id, _)| session_ids.contains(*session_id))
180            .map(|(_, view)| view.clone())
181            .collect();
182        // Match the controller -> lifecycle lock order used by worker polling.
183        // Completion reloads records before publishing its result, so holding
184        // this guard prevents an absent operation paired with older records.
185        let controller = self
186            .controller
187            .lock()
188            .unwrap_or_else(PoisonError::into_inner);
189        let lifecycles = self
190            .lifecycle
191            .lock()
192            .unwrap_or_else(PoisonError::into_inner)
193            .iter()
194            .filter(|(session_id, active)| {
195                (all_workspaces
196                    || session_ids.contains(*session_id)
197                    || active.resume_workspace_id.as_deref() == Some(workspace_id))
198                    && active.is_visible()
199            })
200            .map(|(session_id, active)| RuntimeLifecycleView {
201                operation_id: active.operation_id.clone(),
202                cancellable: active.is_cancellable()
203                    && lifecycle_cancellable(
204                        active.kind,
205                        durable_session_state(&controller, session_id),
206                    ),
207                session_id: session_id.clone(),
208                kind: active.kind.into(),
209                started_at_epoch_seconds: active.started_at_epoch_seconds,
210                active_stages: active
211                    .active_stages
212                    .iter()
213                    .map(|(stage, (_, started_at))| (*stage, *started_at))
214                    .collect(),
215                resume_destination: active.resume_destination.clone(),
216                notice: active.notice.clone(),
217            })
218            .collect();
219        let reviews = self
220            .review_host
221            .views()
222            .into_iter()
223            .filter(|review| session_ids.contains(&review.session_id))
224            .collect();
225        let notices = self
226            .notices
227            .lock()
228            .unwrap_or_else(PoisonError::into_inner)
229            .iter()
230            .filter(|notice| session_ids.contains(&notice.session_id))
231            .cloned()
232            .collect();
233        let records = runtime_records_for_workspace(&controller, &session_ids);
234        let subagents = runtime_subagents_for_workspace(&controller, &records);
235        Ok(RuntimeSnapshot {
236            workspace_names,
237            moves: moves
238                .into_iter()
239                .filter(|operation| session_ids.contains(&operation.selection.session_id))
240                .collect(),
241            revision,
242            config: controller.config.clone(),
243            records,
244            sessions,
245            lifecycles,
246            reviews,
247            notices,
248            subagents,
249        })
250    }
251}