1use super::*;
2
3impl RuntimeState {
4 pub async fn reload_controller(&self) -> Result<()> {
5 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 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 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 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 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 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 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
363pub(super) fn notice_reaches_workspace(
370 notice: &RuntimeNotice,
371 session_ids: &BTreeSet<String>,
372) -> bool {
373 notice.session_id.is_empty() || session_ids.contains(¬ice.session_id)
374}
375
376fn 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}