mj_controller/daemon/
snapshot.rs1use super::*;
2
3pub(super) struct BackgroundPolicyState {
6 quiet: bool,
7 checkpoint_wait: Option<mj_core::activity::CheckpointWait>,
8 latest_completed_turn_ordinal: Option<u64>,
9 worker_build: Option<String>,
10 idle_replacement_blockers: Vec<&'static str>,
13}
14
15impl RuntimeState {
16 fn observe_background_policy(
17 &self,
18 session: &SessionRecord,
19 config: &Config,
20 policy: &BackgroundPolicyState,
21 ) {
22 if policy.quiet {
23 self.worker_upgrade_observer
24 .observe(WorkerUpgradeObservation {
25 session: session.clone(),
26 config: config.clone(),
27 worker_build: policy.worker_build.clone(),
28 quiet: true,
29 });
30 }
31 self.recovery_observer.observe(RecoveryObservation {
32 checkpoint_wait: policy.checkpoint_wait,
33 session: session.clone(),
34 config: config.clone(),
35 latest_completed_turn_ordinal: policy.latest_completed_turn_ordinal,
36 });
37 }
38
39 pub(super) fn refresh_background_policies(&self) {
42 let mut owner = self.owner();
43 let records = owner.controller().state.sessions.clone();
44 let config = owner.controller().config.clone();
45 owner.background_policies.retain(|session_id, policy| {
46 let Some(session) = records
47 .get(session_id)
48 .filter(|session| session.state.has_live_worker())
49 else {
50 return false;
51 };
52 self.observe_background_policy(session, &config, policy);
53 true
54 });
55 }
56
57 pub async fn reload_controller(&self) -> Result<()> {
58 let _mutation = self.config_mutation.lock().await;
61 if self.committed.is_some() {
62 let config = tokio::task::spawn_blocking(Config::load)
63 .await
64 .context("daemon configuration reader panicked")??;
65 self.owner().install_config(config);
66 } else {
67 let controller_loader = self.controller_loader;
68 let controller = tokio::task::spawn_blocking(controller_loader)
69 .await
70 .context("daemon controller loader panicked")??;
71 self.owner().install_controller(controller);
72 }
73 self.publish_revision();
74 Ok(())
75 }
76
77 pub(super) fn missing_target_record(
80 &self,
81 session_id: &str,
82 view: &ManagedSessionView,
83 ) -> Option<(String, String)> {
84 let Some(ViewError::TargetMissing(detail)) = &view.error else {
85 return None;
86 };
87 if view.connected {
88 return None;
89 }
90 let controller_owner = self.owner();
91 if controller_owner
92 .lifecycle
93 .get(session_id)
94 .is_some_and(|active| active.is_running())
95 {
96 return None;
97 }
98 let controller = controller_owner.controller();
99 let session = controller.state.sessions.get(session_id)?;
100 if !matches!(
101 session.state,
102 SessionState::Running | SessionState::Disconnected
103 ) {
104 return None;
105 }
106 Some((detail.clone(), session.updated_at.clone()))
107 }
108
109 pub(super) async fn persist_missing_target(
110 self: &Arc<Self>,
111 session_id: &str,
112 detail: String,
113 observed_updated_at: String,
114 ) -> Result<()> {
115 let changed = blocking({
116 let session_id = session_id.to_owned();
117 let detail = detail.clone();
118 move || {
119 crate::database::mark_session_target_missing_if_current(
120 &session_id,
121 &detail,
122 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
123 &observed_updated_at,
124 )
125 }
126 })
127 .await?;
128 if changed.is_some() {
129 self.reload_controller().await?;
130 if changed == Some(SessionState::Lost) {
134 self.discard_lost_session(session_id.to_owned()).await;
135 } else {
136 self.push_notice(session_id, detail);
137 }
138 }
139 Ok(())
140 }
141
142 pub(super) fn sessions_without_a_usable_harness(&self) -> Vec<UnreadySession> {
149 let now = chrono::Utc::now();
150 let observations = {
151 let owner = self.owner();
152 owner
153 .indexes
154 .running
155 .keys()
156 .filter(|id| {
157 !owner
158 .lifecycle
159 .get(*id)
160 .is_some_and(|active| active.is_running())
161 })
162 .filter(|id| !owner.close_requested.contains(*id))
163 .filter_map(|id| owner.controller().state.sessions.get(id))
164 .map(|record| {
165 let view = owner.sessions.get(&record.id);
166 ReadinessObservation {
167 harness_ready: view.is_some_and(|view| {
168 view.operational.as_ref().is_some_and(
169 mj_core::relay::RelayOperationalState::native_session_is_ready,
170 )
171 }),
172 record_age: record_age(&record.updated_at, now),
173 updated_at: record.updated_at.clone(),
174 detail: view
175 .and_then(|view| view.error.as_ref())
176 .map(|error| error.detail().to_owned()),
177 session_id: record.id.clone(),
178 }
179 })
180 .collect::<Vec<_>>()
181 };
182 self.harness_readiness
183 .lock()
184 .unwrap_or_else(PoisonError::into_inner)
185 .observe(std::time::Instant::now(), observations)
186 }
187
188 pub(super) async fn fail_unready_session(&self, unready: UnreadySession) {
191 let cause = unready.cause();
192 let session_id = unready.session_id.clone();
193 tracing::warn!(%session_id, waited_seconds = unready.waited.as_secs(), "{cause}");
194 let applied = blocking({
195 let session_id = session_id.clone();
196 let cause = cause.clone();
197 move || {
198 let mut controller = Controller::load()?;
199 controller.fail_unready_session(&session_id, &cause, &unready.observed_updated_at)
200 }
201 })
202 .await;
203 match applied {
204 Ok(true) => {
205 if let Err(error) = self.reload_controller().await {
206 tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after failing an unready session");
207 }
208 self.push_notice(&session_id, cause);
209 self.publish_revision();
210 }
211 Ok(false) => {}
212 Err(error) => tracing::warn!(
213 %session_id,
214 error = format!("{error:#}"),
215 "could not record that a session's harness never became usable"
216 ),
217 }
218 }
219
220 pub(super) async fn publish_session(
221 &self,
222 session_id: String,
223 view: ManagedSessionView,
224 ) -> Result<()> {
225 let connected = view.connected;
226 let has_snapshot = view.snapshot.is_some();
227 tracing::debug!(
228 %session_id,
229 connected,
230 has_snapshot,
231 "daemon received a session view"
232 );
233 {
234 let mut owner = self.owner();
235 let controller = owner.controller();
236 if !controller.state.sessions.contains_key(&session_id)
237 || !owner.runs_relay_actor(&session_id)
238 {
239 return Ok(());
240 }
241 if view.connected
242 && let Some(snapshot) = view.snapshot.as_ref()
243 && let Some(session) = controller.state.sessions.get(&session_id)
244 && session.state.has_live_worker()
245 {
246 let quiet = snapshot.operational.quiet();
247 if !quiet.is_yes() {
248 tracing::debug!(
249 session_id = %session_id,
250 reason = quiet.reason(),
251 "session is not quiet; upgrade, checkpoint and move wait"
252 );
253 }
254 let facts = snapshot.operational.facts();
255 let idle_replacement_blockers = if mj_core::activity::can_submit(&facts)
259 && !mj_core::activity::safe_to_replace(&facts, session.harness_kind)
260 {
261 mj_core::activity::replacement_blockers(&facts, session.harness_kind)
262 } else {
263 Vec::new()
264 };
265 if !idle_replacement_blockers.is_empty()
266 && owner
267 .background_policies
268 .get(&session_id)
269 .is_none_or(|previous| {
270 previous.idle_replacement_blockers != idle_replacement_blockers
271 })
272 {
273 tracing::info!(
274 %session_id,
275 harness = ?session.harness_kind,
276 worker_build = ?snapshot.worker_build,
277 reasons = ?idle_replacement_blockers,
278 "an idle worker is not replaced"
279 );
280 }
281 let policy = BackgroundPolicyState {
282 idle_replacement_blockers,
283 quiet: snapshot.operational.safe_to_replace(session.harness_kind),
284 checkpoint_wait: snapshot
285 .operational
286 .routine_checkpoint_wait(session.harness_kind),
287 latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
288 worker_build: snapshot.worker_build.clone(),
289 };
290 self.observe_background_policy(session, &controller.config, &policy);
291 owner.background_policies.insert(session_id.clone(), policy);
292 } else {
293 owner.background_policies.remove(&session_id);
294 }
295 owner.publish_view(session_id, view);
296 }
297 reach_test_hook("relay_projection_before_revision_publication").await?;
298 self.publish_revision();
299 Ok(())
300 }
301}
302
303fn record_age(updated_at: &str, now: chrono::DateTime<chrono::Utc>) -> Option<Duration> {
309 let written = chrono::DateTime::parse_from_rfc3339(updated_at).ok()?;
310 now.signed_duration_since(written.with_timezone(&chrono::Utc))
311 .to_std()
312 .ok()
313}