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_wait: Option<mj_core::activity::CheckpointWait>,
8    latest_completed_turn_ordinal: Option<u64>,
9    worker_build: Option<String>,
10    /// Why the worker is not replaceable while the session is otherwise idle;
11    /// empty otherwise. Kept to log a change once, not on every view.
12    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    /// Quiet workers need no new view to retry a skipped or delayed operation.
40    /// This only queues observations; coordinators perform the actual I/O.
41    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        // Serialize installs so an earlier phone publication cannot overwrite
59        // a later completed lifecycle with the controller snapshot it loaded.
60        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    /// Definitive missing-target evidence belongs to the daemon, including
78    /// when no terminal is attached. Generic connection failures stay transient.
79    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            // A lost session has no checkpoint and no target, so its record
131            // offers nothing but Destroy. Report it once and discard it
132            // instead of leaving a tombstone behind.
133            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    /// One pass of the harness readiness wait: which live sessions have run
143    /// out of time to become usable.
144    ///
145    /// Locks are taken one at a time and nothing here touches the store, so
146    /// this belongs on the daemon's background tick. The write it leads to
147    /// does not; see [`Self::fail_unready_session`].
148    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    /// Fail one session whose harness never became usable, so it stops looking
189    /// like a session a driver can wait for.
190    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    /// Persist a worker-reported preparation failure against the session
221    /// revision observed by the upgrade coordinator. Lifecycle work that has
222    /// since changed that revision remains the owner of the outcome.
223    pub(super) async fn fail_harness_preparation(
224        &self,
225        session_id: String,
226        failure: crate::controller::HarnessPreparationFailure,
227        observed_updated_at: String,
228    ) {
229        let cause = failure.to_string();
230        let applied = blocking({
231            let session_id = session_id.clone();
232            let cause = cause.clone();
233            move || {
234                let mut controller = Controller::load()?;
235                controller.fail_unready_session(&session_id, &cause, &observed_updated_at)
236            }
237        })
238        .await;
239        match applied {
240            Ok(true) => {
241                if let Err(error) = self.reload_controller().await {
242                    tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after harness preparation failed");
243                }
244                self.push_notice(&session_id, cause);
245                self.publish_revision();
246            }
247            Ok(false) => {}
248            Err(error) => tracing::warn!(
249                %session_id,
250                error = format!("{error:#}"),
251                "could not record worker-reported harness preparation failure"
252            ),
253        }
254    }
255
256    pub(super) async fn publish_session(
257        &self,
258        session_id: String,
259        view: ManagedSessionView,
260    ) -> Result<()> {
261        let connected = view.connected;
262        let has_snapshot = view.snapshot.is_some();
263        tracing::debug!(
264            %session_id,
265            connected,
266            has_snapshot,
267            "daemon received a session view"
268        );
269        {
270            let mut owner = self.owner();
271            let controller = owner.controller();
272            if !controller.state.sessions.contains_key(&session_id)
273                || !owner.runs_relay_actor(&session_id)
274            {
275                return Ok(());
276            }
277            if view.connected
278                && let Some(snapshot) = view.snapshot.as_ref()
279                && let Some(session) = controller.state.sessions.get(&session_id)
280                && session.state.has_live_worker()
281            {
282                let quiet = snapshot.operational.quiet();
283                if !quiet.is_yes() {
284                    tracing::debug!(
285                        session_id = %session_id,
286                        reason = quiet.reason(),
287                        "session is not quiet; upgrade, checkpoint and move wait"
288                    );
289                }
290                let facts = snapshot.operational.facts();
291                // A busy session is not replaced, and that needs no comment.
292                // An idle one that is still not replaceable would otherwise
293                // be skipped without a trace (B-5).
294                let idle_replacement_blockers = if mj_core::activity::can_submit(&facts)
295                    && !mj_core::activity::safe_to_replace(&facts, session.harness_kind)
296                {
297                    mj_core::activity::replacement_blockers(&facts, session.harness_kind)
298                } else {
299                    Vec::new()
300                };
301                if !idle_replacement_blockers.is_empty()
302                    && owner
303                        .background_policies
304                        .get(&session_id)
305                        .is_none_or(|previous| {
306                            previous.idle_replacement_blockers != idle_replacement_blockers
307                        })
308                {
309                    tracing::info!(
310                        %session_id,
311                        harness = ?session.harness_kind,
312                        worker_build = ?snapshot.worker_build,
313                        reasons = ?idle_replacement_blockers,
314                        "an idle worker is not replaced"
315                    );
316                }
317                let policy = BackgroundPolicyState {
318                    idle_replacement_blockers,
319                    quiet: snapshot.operational.safe_to_replace(session.harness_kind),
320                    checkpoint_wait: snapshot
321                        .operational
322                        .routine_checkpoint_wait(session.harness_kind),
323                    latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
324                    worker_build: snapshot.worker_build.clone(),
325                };
326                self.observe_background_policy(session, &controller.config, &policy);
327                owner.background_policies.insert(session_id.clone(), policy);
328            } else {
329                owner.background_policies.remove(&session_id);
330            }
331            owner.publish_view(session_id, view);
332        }
333        reach_test_hook("relay_projection_before_revision_publication").await?;
334        self.publish_revision();
335        Ok(())
336    }
337}
338
339/// How long ago a durable record was written, when its timestamp can be read.
340///
341/// A record written in the future — a clock that moved backwards, a store
342/// written by another host — reports no age rather than a wrapped one, so the
343/// readiness wait leaves those sessions alone.
344fn record_age(updated_at: &str, now: chrono::DateTime<chrono::Utc>) -> Option<Duration> {
345    let written = chrono::DateTime::parse_from_rfc3339(updated_at).ok()?;
346    now.signed_duration_since(written.with_timezone(&chrono::Utc))
347        .to_std()
348        .ok()
349}