1use super::*;
2
3pub(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 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 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 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 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 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 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 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
420pub(super) fn notice_reaches_workspace(
427 notice: &RuntimeNotice,
428 session_ids: &BTreeSet<String>,
429) -> bool {
430 notice.session_id.is_empty() || session_ids.contains(¬ice.session_id)
431}
432
433fn 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}