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    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                // A busy session is not replaced, and that needs no comment.
256                // An idle one that is still not replaceable would otherwise
257                // be skipped without a trace (B-5).
258                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
303/// How long ago a durable record was written, when its timestamp can be read.
304///
305/// A record written in the future — a clock that moved backwards, a store
306/// written by another host — reports no age rather than a wrapped one, so the
307/// readiness wait leaves those sessions alone.
308fn 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}